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
相关产品推荐
相关产品推荐

