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

Airflow Task Group执行顺序异常:如何按指定流程配置DAG?

问题原因及解决方案

为什么task_3a会立即执行

你的DAG存在两个核心问题导致task_3a提前执行:

  1. 依赖链断裂:分支任务branching_task仅关联了task_2,未将task_3纳入分支逻辑,且task_3没有明确的上游依赖(原代码中task_3仅被task_2触发,当分支不选择task_2时,task_3无依赖会直接启动),进而导致后续的task_group也提前执行。
  2. task_group任务串行:你在task_after_3_before_4_group中返回task_3a() >> task_3b(),这会让两个任务串行执行,不符合你需要并行的要求。

调整后的完整代码

以下是符合你指定执行顺序的修改版代码:

import logging
from airflow.decorators import task, dag, task_group
from airflow.utils.dates import days_ago

@dag(
    dag_id='taskflow_conditional_dag',
    start_date=days_ago(1),
    schedule_interval=None,
    catchup=False,
)
def my_dag():
    logger = logging.getLogger("airflow.task")

    @task
    def task_1():
        logger.info("Task 1 running")
        return "run_task_"

    @task.branch
    def branching_task(data):
        # 条件满足时执行task_2,否则跳过task_2
        if data == "run_task_2":
            return "task_2"
        # 返回空列表表示不执行任何分支任务(即跳过task_2)
        return []

    @task
    def task_2():
        logger.info("Task 2 running")
        pass

    @task(
        trigger_rule="none_failed_or_skipped"  # 不管task_2是成功还是被跳过,都执行task_3
    )
    def task_3():
        logger.info("Task 3 running")
        pass

    @task
    def task_4():
        logger.info("Task 4 running")
        pass

    @task_group()
    def task_after_3_before_4_group():
        @task
        def task_3a():
            logger.info("Task 3a running.")
            pass

        @task
        def task_3b():
            logger.info("Task 3b running.")
            pass

        # 返回任务列表表示并行执行
        return [task_3a(), task_3b()]

    @task
    def task_5():
        logger.info("Task 5 running")
        pass

    # 任务实例化
    data = task_1()
    decision = branching_task(data)
    task_2_result = task_2()
    task_3_result = task_3()
    task_group_parallel = task_after_3_before_4_group()
    task_4_result = task_4()
    task_5_result = task_5()

    # 构建依赖链
    data >> decision
    decision >> task_2_result
    # task_3同时依赖task_1和task_2,trigger_rule确保跳过task_2时仍能执行
    [data, task_2_result] >> task_3_result
    # task_3完成后启动并行的3a和3b
    task_3_result >> task_group_parallel
    # 最后依次执行task_4和task_5
    task_group_parallel >> task_4_result >> task_5_result

dag = my_dag()

关键修改说明

  1. 修复分支逻辑:branching_task在不需要执行task_2时返回空列表,标记task_2为跳过状态。
  2. 调整task_3的触发规则:设置trigger_rule="none_failed_or_skipped",确保无论task_2是成功执行还是被跳过,task_3都会启动。
  3. 修正依赖关系:让task_3同时依赖task_1和task_2,保证执行顺序符合要求。
  4. 实现任务并行:在task_group中返回任务列表[task_3a(), task_3b()],让两个任务并行执行。
  5. 简化trigger_rule:task_4和task_5无需额外触发规则,默认的all_success即可满足(因为并行任务全部完成后才会触发下游)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 23:00:02