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

