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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 00:02:54