如何在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
相关产品推荐
相关产品推荐

