如何在Airflow出现异常时执行自定义函数并传递异常信息?
在Airflow中捕获异常并执行自定义逻辑(含异常信息传递)
这个需求完全可行,Airflow提供了多种机制捕获任务或DAG运行时的异常,触发自定义逻辑的同时还能直接传递异常相关信息。以下是几种常用实现方式:
1. 任务/全局级别的on_failure_callback
这是最直接的实现方式,可针对单个任务或全局所有任务配置失败回调,异常信息会通过回调函数的context参数传入。
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def custom_failure_handler(context): # 从上下文提取异常及任务信息 exception = context.get('exception') task_instance = context.get('task_instance') print(f"任务 {task_instance.task_id} 执行失败,异常详情: {str(exception)}") # 这里编写你的自定义逻辑:比如发送告警、写入异常日志、调用外部接口等 # 全局默认参数,所有任务失败都会触发该回调 default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), 'on_failure_callback': custom_failure_handler } with DAG('failure_handling_dag', default_args=default_args, schedule_interval='@daily') as dag: def simulate_failing_task(): # 模拟任务抛出异常 raise ValueError("这是一个测试用的任务异常") test_task = PythonOperator( task_id='simulate_failure', python_callable=simulate_failing_task, # 可单独为该任务配置专属回调,覆盖全局设置 # on_failure_callback=custom_failure_handler )
context字典中常用的异常相关字段:
exception: 任务抛出的原始异常对象task_instance: 任务实例对象,可获取任务ID、DAG ID、执行时间等dag_run: DAG运行实例对象,可获取DAG的运行状态、触发方式等
2. TaskFlow API风格的异常处理
如果使用TaskFlow装饰器写法,既可以在任务函数内部捕获异常,也可以通过on_failure_callback配置全局回调:
from airflow.decorators import dag, task from datetime import datetime def custom_failure_handler(context): exception_msg = str(context['exception']) print(f"TaskFlow任务执行失败,异常信息: {exception_msg}") # 执行自定义逻辑 @dag(start_date=datetime(2024, 1, 1), schedule_interval='@daily') def taskflow_failure_dag(): @task(on_failure_callback=custom_failure_handler) def failing_task(): raise TypeError("TaskFlow任务抛出的类型异常") failing_task() dag = taskflow_failure_dag()
3. 触发自定义Activity/外部任务
如果需要调用外部Activity(比如自定义服务、第三方任务流节点等),可在回调函数中直接发起调用,并将异常信息作为参数传递:
def trigger_external_activity(context): exception_details = str(context['exception']) task_id = context['task_instance'].task_id execution_date = context['execution_date'].isoformat() # 调用外部Activity示例:发送HTTP请求传递异常信息 import requests requests.post( "https://your-activity-endpoint.com/handle-failure", json={ "task_id": task_id, "exception": exception_details, "execution_time": execution_date } ) # 将该函数配置为on_failure_callback即可
注意事项
- 回调函数运行在Airflow Worker进程中,需确保依赖的第三方库(如
requests)已在Worker环境安装 - 若回调逻辑耗时较长,建议将其异步化(比如发送到消息队列,由其他服务异步处理),避免阻塞Worker进程
- 可通过
context获取更多上下文信息(如log_url任务日志链接、confDAG配置参数等),丰富自定义逻辑的内容
内容的提问来源于stack exchange,提问作者user1112259
相关产品推荐
相关产品推荐

