工作流调度框架

Data Engineering 入门

Vincent Vankrunkelsven

Data Engineer, DataCamp

一个管道示例

 

使用 Spark 从 csv 提取的数据管道示例

如何调度?

  • 手动
  • cron 调度工具
  • 依赖关系怎么办?
Data Engineering 入门

DAGs

有向无环图

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

DAG 示例

Data Engineering 入门

常用工具

 

  • Linux 的 cron
  • Prefect 和 Dagster
  • Apache Airflow
Data Engineering 入门

Apache Airflow 标志

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

Airflow:DAG 示例

 

Airflow DAG 示例

Data Engineering 入门

Airflow:代码示例

@dag(dag_id="example_dag",
     start_date=datetime(2024, 1, 1),
     schedule="0 * * * *")
def example_dag():

@task def start_cluster(): ... @task def ingest_customer_data(): ... @task def ingest_product_data(): ... @task def enrich_customer_data(): ...
Data Engineering 入门

Airflow:代码示例

@dag(dag_id="example_dag", ...)
def example_dag():
    ...
    # Set up dependency flow
    cluster = start_cluster()
    customers = ingest_customer_data()
    products = ingest_product_data()
    cluster >> [customers, products]
    [customers, products] >> enrich_customer_data()
# Run the DAG
example_dag()
Data Engineering 入门

让我们一起练习吧!

Data Engineering 入门

Preparing Video For Download...