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

如何在Airflow任务失败时将异常信息捕获至变量中?

Airflow任务失败时捕获异常信息到变量的实现方案

核心思路

通过Airflow的on_failure_callback回调机制,在任务失败时自动捕获完整异常信息,并存入Airflow全局变量中,支持按任务ID区分存储。

具体实现步骤

1. 定义失败回调函数

编写回调函数,从任务上下文提取异常信息,存入指定变量:

from airflow.models import Variable
import traceback

def capture_failure_exception(context):
    # 获取完整异常追踪栈(对应你需要的高亮异常内容)
    full_exception = traceback.format_exc()
    # 从上下文获取当前失败任务的ID
    task_id = context['task_instance'].task_id
    # 将异常信息存入Airflow变量,变量名按任务ID区分
    Variable.set(f"failure_exception_{task_id}", full_exception)
    
    # 若仅需异常消息而非完整栈,可替换为:
    # exception_msg = str(context['exception'])
    # Variable.set(f"failure_exception_{task_id}", exception_msg)

2. 为任务绑定回调

可以选择将回调设置在DAG的default_args中,让所有任务自动继承;也可以单独为指定任务设置:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

# 示例任务逻辑,模拟任务B失败
def task_a():
    pass

def task_b():
    raise ValueError("任务B执行出错:数据库连接超时")

def task_c():
    pass

# 默认参数,全局绑定失败回调
default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'on_failure_callback': capture_failure_exception
}

with DAG('failure_exception_capture_dag', 
         default_args=default_args, 
         schedule_interval='@daily') as dag:
    
    task_a = PythonOperator(task_id='task_a', python_callable=task_a)
    task_b = PythonOperator(task_id='task_b', python_callable=task_b)
    task_c = PythonOperator(task_id='task_c', python_callable=task_c)
    
    task_a >> task_b >> task_c

关键说明

  • traceback.format_exc()会返回包含异常类型、消息、调用栈的完整信息,完全匹配你需要的高亮内容
  • 若仅需简洁的异常消息,使用str(context['exception'])即可
  • 通过context['task_instance'].task_id可以区分不同任务的异常,避免变量值被覆盖
  • Airflow Variables为全局存储,可通过Airflow UI的「Admin > Variables」页面查看或修改

内容的提问来源于stack exchange,提问作者Simran Suri

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 01:10:11