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

如何在Airflow中配置Task Group的垂直执行顺序?

Airflow Task Group执行顺序调整方案

当前你的DAG会并行执行两个Task Group的任务,原因是每个workflow_{country}都直接依赖start_task,且内部任务也直接关联start_task,导致两个国家的任务会同时启动。要实现先完成US的所有任务再执行FR的任务,需要调整Task Group之间的依赖关系,具体修改如下:

修改要点

  • 取消Task Group内部任务与外部start_task、end_task的直接关联,仅在Task Group内部定义任务依赖(call_first_task >> call_second_task)
  • 让Task Group之间建立顺序依赖:start_task → workflow_US → workflow_FR → end_task

修改后的完整代码

from airflow.operators.dummy_operator import DummyOperator
from airflow.operators.python_operator import PythonOperator
from airflow import DAG
from airflow.utils.task_group import TaskGroup
from pendulum import datetime
import time


def first_task(country_id: str):
    time.sleep(10)
    print(f"country = {country_id}")


def second_task(country_id: str):
    time.sleep(10)
    print(f"country = {country_id}")


default_args = {
    "start_date": datetime(2023, 10, 8)
}

with DAG(
        dag_id="my_dag",
        concurrency=1,
        schedule_interval='0 0 12 * *',
        default_args=default_args
) as dag:

    start_task = DummyOperator(
        task_id='start',
    )

    end_task = DummyOperator(
        task_id='end',
    )

    # 存储生成的Task Group,用于建立顺序依赖
    workflow_groups = []
    for country in ["US", "FR"]:
        with TaskGroup(group_id=f"workflow_{country}") as workflow:
            call_first_task = PythonOperator(
                task_id=f"call_first_task_{country}",
                python_callable=first_task,
                op_kwargs={
                    "country_id": country
                }
            )

            call_second_task = PythonOperator(
                task_id=f"call_second_task_{country}",
                python_callable=second_task,
                op_kwargs={
                    "country_id": country
                }
            )

            # 仅在Task Group内部定义任务顺序
            call_first_task >> call_second_task
        
        workflow_groups.append(workflow)
    
    # 建立整体依赖链:start -> US任务组 -> FR任务组 -> end
    start_task >> workflow_groups[0] >> workflow_groups[1] >> end_task

说明

修改后,DAG的执行流程变为:

  1. 执行start_task
  2. 执行workflow_US内的call_first_task_US → call_second_task_US,完成整个US任务组
  3. 执行workflow_FR内的call_first_task_FR → call_second_task_FR,完成整个FR任务组
  4. 执行end_task

这样就实现了你需要的"垂直执行"顺序,确保US的所有任务完成后才启动FR的任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 14:45:08