Airflow多并行任务组DAG构建求助:代码生成结构与预期不符
Airflow DAG结构不符合预期的问题分析与修复
问题说明
预期构建的DAG结构为:Start → 第一个TaskGroup(3个并行任务)→ 第二个TaskGroup(3个并行任务)→ 第三个TaskGroup(1个任务)→ Stop(TaskGroup串行执行)
实际生成的DAG结构为:Start同时触发所有TaskGroup,所有TaskGroup完成后触发Stop(TaskGroup并行执行)
问题根源
原代码中直接使用start >> task_group >> stop,其中task_group是TaskGroup对象的列表。Airflow对列表类型的依赖会默认解析为并行关系,即Start会同时触发列表中所有TaskGroup,而非按顺序串行执行。
修复方案
手动构建TaskGroup之间的串行依赖链,具体步骤:
- 将Start与第一个TaskGroup建立依赖
- 遍历TaskGroup列表,依次将前一个TaskGroup与后一个TaskGroup建立依赖
- 将最后一个TaskGroup与Stop建立依赖
修改后的完整代码
import math from airflow import DAG from airflow.operators.bash import BashOperator from airflow.operators.dummy import DummyOperator from airflow.utils.task_group import TaskGroup from airflow.utils.trigger_rule import TriggerRule task_list = [1, 2, 3, 4, 5, 6, 7] parallel_task_per_grp = 3 task_grp_cnt = math.ceil(len(task_list)/parallel_task_per_grp) task_grp_names = [f"grp_{tgc}" for tgc in range(task_grp_cnt)] with DAG( dag_id=dag_name_id, default_args=default_args, start_date=start_date, schedule_interval=None, # 格式为 cron 表达式:分 时 日 月 周 tags=["adhoc_dag_ing_conf_sys_profile_col_upd"], catchup=False ) as dag: start = BashOperator(task_id="start", bash_command='echo "starting batch processing"', do_xcom_push=False) stop = BashOperator(task_id="stop", bash_command='echo "stopping batch processing"', trigger_rule=TriggerRule.NONE_SKIPPED, do_xcom_push=False) start_pos, stop_pos = 0, parallel_task_per_grp task_group = [] for task_grp_name in task_grp_names: with TaskGroup(group_id=task_grp_name) as tg1: # 生成当前TaskGroup内的并行任务 [DummyOperator(task_id=f"task_{tl}") for tl in task_list[start_pos:stop_pos]] task_group.append(tg1) start_pos, stop_pos = start_pos+parallel_task_per_grp, stop_pos+parallel_task_per_grp # 构建串行依赖关系 if task_group: start >> task_group[0] # 依次链接后续TaskGroup for idx in range(1, len(task_group)): task_group[idx-1] >> task_group[idx] # 最后一个TaskGroup链接Stop task_group[-1] >> stop
内容的提问来源于stack exchange,提问作者Santanu Ghosh
相关产品推荐
相关产品推荐

