Airflow 2.3如何标记任务实例为失败且不抛异常仍向下游传值
为什么原写法无效
你的代码无法修改任务状态核心是两个原因:
- 直接给
context['task_instance']赋值State.FAILED本质只是修改了上下文字典的本地键值,没有触发任务状态的持久化更新。Airflow在任务执行退出时,会根据任务是否抛出未捕获异常重新判定并覆写状态,你手动修改的字典值不会被采纳。 - Airflow默认逻辑下,只有任务抛出未捕获异常时才会被标记为
Failed,但这种场景下任务会直接中断,不会将返回值写入XCom,下游任务无法接收到数据,会直接触发跳过/失败逻辑。
正确实现方案
调用上下文里真实任务实例对象自带的set_state()方法主动更新状态即可,这个方法会直接将状态写入Airflow元数据库,不会被任务退出时的默认状态逻辑覆盖。同时因为你是在异常捕获块内正常return返回值,XCom会正常推送数据给下游,不影响下游执行。
修正后的可运行代码如下:
from airflow.utils.state import State from airflow.operators.python import get_current_context from datetime import datetime @dag( schedule_interval=None, start_date=datetime(2021, 1, 1) ) def myDag(): @task(retries=0) # 记得关自动重试,避免标记失败后任务被拉起重跑 def getData(): try: result = something_that_might_fail() return result except HorribleException as error: context = get_current_context() # 调用实例方法更新状态,直接持久化到元数据库 context['task_instance'].set_state(State.FAILED) return {"error": str(error)} @task() def transform(data_dict: dict): # 下游可自行判断返回类型,区分正常数据和错误返回做对应处理 if "error" in data_dict: print(f"上游getData任务执行出错,错误信息: {data_dict['error']}") for data in data_dict: print(data) # 其余业务逻辑 transform(getData()) run = myDag()
注意事项
- 捕获异常后不要重新
raise错误,否则任务会直接中断,XCom无法正常写入,下游拿不到返回数据。 - 该方式标记的
Failed状态会正常触发Airflow自带的失败告警、失败统计逻辑,完全满足监控需求。 - 如果集群开启了全局任务失败自动重试,一定要给
getData任务单独配置retries=0,避免任务标记失败后被调度器自动拉起重跑。
内容的提问来源于stack exchange,提问作者Flo
相关产品推荐
相关产品推荐

