Apache Airflow 回调机制(Callbacks)完整指南:DAG 与 Task 状态事件驱动的监控告警
Airflow 的日志与监控体系(见 airflow-core/docs/administration-and-deployment/logging-monitoring)中,回调(Callbacks) 是响应 DAG/Task 状态变化的关键手段:你可以让某个 DAG 成功、某个 Task 失败或重试时,自动触发预先定义的函数(如发送告警、清理资源、更新外部系统)。本文以 callbacks.rst 为骨架,结合仓库源码与配置示例,系统讲解回调的三种定义层级、五种事件类型、Context 映射规则、Notifier 用法与 Deadline Alert(截止时间告警),帮助你构建一套可落地的监控告警方案。
回调的本质:把「状态变化」变成「可编程动作」
回调函数的用途非常直观:当某些 Task 失败时发告警、当整个 DAG 成功时触发后续通知。在 Airflow 中,回调可以在三个不同的位置定义,作用范围由细到粗:
- 在单个 Task 定义内部直接设置回调 —— 只作用于这一个 Task;
- 通过
default_args为 DAG 中每个 Task 统一设置回调; - 在 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_callback 与 on_failure_callback 同时支持 Dag 和 Task 两种层级,而 on_retry_callback、on_execute_callback、on_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 当前的状态:
- 常规失败时:选中最新失败的那个 Task;
- DAG Run 超时(timeout)时:选中最新开始但尚未结束的 Task;
- 出现死锁(deadlock) 时:选中本应继续运行、却因死锁无法运行的 Task;
- 成功时:选中最新成功完成的那个 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]——task2、task3没有覆盖它,但它们依然继承了default_args中的on_execute_callback; - 这里的
EmptyOperator来自airflow.providers.standard.operators.empty,其实现位于 empty.py(本质上是一个什么都不做、常用来组织依赖关系的算子)。
如果把 on_failure_callback 换成一次真实的告警集成,只需要在回调函数中调用你团队的消息通道 API 即可——context 中的 task_instance_key_str、run_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 级 Notifier(example_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)。
排查与最佳实践小结
结合仓库源码与上述文档内容,使用回调时有几个关键实践值得固化:
- 选择正确的定义层级:全局性的流程成功/失败通知放 DAG 级;单算子特有逻辑(如某关键步骤失败需要立刻告警)放 Task 级;团队统一默认行为用
default_args下发。 - 区分时间型与事件型告警:失败、重试、跳过属于事件型,用五类
on_*_callback;担心"卡死不报错、迟迟不结束"的运行要用 Deadline Alert 兜底,两者互补才能形成完整监控闭环。 - 告警逻辑保持轻量:回调在 Task 完成之后执行,其异常只写入 dag processor 日志(默认不显示在 UI,需查看
$AIRFLOW_HOME/logs/dag_processor/latest/dags-folder/.../DAG_FILE.py.log),容易漏看,因此回调内建议只做幂等、短小的通知动作。 - 不要依赖 Dag 级回调的 Task context:Dag 回调中选中的 Task 只反映局部状态,超时或死锁场景下仅代表"被选中"的一个实例,不能据此推断全局原因。
- 能用 Notifier 就不自己写:优先复用社区 Notifier(如
SlackWebhookNotifier),需要深度定制再参考 如何编写自定义 Notifier;普通自定义回调至少把 context 参数预留出来,方便后续扩展为发送带task_instance_key_str、run_id的富文本告警。
把回调当作"状态机上的钩子"来设计,再叠加 Deadline Alert 的时间维度,就能用很少的代码把 Airflow 的工作流运行状态无缝接入你的监控告警体系。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00