如何让Airflow中并行的weather至end任务流按顺序执行?
问题:让Airflow中多组任务流串行执行
当前我的Airflow DAG中,从weather到weather_end存在多组并行任务流(对应summers、winters、spring三个分组),触发DAG时这些任务流会同时运行。我希望改成逐个触发执行——完成第一组后再运行第二组,以此类推。
原DAG代码如下:
start = DummyOperator(task_id='start') finish = DummyOperator(task_id='finish', dag=dag, trigger_rule=TriggerRule.NONE_FAILED) def running_tasks(): print("boom") weather = {'summers', 'winters', 'spring'} with TaskGroup(group_id='compute') as compute: w = DummyOperator(task_id='weather', dag=dag) end = DummyOperator(task_id='weather_end', dag=dag) with TaskGroup(group_id='weather_based_tasks') as weather_based_tasks: for name in weather: with TaskGroup(group_id=name, tooltip=name) as weather_name: weather_specific = PythonOperator(dag=dag, task_id=name, python_callable=running_tasks, pool="weather_pool" ) with TaskGroup(group_id='tasks', tooltip='tasks') as task_name: task01 = DummyOperator(task_id='task_01') task02 = DummyOperator(task_id='task_02') [task01, task02] w >> weather_specific >> task_name >> end start >> compute >> finish
解决方案
要实现串行执行,核心是让每个天气分组的任务流依赖前一个分组的结束节点,同时保证任务执行顺序固定。具体修改如下:
- 将无序的
set类型天气集合改为list,确保执行顺序可预测 - 新增变量跟踪前一个分组的结束任务,让当前分组的起始任务依赖该变量
- 调整全局
end节点的依赖,改为依赖最后一个分组的结束任务
修改后的完整代码:
start = DummyOperator(task_id='start') finish = DummyOperator(task_id='finish', dag=dag, trigger_rule=TriggerRule.NONE_FAILED) def running_tasks(): print("boom") # 改为有序列表,保证执行顺序固定 weather = ['summers', 'winters', 'spring'] with TaskGroup(group_id='compute') as compute: w = DummyOperator(task_id='weather', dag=dag) end = DummyOperator(task_id='weather_end', dag=dag) with TaskGroup(group_id='weather_based_tasks') as weather_based_tasks: # 初始化前一个任务为weather节点,第一个分组直接依赖它 prev_task = w for name in weather: with TaskGroup(group_id=name, tooltip=name) as weather_name: weather_specific = PythonOperator(dag=dag, task_id=name, python_callable=running_tasks, pool="weather_pool" ) with TaskGroup(group_id='tasks', tooltip='tasks') as task_name: task01 = DummyOperator(task_id='task_01') task02 = DummyOperator(task_id='task_02') # 建立当前分组内的任务依赖 weather_specific >> [task01, task02] # 当前分组的起始任务依赖前一个任务 prev_task >> weather_specific # 更新前一个任务为当前分组的结束节点(task_name) prev_task = task_name # 全局end节点依赖最后一个分组的结束任务 prev_task >> end start >> compute >> finish
关键改动说明
- 把
weather从set改为list:集合是无序的,Airflow无法保证执行顺序,列表可以固定分组的执行顺序 - 新增
prev_task变量:用于跟踪上一个分组的结束节点,让每个新分组的起始任务(weather_specific)依赖它,实现串行 - 调整全局
end的依赖:确保所有分组执行完成后才触发end节点 - 修正分组内的任务依赖:原代码中
[task01, task02]没有明确依赖,改为weather_specific >> [task01, task02]让逻辑更清晰
内容的提问来源于stack exchange,提问作者sammy morgan
相关产品推荐
相关产品推荐

