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

Airflow中DAG级别on_failure_callback不生效,仅任务级别可行的问题咨询

Airflow DAG级on_failure_callback不触发的原因及解决办法

核心差异:触发时机与日志位置

你遇到的问题本质是对DAG级和任务级on_failure_callback的触发逻辑理解偏差:

  • 任务级回调:单个任务执行失败时立即触发,运行在任务执行进程中,日志会输出到任务实例的日志页面。
  • DAG级回调:仅当整个DAG Run的状态变为failed时触发,由Airflow Scheduler进程执行,日志会输出到Scheduler的日志文件中,而非任务日志页面。

你的代码本身没有问题,DAG级回调其实已经触发了,但因为你用print输出内容,而这些内容只出现在Scheduler日志里,你可能只查看了任务日志,所以误以为没触发。

验证方法

把回调函数中的print替换为标准日志记录,这样Scheduler日志会明确记录回调执行情况:

import sys
import logging
from datetime import datetime
from airflow import DAG
from airflow.operators.python_operator import PythonOperator

def on_failure(ctx):
    logging.error("DAG Run failed callback triggered!")
    logging.error(f"Context details: {ctx}")

def always_fails():
    sys.exit(1)

dag = DAG(
    dag_id='always_fails',
    description='dag that always fails',
    schedule_interval=None,
    catchup=False,
    start_date=datetime(2021,7,12),
    on_failure_callback=on_failure
)

task = PythonOperator(task_id='test-error-notifier', python_callable=always_fails, dag=dag)

运行后查看Airflow Scheduler的日志文件(通常位于$AIRFLOW_HOME/logs/scheduler目录下),就能看到对应的错误日志,证明回调已执行。

替代方案:全局任务失败回调

如果需要在任意任务失败时触发统一回调,更推荐使用default_args配置,而非DAG级回调:

default_args = {
    'on_failure_callback': on_failure
}

dag = DAG(
    dag_id='always_fails',
    description='dag that always fails',
    schedule_interval=None,
    catchup=False,
    start_date=datetime(2021,7,12),
    default_args=default_args
)

task = PythonOperator(task_id='test-error-notifier', python_callable=always_fails, dag=dag)

这种方式会自动将回调应用到DAG下的所有任务,且触发逻辑和任务级回调一致,日志直接出现在任务日志页面,更便于调试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 14:25:28