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

如何让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

解决方案

要实现串行执行,核心是让每个天气分组的任务流依赖前一个分组的结束节点,同时保证任务执行顺序固定。具体修改如下:

  1. 将无序的set类型天气集合改为list,确保执行顺序可预测
  2. 新增变量跟踪前一个分组的结束任务,让当前分组的起始任务依赖该变量
  3. 调整全局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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 04:43:29