首页
/ Apache Airflow 回调机制(Callbacks)完整指南:DAG 与 Task 状态事件驱动的监控告警

Apache Airflow 回调机制(Callbacks)完整指南:DAG 与 Task 状态事件驱动的监控告警

2026-09-08 10:14:58作者:裘旻烁

Airflow 的日志与监控体系(见 airflow-core/docs/administration-and-deployment/logging-monitoring)中,回调(Callbacks) 是响应 DAG/Task 状态变化的关键手段:你可以让某个 DAG 成功、某个 Task 失败或重试时,自动触发预先定义的函数(如发送告警、清理资源、更新外部系统)。本文以 callbacks.rst 为骨架,结合仓库源码与配置示例,系统讲解回调的三种定义层级、五种事件类型、Context 映射规则、Notifier 用法与 Deadline Alert(截止时间告警),帮助你构建一套可落地的监控告警方案。

回调的本质:把「状态变化」变成「可编程动作」

回调函数的用途非常直观:当某些 Task 失败时发告警、当整个 DAG 成功时触发后续通知。在 Airflow 中,回调可以在三个不同的位置定义,作用范围由细到粗:

  1. 在单个 Task 定义内部直接设置回调 —— 只作用于这一个 Task;
  2. 通过 default_args 为 DAG 中每个 Task 统一设置回调;
  3. DAG 定义层面设置回调 —— 作用于整个 DAG Run 的生命周期事件(例如整条 DAG 的成功/失败)。

这三个层级可以在同一条 DAG 中混合使用,覆盖面相互叠加,例如任务级回调精确到单个算子,DAG 级回调则关心整条工作流的成败。这种分层设计让告警粒度可以从「全流程」精确到「单步骤」。

注意:回调函数只在 DAG 或 Task 因为 Worker 执行而发生状态变化时被调用。因此,通过命令行(CLI 文档)或用户界面(UI 文档)手动设置的状态变更,不会执行任何回调函数。例如你用 airflow dags backfill 或手动 mark 某个 Task 为 failed,都不会触发对应回调。

另一个重要约束是回调的执行位置与日志归属:回调函数在 Task 完成之后执行,回调函数内部若抛错,错误会出现在 dag processor 日志中而非 task 日志中。默认情况下 dag processor 日志不会展示在 UI 中,需要到文件系统查看:

$AIRFLOW_HOME/logs/dag_processor/latest/dags-folder/<the_path_for_your_dag>/DAG_FILE.py.log

Airflow 2.6.0 开始,回调还支持传入函数列表,即在同一个事件上挂多个回调函数,例如:

on_failure_callback=[callback_func_1, callback_func_2]

五种回调类型:触发事件与适用层级

Airflow 共定义了 5 种由特定事件触发的回调:

名称 描述 可用层级
on_success_callback 当 DAG Run 成功或 Task Instance 成功时调用 Dag 或 Task
on_failure_callback 当 DAG Run 失败或 Task Instance 失败时调用 Dag 或 Task
on_retry_callback 当 Task 进入重试(up for retry)时调用 Task
on_execute_callback 在 Task 开始执行之前调用 Task
on_skipped_callback 当 Task 正在运行时抛出了 AirflowSkipException 时调用。注意:如果 Task 因为 DAG 中前序分支决策或 trigger rule 导致从未被调度执行而被跳过,则不会调用该回调 Task

可以看到,on_success_callbackon_failure_callback 同时支持 Dag 和 Task 两种层级,而 on_retry_callbackon_execute_callbackon_skipped_callback 是针对 Task 执行细节的钩子。表里提到的各种状态语义可参考 Airflow 文档的 Dag Run 状态Task Instance 状态 章节。

在源码层面,这些回调标志会被序列化并在不同进程间传递。例如 baseoperator.py 中定义了一组布尔字段记录每个算子是否携带回调:

has_on_execute_callback: bool = False
has_on_failure_callback: bool = False
has_on_retry_callback: bool = False
has_on_success_callback: bool = False
has_on_skipped_callback: bool = False

这些标记帮助调度器判断某个 Task 是否需要进入带回调的执行路径(参见 callback.py 中关于执行器回调 workload 的设计)。

Context 映射:每个回调都能拿到哪些运行时信息

每次回调被触发时,Airflow 都会向回调函数传入一个 context 映射(Context Mapping),其中包含关于 Task Instance 的运行时信息。context 中可用变量的完整清单见 模板与变量参考,其类型定义位于 context.py。从源码看,Context 是一个 TypedDict(见 task-sdk/src/airflow/sdk/definitions/context.py 中的 class Context(TypedDict, total=False)),常见字段包括:

  • dag_run:当前 DAG Run 对象;
  • run_id:本次运行的 ID;
  • task_instance_key_str:Task Instance 的字符串标识(常用于日志和告警文案);
  • 以及 Task Instance、DAG、逻辑日期(logical date)、conf 等大量运行时信息。

