如何在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的执行流程变为:
- 执行
start_task - 执行
workflow_US内的call_first_task_US→call_second_task_US,完成整个US任务组 - 执行
workflow_FR内的call_first_task_FR→call_second_task_FR,完成整个FR任务组 - 执行
end_task
这样就实现了你需要的"垂直执行"顺序,确保US的所有任务完成后才启动FR的任务。
内容的提问来源于stack exchange,提问作者Antoine F
相关产品推荐
相关产品推荐

