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先执行的情况,通常是以下两种原因:
- 并行度配置限制
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: # 后续代码不变
- 任务执行时间过短
如果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
相关产品推荐
相关产品推荐

