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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 21:36:24