Dag 级回调的 Task 选择规则

由于 context 本质上描述的是某个 Task Instance 的执行情况,即便你写的是 Dag 级回调,context 中仍会包含 Task Instance 变量——只不过具体选中哪一个 Task 取决于 DAG Run 当前的状态:

  1. 常规失败时:选中最新失败的那个 Task;
  2. DAG Run 超时(timeout)时:选中最新开始但尚未结束的 Task;
  3. 出现死锁(deadlock) 时:选中本应继续运行、却因死锁无法运行的 Task;
  4. 成功时:选中最新成功完成的那个 Task。

因此原文档特别建议:不要在 Dag 级回调中依赖 Task Instance 变量来做业务判断(仅适合人工分析场景),因为它只反映 DAG 状态的局部信息。例如一次超时可能是多个 Task 停滞共同造成的,但最终只会有一个 Task 被选中放入 context。此外,在 Airflow 3.2.0 之前,上述状态关联规则并不存在——当时传入 Dag 回调的 Task 与 DAG 状态无关,只是按字典序选出的 DAG 中"最新"的那个 Task。升级后如果你的 Dag 回调逻辑恰好依赖旧行为,需要格外留意这一变化。

实战一:使用自定义回调方法

原文档给出一个完整的可运行示例:task1 失败时调用 task_failure_alert,DAG 整体成功时调用 dag_success_alert,而在每个 Task 开始执行前都会先调用 task_execute_callback。注意它同时演示了三种定义层级与回调函数列表的用法:

from airflow.sdk import DAG
from airflow.providers.standard.operators.empty import EmptyOperator


def task_execute_callback(context):
    print(f"Task has begun execution, task_instance_key_str: {context['task_instance_key_str']}")


def task_failure_alert(context):
    print(f"Task has failed, task_instance_key_str: {context['task_instance_key_str']}")


def dag_success_alert(context):
    print(f"Dag has succeeded, run_id: {context['run_id']}")


with DAG(
    dag_id="example_callback",
    on_success_callback=dag_success_alert,
    default_args={"on_execute_callback": task_execute_callback},
):
    task1 = EmptyOperator(task_id="task1", on_failure_callback=[task_failure_alert])
    task2 = EmptyOperator(task_id="task2")
    task3 = EmptyOperator(task_id="task3")
    task1 >> task2 >> task3

几点值得注意的写法:

  • Dag 级回调写在 DAG(...) 构造参数里:on_success_callback=dag_success_alert
  • 全任务默认回调通过 default_args 下发:default_args={"on_execute_callback": task_execute_callback} 会让 DAG 内所有 Task 在开始执行前都调用该函数;
  • 任务级回调写在算子参数里,且演示了 Airflow 2.6+ 的列表语法:on_failure_callback=[task_failure_alert]——task2task3 没有覆盖它,但它们依然继承了 default_args 中的 on_execute_callback
  • 这里的 EmptyOperator 来自 airflow.providers.standard.operators.empty,其实现位于 empty.py(本质上是一个什么都不做、常用来组织依赖关系的算子)。

如果把 on_failure_callback 换成一次真实的告警集成,只需要在回调函数中调用你团队的消息通道 API 即可——context 中的 task_instance_key_strrun_id 等字段能直接拼进告警正文。

实战二:使用 Notifier 作为回调

除了普通 Python 函数,你还可以把 Notifier 作为参数传给 on_*_callbacks。Notifier 是 Airflow 中对"通知动作"的封装对象,通常在 on_success_callback / on_failure_callback 中使用,按 Task 或 DAG Run 的状态发送通知。它的好处是开箱即用——社区为 Slack、Email、Microsoft Teams 等渠道维护了大量现成 Notifier,无需手写 HTTP 调用。

以下是原文档给出的自定义 Notifier 用法示例:

from airflow.sdk import DAG
from airflow.providers.standard.operators.bash import BashOperator

from myprovider.notifier import MyNotifier

with DAG(
    dag_id="example_notifier",
    on_success_callback=MyNotifier(message="Success!"),
    on_failure_callback=MyNotifier(message="Failure!"),
):
    task = BashOperator(
        task_id="example_task",
        bash_command="exit 1",
        on_success_callback=MyNotifier(message="Task Succeeded!"),
    )

注意这里同时出现了 Dag 级 Notifier(DAG 成功/失败时通知)和 Task 级 Notifierexample_task 自身成功时通知),说明 Notifier 与普通函数在三种定义层级上的用法完全一致。由于 bash_command="exit 1" 会让该 Task 必然失败,实际运行时触发的是 Dag 的 on_failure_callback

关于社区维护的 Notifier 清单,以及如何编写自定义 Notifier(实现通知逻辑与消息模板),参见 通知(Notifications)how-to

进阶:Deadline Alert 回调(基于时间的超时告警)

