You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow嵌套循环内DAG迭代无法顺序执行问题求助

问题描述

在Airflow中通过嵌套循环为每个分类-指标组合生成一组任务,但所有迭代的任务组并行执行,不符合预期的串行顺序。具体需求如下:

  • 加载包含分类和对应指标的JSON配置
  • 每个指标需依次执行:create_table → insert_data_into_tmp → bq_to_gcs → gcs_to_S3
  • 期望所有指标任务组串行执行:完成一个指标的全流程后,再启动下一个指标的任务组

当前DAG运行时,所有任务组的create_table并行启动,随后所有insert_data_into_tmp并行,以此类推,完全不符合串行要求。

示例配置JSON:

{
    "reports": [
        {
            "name": "category_1",
            "metrics": [
                {
                    "period": "day",
                    "schema": "metric_1.json",
                    "kpi": "metric_1"
                },
                {
                    "period": "month",
                    "schema": "metric_2.json",
                    "kpi": "metric_2"
                }
            ]
        },
        {
            "name": "category_2",
            "metrics": [
                {
                    "period": "month",
                    "schema": "metric_3.json",
                    "kpi": "metric_3"
                },
                {
                    "period": "month",
                    "schema": "metric_4.json",
                    "kpi": "metric_4"
                }
            ]
        }
    ]
}

当前DAG代码:

with DAG(
          dag_id='daily',
          default_args=default_args,
          start_date=pendulum.datetime(2024, 4, 10)
)  as dag:
    dag.doc_md = doc_md

    # Define start and end tasks
    start_task = EmptyOperator(task_id="start")
    end_task = EmptyOperator(task_id="end")

    file_path = "config.json"
    with open(file_path, 'r') as f:
        data = json.load(f)
        print(data)
        
    config_data = data['reports']

    for name in config_data:
        category = name['name']
        metrics = name['metrics']
        for metric in metrics:
            schema =metric['schema']
            kpi = metric['kpi']

            with TaskGroup(f'Writing_{kpi}_table') as export:
            
                schema_file_path = "/schema/"+schema
                with open(schema_file_path, 'r') as f:
                    loadedSchema = json.load(f)

                create_table = BigQueryCreateEmptyTableOperator(
                    task_id=f"create_{kpi}_tmp",
                    dataset_id=DATASET,
                    table_id=f"tmp_{kpi}",
                    project_id=PROJECT_ID,
                    schema_fields=loadedSchema,
                    gcp_conn_id='gcs',
                    exists_ok=True,
                    dag=dag
                )

                insert_data_into_tmp = BigQueryInsertJobOperator(
                    task_id=f"insert_into_{kpi}_tmp",
                    configuration={
                            "query": {
                                "query": f""" INSERT INTO {PROJECT_ID}.{DATASET}.tmp_{kpi}
                                SELECT * FROM {PROJECT_ID}.{DATASET}.{kpi}  """,
                                "useLegacySql": False,
                            }
                        },
                    dag=dag
                )

         
                bq_to_gcs = BigQueryToGCSOperator (
                    task_id=f"{kpi}_to_GCS",
                    source_project_dataset_table=f"{PROJECT_ID}.{DATASET}.tmp_{kpi}",
                    destination_cloud_storage_uris=[f'gs://test/{kpi}.parquet'],
                    export_format='PARQUET',
                    compression='SNAPPY',
                    gcp_conn_id='gcs',
                    dag=dag
                )

                gcs_to_S3 = GCSToS3Operator (
                    task_id=f"{category}_to_S3",
                    gcs_bucket="test-dev",
                    gcp_conn_id='gcs',
                    dest_aws_conn_id='aws_s3',
                    dest_s3_key =f"s3:test",
                    keep_directory_structure =True ,
                    replace=False,
                    dag=dag
                )

                start_task >> export  >> end_task
解决方案

当前代码的核心问题是所有TaskGroup都直接绑定到start_task和end_task,导致Airflow认为所有任务组都是并行分支。要实现串行执行,需要将每个TaskGroup按顺序串联起来,即上一个任务组完成后再启动下一个。同时,每个TaskGroup内部也需要明确任务的依赖关系,确保组内任务按顺序执行。

修改后的完整代码:

with DAG(
          dag_id='daily',
          default_args=default_args,
          start_date=pendulum.datetime(2024, 4, 10)
)  as dag:
    dag.doc_md = doc_md

    # Define start and end tasks
    start_task = EmptyOperator(task_id="start")
    end_task = EmptyOperator(task_id="end")

    file_path = "config.json"
    with open(file_path, 'r') as f:
        data = json.load(f)
        
    config_data = data['reports']

    # 初始化上游任务为start_task
    upstream_task = start_task

    for name in config_data:
        category = name['name']
        metrics = name['metrics']
        for metric in metrics:
            schema = metric['schema']
            kpi = metric['kpi']

            with TaskGroup(f'Writing_{kpi}_table') as export:
                schema_file_path = "/schema/" + schema
                with open(schema_file_path, 'r') as f:
                    loadedSchema = json.load(f)

                create_table = BigQueryCreateEmptyTableOperator(
                    task_id=f"create_{kpi}_tmp",
                    dataset_id=DATASET,
                    table_id=f"tmp_{kpi}",
                    project_id=PROJECT_ID,
                    schema_fields=loadedSchema,
                    gcp_conn_id='gcs',
                    exists_ok=True
                )

                insert_data_into_tmp = BigQueryInsertJobOperator(
                    task_id=f"insert_into_{kpi}_tmp",
                    configuration={
                            "query": {
                                "query": f""" INSERT INTO {PROJECT_ID}.{DATASET}.tmp_{kpi}
                                SELECT * FROM {PROJECT_ID}.{DATASET}.{kpi}  """,
                                "useLegacySql": False,
                            }
                        }
                )

                bq_to_gcs = BigQueryToGCSOperator (
                    task_id=f"{kpi}_to_GCS",
                    source_project_dataset_table=f"{PROJECT_ID}.{DATASET}.tmp_{kpi}",
                    destination_cloud_storage_uris=[f'gs://test/{kpi}.parquet'],
                    export_format='PARQUET',
                    compression='SNAPPY',
                    gcp_conn_id='gcs'
                )

                gcs_to_S3 = GCSToS3Operator (
                    task_id=f"{category}_to_S3",
                    gcs_bucket="test-dev",
                    gcp_conn_id='gcs',
                    dest_aws_conn_id='aws_s3',
                    dest_s3_key=f"s3:test",
                    keep_directory_structure=True,
                    replace=False
                )

                # 定义TaskGroup内部的任务依赖:按顺序执行
                create_table >> insert_data_into_tmp >> bq_to_gcs >> gcs_to_S3

            # 将当前TaskGroup与上游任务串联
            upstream_task >> export
            # 更新上游任务为当前TaskGroup,供下一次循环使用
            upstream_task = export

    # 将最后一个TaskGroup连接到end_task
    upstream_task >> end_task

关键修改说明:

  • 新增upstream_task变量,用于跟踪当前需要连接的上游节点,实现任务组之间的串行
  • 每个TaskGroup内部明确任务依赖:create_table >> insert_data_into_tmp >> bq_to_gcs >> gcs_to_S3,确保组内任务按指定顺序执行
  • 循环中每次将upstream_task指向当前TaskGroup,保证一个任务组完成后才启动下一个
  • 循环结束后将最后一个TaskGroup连接到end_task,完成整个DAG的流程闭环

内容的提问来源于stack exchange,提问作者user3455979

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 21:30:31