工作流调度框架

Data Engineering 入门

Vincent Vankrunkelsven

Data Engineer @ DataCamp

示例流水线

 

使用 Spark 从 CSV 抽取的简单流水线示例

如何调度?

  • 手动
  • cron 调度工具
  • 依赖如何处理?
Data Engineering 入门

DAG

有向无环图

  • 一组节点
  • 有向边
  • 无环

DAG 示例

Data Engineering 入门

常用工具

 

  • Linux 的 cron
  • Spotify 的 Luigi
  • Apache Airflow
Data Engineering 入门

Apache Airflow 徽标

  • 诞生于 Airbnb
  • 使用 DAG
  • 基于 Python
Data Engineering 入门

Airflow:DAG 示例

 

Airflow DAG 示例

Data Engineering 入门

Airflow:代码示例

# Create the DAG object
dag = DAG(dag_id="example_dag", ..., schedule_interval="0 * * * *")

# Define operations start_cluster = StartClusterOperator(task_id="start_cluster", dag=dag) ingest_customer_data = SparkJobOperator(task_id="ingest_customer_data", dag=dag) ingest_product_data = SparkJobOperator(task_id="ingest_product_data", dag=dag) enrich_customer_data = PythonOperator(task_id="enrich_customer_data", ..., dag = dag)
# Set up dependency flow start_cluster.set_downstream(ingest_customer_data) ingest_customer_data.set_downstream(enrich_customer_data) ingest_product_data.set_downstream(enrich_customer_data)
Data Engineering 入门

Passons à la pratique !

Data Engineering 入门

Preparing Video For Download...