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

Airflow中能否同时配置DAG级与任务级on_failure_callback?

Airflow DAG级与任务级on_failure_callback的执行逻辑

核心结论

任务级的on_failure_callback会覆盖DAG级的回调,不会同时触发。具体执行逻辑分两种场景:

  • 当单个任务失败时:如果任务自身配置了on_failure_callback,只会执行该任务的回调;若任务未配置,则依次 fallback 到default_args中的回调、DAG级的回调。
  • 当DAG本身出现调度层面的失败(无任务执行就宣告失败):只会触发DAG级的on_failure_callback。

优先级规则

Airflow回调的优先级从高到低为:

  • 任务级单独配置的on_failure_callback
  • default_args中定义的on_failure_callback
  • DAG对象直接配置的on_failure_callback

结合你的示例说明

你的DAG定义中,default_args包含独立的失败回调,同时DAG本身配置了failure_callback_dag:

with DAG(
        default_args={"on_failure_callback": 自定义任务默认回调},
        dag_id = "test_dag",
        schedule_interval=None,
        start_date=datetime.datetime(2021, 1, 1),
        catchup=False,
        on_failure_callback=failure_callback_dag
    ) as dag:
    # 任务定义
  • 如果任务未单独配置回调,会优先执行default_args里的回调,而非DAG级的failure_callback_dag
  • 如果任务自己指定了on_failure_callback,则只会执行该任务的回调

实现两者都触发的方法

如果需要任务失败时同时执行任务级和DAG级的回调,不能依赖Airflow默认逻辑,需手动在任务级回调中调用DAG级回调:

def task_failure_callback(context):
    # 执行任务级失败操作
    handle_task_failure(context)
    # 调用DAG级失败回调
    failure_callback_dag(context)

# 任务配置示例
test_task = PythonOperator(
    task_id="test_task",
    python_callable=your_task_func,
    on_failure_callback=task_failure_callback
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 00:25:58