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

如何在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任务日志链接、conf DAG配置参数等),丰富自定义逻辑的内容

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

相关产品推荐
方舟 Agent Plan

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

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