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

Airflow多链式Expand调用下动态任务映射的分组执行问题

解决方案:用动态生成TaskGroup实现分组串行执行

你说得对,直接用expand确实会导致所有task_2实例执行完毕后才会启动task_3,而要实现每个task_2完成后立刻触发对应的task_3,我们可以通过动态生成TaskGroup的方式来实现——虽然TaskGroup没有expand方法,但我们可以遍历task_1的返回结果,为每个task_num单独创建一个包含task_2和task_3的TaskGroup,让组内任务串行,组间任务并行。

下面是完整的实现代码:

from airflow import DAG
from airflow.decorators import task, task_group
from pendulum import datetime, now

@task
def task_1():
    # 这里返回动态数量的任务标识,示例返回[0,1,2,3,4]
    return list(range(5))

@task
def task_2(task_num):
    print(f"Executing task_2 for task_num: {task_num}")
    return task_num

@task
def task_3(task_num):
    print(f"Executing task_3 for task_num: {task_num}")
    return task_num

@task_group
def create_task_group(task_num):
    # 每个TaskGroup内,task_2完成后立即执行task_3
    t2 = task_2(task_num=task_num)
    t3 = task_3(task_num=task_num)
    t2 >> t3

with DAG(dag_id="my_dag", start_date=now(), schedule_interval=None) as dag:
    # 先执行task_1获取动态任务数量列表
    task_1_result = task_1()
    
    # 遍历task_1的返回结果,为每个task_num创建独立的TaskGroup
    for num in task_1_result:
        create_task_group(task_num=num)

代码解释:

  1. TaskGroup工厂函数:create_task_group是一个被@task_group装饰的函数,接收task_num作为参数,内部定义了该组内的task_2和task_3,并通过>>设置了它们的上下游依赖,确保task_2完成后立刻启动task_3。
  2. 动态生成分组:在DAG定义中,我们先执行task_1拿到返回的任务列表,然后遍历这个列表,为每个task_num实例化一个create_task_group——这样每个分组都是独立的,组内串行、组间并行,完全匹配你想要的任务图结构。
  3. 执行逻辑对比:和你之前的方案不同,这个实现不会等待所有task_2完成,只要某个分组内的task_2执行完毕,对应的task_3就会立即启动,完美解决了你的需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 14:07:39