除了上面五类面向 DAG/Task 生命周期事件的回调,Airflow 还支持 Deadline Alert(截止时间告警)回调:当某次 DAG Run 超过设定时间阈值仍未结束时触发。它与传统回调的关键区别在于:

  • 传统回调绑定"状态变化"事件(成功/失败/重试/跳过);
  • Deadline Alert 绑定"时间流逝"事件(运行超时/即将超时)。

Deadline Alert 在 DAG 上通过 deadline 参数配置,并使用两种回调载体:

  • airflow.sdk.AsyncCallback —— 异步执行,运行在 Triggerer 进程中;
  • airflow.sdk.SyncCallback —— 同步执行,运行在 Executor 中。

一个典型的 Deadline Alert 由三要素组成:参考点(reference)(从何时开始计时,如 DAG Run 入队时刻)、时间间隔(interval)(在参考点前后偏移多少触发)与回调(callback)(超时后执行的动作)。其计算方式可用下面的时间线表达:

[Reference] ------ [Interval] ------> [Deadline]
    ^                                     ^
    |                                     |
 Start time                          Trigger point

示例如下——如果 DAG 在被排入队列 15 分钟后仍未完成,则通过 Slack 发送消息:

from datetime import datetime, timedelta
from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference
from airflow.providers.slack.notifications.slack_webhook import SlackWebhookNotifier
from airflow.providers.standard.operators.empty import EmptyOperator

with DAG(
    dag_id="deadline_alert_example",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_QUEUED_AT,
        interval=timedelta(minutes=15),
        callback=AsyncCallback(
            SlackWebhookNotifier,
            kwargs={
                "text": "🚨 Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}"
            },
        ),
    ),
):
    EmptyOperator(task_id="example_task")

DeadlineReference 提供了多种内置参考点,例如:DAGRUN_QUEUED_AT(从入队时刻计时,常用于监控资源不足导致的排队)、DAGRUN_LOGICAL_DATE(按计划开始时刻计时,确保排定的 DAG 在下一次调度前跑完)、FIXED_DATETIME(固定截止时刻,配合负 interval 可做到"提前 N 分钟提醒")、AVERAGE_RUNTIME(依据历史成功运行的平均耗时推算,历史数据不足时不创建 deadline)。Deadline Alert 的 interval 支持正负 timedelta,因此既能在超时后告警,也能在截止点前预告。

一个 DAG 还可以传入 DeadlineAlert 列表实现分级告警(例如 30 分钟先 Slack 提醒、60 分钟再升级调用 on-call);与 Deadline Alert 相关的完整示例、内置参考点对比、自定义参考点(通过插件注册 deadline_references)、自定义 Sync/Async 回调及 context 传递规则,详见 Deadline Alerts 专题文档;从 SLA 迁移到 Deadline 的指引参见 SLA 迁移指南

说明:Deadline Alert 在 Airflow 3.1 引入且标记为实验性,接口可能随用户反馈调整;AsyncCallback 的导入路径在 Airflow 3.2 起从 airflow.sdk.definitions.deadline 变更为 airflow.sdk(相关类定义见 task-sdk/src/airflow/sdk/definitions/deadline.py)。

排查与最佳实践小结

结合仓库源码与上述文档内容,使用回调时有几个关键实践值得固化:

  1. 选择正确的定义层级:全局性的流程成功/失败通知放 DAG 级;单算子特有逻辑(如某关键步骤失败需要立刻告警)放 Task 级;团队统一默认行为用 default_args 下发。
  2. 区分时间型与事件型告警:失败、重试、跳过属于事件型,用五类 on_*_callback;担心"卡死不报错、迟迟不结束"的运行要用 Deadline Alert 兜底,两者互补才能形成完整监控闭环。
  3. 告警逻辑保持轻量:回调在 Task 完成之后执行,其异常只写入 dag processor 日志(默认不显示在 UI,需查看 $AIRFLOW_HOME/logs/dag_processor/latest/dags-folder/.../DAG_FILE.py.log),容易漏看,因此回调内建议只做幂等、短小的通知动作。
  4. 不要依赖 Dag 级回调的 Task context:Dag 回调中选中的 Task 只反映局部状态,超时或死锁场景下仅代表"被选中"的一个实例,不能据此推断全局原因。
  5. 能用 Notifier 就不自己写:优先复用社区 Notifier(如 SlackWebhookNotifier),需要深度定制再参考 如何编写自定义 Notifier;普通自定义回调至少把 context 参数预留出来,方便后续扩展为发送带 task_instance_key_strrun_id 的富文本告警。

把回调当作"状态机上的钩子"来设计,再叠加 Deadline Alert 的时间维度,就能用很少的代码把 Airflow 的工作流运行状态无缝接入你的监控告警体系。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
899
5.82 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
920
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.8 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
532
596
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.02 K
521
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.36 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
392