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

Airflow动态映射任务未按预期顺序执行问题排查

Airflow Task Group Expand 执行顺序问题:实现(task1>>task2)x20的并行链

你在使用Airflow Task Group的expand功能实现动态任务映射时,遇到了执行顺序不符合预期的问题:所有task1实例先执行完毕,才会启动所有task2实例,但你期望每个数据项对应的task1执行完成后立即触发对应的task2,形成20条独立的task1>>task2串行链并行执行。

你的原始代码如下:

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

with DAG(
    dag_id="sandbox",
    start_date=datetime(2024, 1, 1, 9),
    schedule_interval=None,
    catchup=False,
    max_active_runs=1,
    default_args={
    }
) as dag:
    @task
    def return_a_list():
        return [i for i in range(20)]
    
    @task_group
    def tg(n):
        @task
        def task1(n):
            print("task1 received " + str(n) + " at " + str(datetime.now()))
            return n
        
        @task
        def task2(n):
            print("task2 received " + str(n) + " at " + str(datetime.now()))

        task2(task1(n))

    nbrs = return_a_list()
    tg.expand(n=nbrs)

问题原因与解决方法

你的expand用法本身是正确的,每个Task Group实例内的task1和task2已经通过task2(task1(n))建立了串行依赖,理论上应该实现task1>>task2的链式执行。出现所有task1先执行的情况,通常是以下两种原因:

  1. 并行度配置限制
    Airflow的默认并行任务数可能不足,导致所有task1排队执行,没有剩余资源启动task2。你可以在DAG定义中增加max_active_tasks参数,允许更多任务同时运行:
with DAG(
    dag_id="sandbox",
    start_date=datetime(2024, 1, 1, 9),
    schedule_interval=None,
    catchup=False,
    max_active_runs=1,
    default_args={},
    max_active_tasks=20  # 调整为合适的并行数
) as dag:
    # 后续代码不变
  1. 任务执行时间过短
    如果task1执行速度极快,视觉上会呈现出所有task1先完成的效果,但实际上是多个task1几乎同时执行完毕,随后task2批量启动。可以给task1添加短暂延迟,便于观察真实的执行顺序:
import time

@task
def task1(n):
    print(f"task1 received {n} at {datetime.now()}")
    time.sleep(1)  # 添加1秒延迟
    return n

代码优化建议

将task1和task2的定义移到Task Group外部,避免重复定义,代码结构更清晰:

from datetime import datetime
from airflow import DAG
from airflow.decorators import task, task_group
import time

with DAG(
    dag_id="sandbox",
    start_date=datetime(2024, 1, 1, 9),
    schedule_interval=None,
    catchup=False,
    max_active_runs=1,
    default_args={},
    max_active_tasks=20
) as dag:
    @task
    def return_a_list():
        return [i for i in range(20)]
    
    @task
    def task1(n):
        print(f"task1 received {n} at {datetime.now()}")
        time.sleep(1)
        return n
    
    @task
    def task2(n):
        print(f"task2 received {n} at {datetime.now()}")

    @task_group
    def tg(n):
        t1_instance = task1(n)
        task2(t1_instance)

    nbrs = return_a_list()
    tg.expand(n=nbrs)

结论

不需要改用循环创建任务,Airflow的expand是处理动态任务映射的标准方式。调整并行配置或添加延迟后,即可实现每个数据项对应的task1完成后立即执行task2的预期效果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:26:02