Airflow TaskGroup中映射任务执行顺序异常问题求助
问题根源
你遇到的核心问题是:Task Group结合Dynamic Mapping时,同一映射实例内的下游任务(add_42)未等待上游任务(print_num)完成就被触发,同时伴随任务状态异常报错。这并非depends_on_past或wait_for_downstream参数的问题,而是代码中任务依赖的绑定逻辑未被Airflow正确识别,或Airflow版本对Task Group动态映射的支持存在缺陷。
修复步骤
1. 修正任务依赖的绑定方式
原代码中add_42直接使用外部传入的my_num参数,而非依赖print_num的输出结果,这会导致Airflow无法准确识别同一映射实例内的上下游依赖关系。修改为通过任务输出传递数据,强制Airflow维护依赖:
@dag(dag_id='chore_task_group_stage3', catchup=False) def pipeline(): @task_group(group_id="channel_demo_tg") def tg1(my_num): @task() def print_num(num): return num @task() def add_42(num): return num + 42 # 用print_num的输出作为add_42的输入,而非直接传my_num num_output = print_num(my_num) add_42(num_output) tg1_object = tg1.expand(my_num=[19, 23, 42, 8, 7, 108]) pipeline()
2. 检查Airflow版本兼容性
若使用的Airflow版本低于2.4.0,Task Group与Dynamic Mapping的组合存在已知的依赖解析bug。MWAA环境中可通过查看环境配置确认版本,若版本过低,建议升级至2.4.0及以上版本。
3. 排查MWAA资源配置异常
日志中的报错Executor reports task instance finished (failed) although the task says it's queued通常与资源不足或Executor配置有关:
- 检查MWAA环境的Worker资源(CPU/内存)是否足够支撑并行任务
- 确认Executor类型(如CeleryExecutor)的并发数配置是否合理,避免任务被强制终止
4. 移除无效参数
depends_on_past=True和wait_for_downstream=True是用于跨DAG运行的依赖控制,对同一DAG内的动态映射任务依赖无效,可直接移除这些参数。
验证方法
修改代码后手动触发DAG,确认每个print_num的映射任务(对应map_index)完成后,对应的add_42任务才会启动,且无状态异常报错。
内容的提问来源于stack exchange,提问作者marcin2x4

