Apache Airflow问题:任务组执行后do_third任务从排队转失败
问题分析与解决方案
问题原因
你遇到的do_third任务失败,核心原因是:
- 使用
task_group.expand()生成多个动态任务组实例后,返回的是多实例任务的XCom输出集合(一个XComArg列表)。直接将这个集合作为依赖传递给无参数的do_third时,Airflow会尝试把所有上游输出传入do_third,但do_third未定义接收参数,导致任务执行报错。 - 另外,代码中
log.info(f'zip_file_path: {tg_result}')是在DAG解析阶段执行的,此时tg_result是未运行的XComArg对象,这行代码无实际业务意义,属于冗余代码。
解决方案
有两种简单的修复方式,任选其一即可:
方式一:直接绑定动态任务组依赖
不需要将task_group_one.expand()的结果赋值给变量,直接将动态任务组实例作为do_third的上游依赖,Airflow会自动等待所有动态生成的任务组完成后再执行do_third:
# 替换原有的tg_result相关代码 task_group_one.expand(index=details) >> do_third()
方式二:修改do_third任务触发规则或接收参数
如果需要保留tg_result变量,可以通过以下两种方式适配:
- 给
do_third添加all_done触发规则,确保所有上游任务完成后执行,无需接收输入:@task(trigger_rule="all_done") def do_third(): log.info(f"do_third executed") - 让
do_third接收任意输入(即使不用),适配上游的多输出集合:@task def do_third(*args): log.info(f"do_third executed")
完整修正代码
以下是采用方式一的完整可运行代码:
from datetime import datetime from airflow.sdk import dag, task_group, task import logging as log @dag( description='Test DAG', schedule=None, catchup=False, start_date=datetime(2025, 6, 13) ) def test_dag(): @task def do_third(): log.info(f"do_third executed") @task def group_task_one(index): log.info(f'group_task_one executed {index}') return {"task1": "group_task_one"} @task def group_task_two(index): log.info(f'group_task_two executed {index}') return {"task2": "group_task_two"} @task_group(group_id="task_group_one") def task_group_one(index): t1 = group_task_one(index) t2 = group_task_two(index) t1 >> t2 return t2 @task def make_list(): return [1, 2] details = make_list() # 直接绑定动态任务组到do_third的上游 task_group_one.expand(index=details) >> do_third() test_dag()
内容的提问来源于stack exchange,提问作者Abhilash Mangalan
相关产品推荐
相关产品推荐

