Apache Airflow 从 SLA 迁移到 Deadline Alerts 完整实战指南:范式对比、迁移路径与源码级原理解析
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
SLA(Service Level Agreement)与 Deadline Alerts(截止时间告警)都是为 Dag Run 设置"最晚完成时间"并对超时做出响应的机制,但二者在触发时机、检查方式与底层实现上截然不同。本指南以 Airflow 官方迁移文档(airflow-core/docs/howto/sla-to-deadlines.rst)为核心骨架,结合 Deadline Alerts 在 Task SDK 与核心调度器中的源码实现,帮助你彻底理解两种范式的差异,掌握最直接的迁移路径,并学会为你的使用场景挑选合适的 Deadline 基准点(Reference)与间隔(Interval)。
读完本文,你将能够:把基于sla/sla_miss_callback的旧式 DAG 平滑迁移为基于DeadlineAlert的新式告警 DAG;理解 Deadline 从"创建 → 计算 → 检查 → 回调"的完整生命周期;并且知道如何利用内置的四种 Reference(队列时间、逻辑日期、固定时间、历史平均运行时长)以及自定义 Reference 与回调来构建贴合业务的分级告警体系。
两种截然不同的范式
虽然SLA与Deadline Alerts的目标非常相似——都是"规定一个最晚时间点,超时即告警"——但它们采用的是两条完全不同的技术路线。理解这两条路线的差异,是做出迁移决策的前提。
SLA:Dag Run 结束后才检查
SLA 的工作方式如下:
- 当 Dag Run结束(finishes)时,检查当前时间;
- 如果当前时间大于
(logical_date + sla),则执行sla_miss_callback; - 如果 Dag Run 永远没有结束,SLA 永远不会被检查。
也就是说,SLA 是一种"事后判定"机制:它把检查动作挂在了 Dag Run 的结束事件上,靠的是任务完成后的一次性判断,而不是持续的轮询。这意味着两个天然的局限:
- 运行中的 Dag 无论超时多久,在它结束之前你都收不到任何提醒;
- 一个"卡死"(既不失败也不结束)的 Dag 会完全绕过 SLA 检查。
Deadline Alerts:Dag Run 开始时即计算,调度器周期性检查
Deadline Alerts 的工作方式则是完全不同的"事前 + 持续监控"范式:
- 当 Dag Run开始(starts)时,立即计算并存储截止时间:
DeadlineReference 的基准时间 + interval; - 调度器循环(scheduler loop)随后周期性检查(默认每 5 秒,由
scheduler_heartbeat_sec配置项控制)这些时间点是否已经过去; - 一旦发现过期,立即执行
callback(**kwargs)。
从源码层面看,这一过程对应着核心调度器中的一段独立逻辑。在 scheduler_job_runner.py 中,调度器每个心跳周期都会执行一次查询:
deadline_query = ( select(Deadline) .where(Deadline.deadline_time < datetime.now(timezone.utc)) .where(~Deadline.missed) .options(selectinload(Deadline.callback), selectinload(Deadline.dagrun)) ) for deadline in session.scalars( with_row_locks( deadline_query, of=Deadline, session=session, skip_locked=True, key_share=False, ) ): deadline.handle_miss(session)这段代码中值得注意的实现细节是with_row_locks(..., skip_locked=True):在启用 HA(高可用)多调度器副本时,通过行级锁FOR UPDATE SKIP LOCKED保证同一行 Deadline 只会被一个调度器处理,避免重复创建回调。这正是 Deadline 范式能够在"运行中"就感知超时的底层保障。
在 config.yml 中可以看到该配置项的定义:
| 配置项 | 所属 Section | 类型 | 默认值 | 说明 |
|---|---|---|---|---|
scheduler_heartbeat_sec | [scheduler] | integer | 5 | 调度器尝试触发新任务 / 执行调度循环的频率(秒) |
也就是说,Deadline 触发的最大延迟约等于一个心跳周期:一旦deadline_time已过,最迟在下一个scheduler_heartbeat_sec(默认 5 秒)内回调就会被执行,不需要等待 Dag 结束。
最直接的迁移路径
官方文档给出的最直接迁移路径是使用DeadlineReference.DAGRUN_LOGICAL_DATE基准点,它最接近 SLA 的logical_date + sla语义。但必须清醒地认识到其中的重大行为差异:
Deadline 的回调会在计算出的过期时间到达后"立即"(在
scheduler_heartbeat_sec之内)执行,而不是等待 Dag 先结束。
换句话说,同样的"1 小时超时"配置,SLA 会在 Dag 结束后才告诉你"这次跑超了",而 Deadline 会在运行到第 60 分钟的那一刻就发出告警。如果团队依赖的是"事后统计",迁移后告警的到达时机将显著提前,这通常是期望中的改进,但也需要提前与告警接收方对齐预期。
等价示例 DAG:1 小时 SLA vs 1 小时 Deadline
下面先给出一个使用 1 小时 SLA 的 Dag,再给出一个功能等价的、使用 Deadline Alerts 的 Dag,两者可以并排对照。
SLA 版本
with DAG( "minimal_sla_example", default_args={"sla": timedelta(hours=1)}, sla_miss_callback=SlackWebhookNotifier( text="SLA missed for {{ dag_run.dag_id }}", ), ): BashOperator(task_id="long_task", bash_command="sleep 3600")Deadline Alerts 版本
with DAG( "minimal_deadline_example", deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_LOGICAL_DATE, interval=timedelta(hours=1), callback=AsyncCallback( SlackWebhookNotifier, kwargs={ "text": "Deadline missed for {{ dag_run.dag_id }}", }, ), ), ): BashOperator(task_id="long_task", bash_command="sleep 3600")两个版本在语义上的对照关系如下:
| 概念 | SLA 版本 | Deadline 版本 |
|---|---|---|
| 超时阈值 | default_args={"sla": timedelta(hours=1)} | interval=timedelta(hours=1) |
| 基准点 | 隐含使用logical_date | 显式声明reference=DeadlineReference.DAGRUN_LOGICAL_DATE |
| 回调 | sla_miss_callback= | callback=AsyncCallback(...)(包一层 Callback 对象) |
| 回调入参 | 直接传text= | 通过kwargs={...}传给 Callback |
| 检查时机 | Dag Run 结束后一次性检查 | Dag Run 运行期间每scheduler_heartbeat_sec(默认 5 秒)轮询 |
需要注意:Deadline 版本的回调被包在了AsyncCallback中,这是 Deadline 回调机制的统一接口——无论是内置 Notifier、自定义同步函数还是自定义异步函数,都必须以AsyncCallback或SyncCallback的形式传入(见下文"回调"一节)。
深入 Deadline Alerts:从配置到源码
在动手迁移之前,先建立对 Deadline Alerts 机制的完整认知。该特性在 Airflow 3.1 中引入,目前标记为experimental(实验性),未来版本可能根据用户反馈调整,官方文档在 deadline-alerts.rst 顶部有明确的 warning 提示。
Deadline 的计算模型
创建一条 Deadline Alert 需要三个必填参数,它们共同决定"何时算超时":
- Reference(基准点):从什么时候开始计时;
- Interval(间隔):在基准点之前或之后多远触发告警(可以是
timedelta,也可以是VariableInterval这样的动态间隔); - Callback(回调):一个 Callback 对象,包含指向可调用对象的路径以及(可选的)kwargs,超时后执行。
Deadline 的计算公式可以表示为:
[Reference] ------ [Interval] ------> [Deadline] ^ ^ | | Start time Trigger point即:deadline_time = reference 返回的时间 + interval。
下面是一个完整示例:如果 Dag 在被queued(入队)后 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")该例的时间线示意:
|------|-----------|---------|-----------|--------| Scheduled Queued Started Deadline 00:00 00:03 00:05 00:18注意:这里的AsyncCallback导入路径在 Airflow 3.2 中从airflow.sdk.definitions.deadline变更为airflow.sdk,本文示例统一使用新路径。
存储与生命周期:Deadline、DeadlineAlert 两张表
从源码结构看,Deadline 机制在数据库中对应两张核心表:
deadline_alert表(对应 models/deadline_alert.py 中的DeadlineAlert模型):保存定义—— 即 DAG 作者在DAG(deadline=...)中声明的告警配置。字段包括name、description、reference(JSON)、interval(JSON)与callback_def(JSON),并通过外键关联到serialized_dag。deadline表(对应 models/deadline.py 中的Deadline模型):保存实例—— 即每个 Dag Run 计算出的具体到期时间点。关键字段为deadline_time(过期时间)、missed(是否已被标记为错过)、callback_id(外键到callback表),并通过dagrun_id关联到具体的 Dag Run。表上还建有deadline_missed_deadline_time_idx索引,服务于调度器的过期扫描查询。
生命周期中的两个关键动作:
- 创建/清理:Dag Run 结束后,dagrun.py 会调用
Deadline.prune_deadlines(session=session, conditions={DagRun.id: self.id})。prune_deadlines(见 models/deadline.py)会删除"在 deadline 之前正常结束"的 Deadline 记录,并上报deadline_alerts.deadline_not_missed指标——注意它不会触碰已被标记missed的记录,那些回调的所有权在调度器手中。 - 触发:上文已经看到,调度器每个心跳周期扫描
deadline_time < now AND NOT missed的记录并调用handle_miss。handle_miss(见 models/deadline.py)会把TriggererCallback或ExecutorCallback加入队列,注入一个简化版的 Airflow context(包含dag_run与deadline信息),将记录标记为missed,并上报deadline_alerts.deadline_missed指标。
内置 Reference 全解
Airflow 提供四个开箱即用的内置基准点(task-sdk 定义,对应核心侧的序列化实现见 serialization/definitions/deadline.py):
DeadlineReference.DAGRUN_QUEUED_AT
从 Dag Run入队(queued)的时刻开始计时。适合监控资源受限、任务迟迟无法拿到执行 slot 的场景。在核心侧由DagRunQueuedAtDeadline实现,它从 DagRun 表中读取queued_at列(required_kwargs = {"dag_id", "run_id"})。仓库自带的示例 DAG example_deadline_alert.py 使用的正是这一基准点:
with DAG( dag_id="example_deadline_alert", ... deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(seconds=30), callback=AsyncCallback(notify_deadline_missed), name="example_deadline", ), ) as dag: @task def hello_deadline(): time.sleep(60) # 故意睡过 30 秒的 deadline,触发告警DeadlineReference.DAGRUN_LOGICAL_DATE
引用 Dag Run计划开始(scheduled to start)的时间,即迁移文档中推荐的、与 SLA 语义最接近的基准点。例如设置interval=timedelta(minutes=15),则无论 Dag 实际何时开始(甚至从未开始),只要在计划开始时间后 15 分钟尚未完成就会触发告警。适用于确保定时 DAG 在下一轮调度前完成。
DeadlineReference.FIXED_DATETIME
指定一个固定的时间点。适用于业务上存在硬性完成时间要求的场景(如"每天 10:00 前必须出报表")。其核心实现FixedDatetimeDeadline直接返回构造时传入的 datetime,不依赖数据库查询。
DeadlineReference.AVERAGE_RUNTIME
基于历史成功运行的耗时平均值动态计算 deadline。它分析历史执行数据来预测当前运行应该在何时完成:deadline = 当前时间 + 平均运行时长 + interval。如果历史数据不足,则不创建 deadline(也就不会误报)。
参数说明:
max_runs(int,可选):纳入统计的最近成功运行次数上限,默认 10;min_runs(int,可选):计算平均值所需的最少成功运行次数,默认与max_runs相同。
# 使用默认设置(分析最近 10 次运行,且要求至少 10 次) DeadlineReference.AVERAGE_RUNTIME() # 分析最近 20 次运行,但只要有 5 次即可计算 DeadlineReference.AVERAGE_RUNTIME(max_runs=20, min_runs=5) # 严格模式:必须恰好有 15 次运行才计算 DeadlineReference.AVERAGE_RUNTIME(max_runs=15, min_runs=15)从源码实现看,AverageRuntimeDeadline._evaluate_with是一个值得细读的示例(见 models/deadline.py 与序列化版本 serialization/definitions/deadline.py):
- 它按数据库方言生成时长表达式:PostgreSQL 使用
EXTRACT(EPOCH FROM end_date - start_date),MySQL 使用TIMESTAMPDIFF(SECOND, start_date, end_date),SQLite 使用julianday差值换算秒数; - 查询只筛选成功(SUCCESS)的 Dag Run 且
start_date、end_date均非空,按logical_date倒序取最近max_runs条——官方注释明确指出:失败快速退出或长时间挂起后失败的运行会扭曲平均值,导致 deadline 过短(误报)或过长(真慢也触发不了),因此必须排除; - 计算使用
Decimal高精度求和再转 float,兼容 MySQL 的Decimal类型; - 若成功运行数不足
min_runs,返回None(不创建 deadline)。
使用平均运行时的示例(历史平均 30 分钟,间隔 30 分钟):
with DAG( dag_id="average_runtime_deadline", deadline=DeadlineAlert( reference=DeadlineReference.AVERAGE_RUNTIME(max_runs=15, min_runs=5), interval=timedelta(minutes=30), # 超过平均耗时 30 分钟即告警 callback=AsyncCallback( SlackWebhookNotifier, kwargs={"text": "Dag {{ dag_run.dag_id }} is running longer than expected!"}, ), ), ): EmptyOperator(task_id="data_processing")对应时间线:
|------|----------|--------------|--------------|--------| Queued Start | Deadline 09:00 09:05 09:35 10:05 | | | |--- Average --|-- Interval --| (30 min) (30 min)使用固定时间的示例(负间隔实现"提前告警"):
tomorrow_at_ten = datetime.combine(datetime.now().date() + timedelta(days=1), time(10, 0)) with DAG( dag_id="fixed_deadline_alert", deadline=DeadlineAlert( reference=DeadlineReference.FIXED_DATETIME(tomorrow_at_ten), interval=timedelta(minutes=-30), # 在基准点前 30 分钟告警 callback=AsyncCallback( SlackWebhookNotifier, kwargs={ "text": "Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}" }, ), ), ): EmptyOperator(task_id="example_task")时间线示意(注意 interval 为负值,Deadline 位于 Reference 之前):
|------|----------|---------|------------|--------| Queued Start Deadline Reference 09:15 09:17 09:30 10:00回调:Async 与 Sync 两种执行路径
超时后执行的回调必须封装为AsyncCallback或SyncCallback之一(sdk 定义)。二者的区别在于执行者不同:
AsyncCallback:回调在Triggerer(触发器)中运行,适合与内置 Notifier(如SlackWebhookNotifier)配合;SyncCallback:回调被发送给executor(执行器),像最高优先级的普通任务一样运行,在 Airflow 3.2 中引入。
下面两个示例实现完全相同的功能——Dag 入队 30 分钟内未完成则发 Slack 告警,区别只在回调类型:
# 异步版本:回调运行在 Triggerer with DAG( dag_id="slack_deadline_alert_async", deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(minutes=30), callback=AsyncCallback( SlackWebhookNotifier, kwargs={ "text": "Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}" }, ), ), ): EmptyOperator(task_id="example_task") # 同步版本:回调运行在 executor with DAG( dag_id="slack_deadline_alert_sync", deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(minutes=30), callback=SyncCallback( SlackWebhookNotifier, kwargs={ "text": "Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }}" }, ), ), ): EmptyOperator(task_id="example_task")自定义回调的注意点:
- 自定义 callable 若要接收
kwargs,直接在Callback中传入即可; - 异步回调必须位于 Triggerer 的系统路径上。简单做法是把 callable 作为顶层函数放在 plugins 目录(如
$AIRFLOW_HOME/plugins/deadline_callbacks.py)的新文件中;嵌套函数暂不支持;新增或修改回调后需要重启 Triggerer以重新加载文件; - 同步回调必须能被执行它的 worker 导入;
- 超时触发时,Airflow 会自动向回调注入一个
contextkwarg,包含 Dag Run 与 deadline 的信息(通过kwargs["context"]访问,或声明一个名为context的参数接收)。不需要 context 的回调可以省略它——Airflow 只会传入 callable 能接受的参数。context是保留关键字,不能出现在Callback的kwargs中,否则在 DAG 解析期就会抛出ValueError。
自定义同步回调示例(第 1 步:放入 plugins 文件夹,例如$AIRFLOW_HOME/plugins/deadline_callbacks.py):
def custom_sync_callback(**kwargs): """Handle deadline violation with custom logic.""" context = kwargs.get("context", {}) print(f"Deadline exceeded for Dag {context.get('dag_run', {}).get('dag_id')}!") print(f"Context: {context}") print(f"Alert type: {kwargs.get('alert_type')}") # Additional custom handling here第 2 步:在 Dag 文件中引用:
from datetime import timedelta from deadline_callbacks import custom_sync_callback from airflow.providers.standard.operators.empty import EmptyOperator from airflow.sdk import DAG, DeadlineAlert, DeadlineReference, SyncCallback with DAG( dag_id="custom_sync_deadline_alert", deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(minutes=15), callback=SyncCallback( custom_sync_callback, kwargs={"alert_type": "time_exceeded"}, ), ), ): EmptyOperator(task_id="example_task")自定义异步回调示例(第 1 步:放入 plugins 文件夹):
async def custom_async_callback(**kwargs): """Handle deadline violation with custom logic.""" context = kwargs.get("context", {}) print(f"Deadline exceeded for Dag {context.get('dag_run', {}).get('dag_id')}!") print(f"Context: {context}") print(f"Alert type: {kwargs.get('alert_type')}") # Additional custom handling here第 2 步:重启 Triggerer;第 3 步:在 Dag 文件中引用:
from datetime import timedelta from deadline_callbacks import custom_async_callback from airflow.providers.standard.operators.empty import EmptyOperator from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference with DAG( dag_id="custom_deadline_alert", deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(minutes=15), callback=AsyncCallback( custom_async_callback, kwargs={"alert_type": "time_exceeded"}, ), ), ): EmptyOperator(task_id="example_task")进阶提示:
SyncCallback支持可选的executor参数,用于把回调路由到指定 executor;不指定则使用默认 executor:SyncCallback( my_callback, kwargs={"msg": "deadline missed"}, executor="celery_executor", )AsyncCallback支持可选的queue参数,把由此产生的 trigger 分配到特定 trigger queue;不指定则运行在任意未加--queues限制的 triggerer 上:AsyncCallback( my_callback, kwargs={"msg": "deadline missed"}, queue="alerts", )
模板化与简化 Context
Deadline 回调当前收到的是一份简化版的 Airflow context,并且 Airflow不会对 Callback 的 kwargs 做 Jinja 模板渲染。但内置 Notifier 在执行时本身会基于收到的 context 做模板渲染,因此只要被模板化的变量包含在简化 context 中,模板语法在 Notifier 场景下依然可用。简化 context 目前包含:Deadline Alert 的 ID 与计算出的 deadline 时间,以及 Dag Run 的GETREST API 响应中包含的数据(因此示例中{{ dag_run.dag_id }}、{{ deadline.deadline_time }}均可正常渲染)。更完整的 context 与模板支持会在未来版本中增强。
如何为你的场景选择合适的 Deadline
迁移或新建 Deadline Alert 时,关键决策点是选哪个 Reference。下表汇总了四个内置基准点的适用场景:
| Reference | 基准时间来源 | 典型场景 | 常用 interval |
|---|---|---|---|
DAGRUN_QUEUED_AT | DagRun 表的queued_at | 监控资源挤压、入队后迟迟不启动 | 正数(入队后 X 分钟) |
DAGRUN_LOGICAL_DATE | DagRun 表的logical_date | 定时任务需在下一轮调度前完成(最接近 SLA 语义) | 正数(计划时间后 X 分钟) |
FIXED_DATETIME | 固定的 datetime | 硬性完成时间(报表截止、会议开始) | 可为负数(提前告警) |
AVERAGE_RUNTIME | 历史成功运行耗时平均值 | 运行时长波动大、需要"自适应"阈值 | 正数(平均耗时后 X 分钟) |
在"Deadline 计算"一节中,官方文档还给出了两个精炼的范例,展示了正负 interval 的灵活组合:
# 场景 A:会议开始前 2 小时提醒(FIXED_DATETIME + 负 interval) next_meeting = datetime(2025, 6, 26, 9, 30) DeadlineAlert( reference=DeadlineReference.FIXED_DATETIME(next_meeting), interval=timedelta(hours=-2), callback=notify_team, ) # 场景 B:计划 1 小时内未完成即告警(DAGRUN_LOGICAL_DATE + 正 interval) DeadlineAlert( reference=DeadlineReference.DAGRUN_LOGICAL_DATE, interval=timedelta(hours=1), callback=notify_team, )场景 B 的含义是:如果 Dag 计划每天 0 点运行,那么只要 1:00 还没完成就会触发告警——这正是迁移文档中推荐的 SLA 等价替代。把不同 Reference 与正负 Interval 自由组合,几乎可以覆盖所有运维告警需求。
自定义 Reference:扩展 Deadline 到你的业务数据
内置 Reference 覆盖了绝大多数通用场景,但当业务要求"以日历上的截止时间""以外部系统返回的时间戳"等作为基准时,就需要自定义 Reference。
创建并注册自定义 Reference
要创建自定义 Reference,需要三步:继承BaseDeadlineReference、加上@deadline_reference装饰器、实现_evaluate_with()方法;随后把类注册到插件的deadline_references列表中(与自定义 Timetable 的注册方式相同),调度器才能在反序列化 Dag 时解析到它。
把下面的类与插件放入 plugins 目录(例如$AIRFLOW_HOME/plugins/deadline_references.py):
from sqlalchemy.orm import Session from airflow.plugins_manager import AirflowPlugin from airflow.sdk import BaseDeadlineReference, DeadlineReference, deadline_reference from airflow.sdk.timezone import datetime # 默认在 Dag Run 创建时执行 evaluate_with @deadline_reference() class MyCustomDecoratedReference(BaseDeadlineReference): """A custom reference evaluated when Dag runs are created.""" def _evaluate_with(self, *, session: Session, **kwargs) -> datetime: # Add your business logic here return your_datetime # 通过 DeadlineReference.TYPES 指定求值时机:入队时执行 @deadline_reference(DeadlineReference.TYPES.DAGRUN_QUEUED) class MyQueuedReference(BaseDeadlineReference): """A custom reference evaluated when Dag runs are queued.""" # 声明需要 Airflow 传入的 Dag Run context 值 required_kwargs = {"dag_id", "run_id"} def _evaluate_with(self, *, session: Session, **kwargs) -> datetime: dag_id = kwargs["dag_id"] run_id = kwargs["run_id"] # Use dag_id and run_id in your calculation return your_datetime # 注册类,让调度器在反序列化 Dag 时能够解析 class MyDeadlineReferencePlugin(AirflowPlugin): name = "my_deadline_reference_plugin" deadline_references = [MyCustomDecoratedReference, MyQueuedReference]在 Dag 中使用自定义 Reference
注册完成后,即可像内置 Reference 一样在 Dag 定义中使用:
from datetime import timedelta from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference with DAG( dag_id="custom_reference_example", deadline=DeadlineAlert( reference=DeadlineReference.MyCustomDecoratedReference, interval=timedelta(hours=2), callback=AsyncCallback(my_callback), ), ): # Your tasks here ...自定义 Reference 的五条硬性约束
官方文档对自定义 Reference 提出了明确要求,违反任何一条都会在注册或求值时直接报错:
- 时区感知:
_evaluate_with必须返回 timezone-aware 的 datetime 对象; - 无参构造:自定义 Reference 在注册时会被实例化,因此必须能用无参数构造。如果确实需要参数,用
@dataclass装饰并给每个字段默认值; - 插件注册:必须列入某个
AirflowPlugin的deadline_references属性。只调用register_custom_reference(装饰器内部行为)只会影响运行 Dag 文件的进程,不等于注册了插件;未注册插件会在反序列化时抛出DeadlineReferenceNotRegistered; - API Server 重启:新增或修改自定义 Reference 后需要重启 Airflow API Server;
required_kwargs白名单:required_kwargs声明 Airflow 应向_evaluate_with()转发的 Dag Run context 值,目前只有dag_id和run_id可用,声明其他键会在求值时抛出ValueError。若要给 Reference 自身传配置,应使用构造字段或读取 Airflow Variable;需要查库时使用session参数。
从源码看,装饰器@deadline_reference(sdk 定义)既可以裸用(@deadline_reference,等价于@deadline_reference(),Dag Run 创建时求值),也可以带参指定求值时机(如@deadline_reference(DeadlineReference.TYPES.DAGRUN_QUEUED))。它内部调用DeadlineReference.register_custom_reference,把类注册为DeadlineReference.<ClassName>并加入对应的 TYPES 分类元组;核心侧则通过 plugins_manager.py 的get_deadline_references_plugins()收集插件中注册的类,供反序列化时解析。
进阶:一个 DAG 挂多个 Deadline,构建分级告警
Dag 的deadline参数既可以传单个DeadlineAlert,也可以传一个列表。列表中的每个告警独立求值、互不影响,且可以自由混用不同的 Reference 与回调类型——这是构建分级(tiered)告警策略的官方推荐姿势:
from datetime import timedelta from airflow.sdk import AsyncCallback, DAG, DeadlineAlert, DeadlineReference, SyncCallback from airflow.providers.slack.notifications.slack_webhook import SlackWebhookNotifier from airflow.providers.standard.operators.empty import EmptyOperator with DAG( dag_id="multiple_deadline_alerts", deadline=[ # 第一级:入队 30 分钟后未完成,Slack 异步提醒 DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(minutes=30), callback=AsyncCallback( SlackWebhookNotifier, kwargs={"text": "Dag {{ dag_run.dag_id }} is approaching its deadline."}, ), ), # 第二级:入队 60 分钟后仍未完成,自定义同步回调升级告警 DeadlineAlert( reference=DeadlineReference.DAGRUN_QUEUED_AT, interval=timedelta(minutes=60), callback=SyncCallback( "my_plugins.escalation.escalate_to_oncall", kwargs={"severity": "high"}, ), ), ], ): EmptyOperator(task_id="example_task")这个模式的意义在于:先用低门槛的异步通知做"预警",再用高门槛的同步回调做"升级",把告警从信息噪音中区分出来。注意这里SyncCallback的第一个参数也可以直接传点路径字符串(如"my_plugins.escalation.escalate_to_oncall"),Core 侧的序列化机制会按点路径解析可调用对象。
从 SLA 迁移到 Deadline 的决策清单
综合官方迁移文档与源码实现,给出迁移时的完整决策清单:
- 确认行为差异可接受:Deadline 在超时瞬间(最迟
scheduler_heartbeat_sec,默认 5 秒)就触发回调,而 SLA 要等 Dag Run 结束。告警会明显提前,请与接收方对齐预期; - 选择基准点:默认首选
DeadlineReference.DAGRUN_LOGICAL_DATE(最接近logical_date + sla);若关心"入队后是否及时跑起来"选DAGRUN_QUEUED_AT;有硬性截止时间用FIXED_DATETIME;运行时长波动大用AVERAGE_RUNTIME; - 设置间隔:
interval可为正(基准点之后)可为负(基准点之前),timedelta即可满足绝大多数场景,动态场景可改用VariableInterval(变量值以秒为单位的整数,解析发生在 Dag Run 创建时;修改变量只影响新解析的 DAG 与未来的 Dag Run,不会回溯更新已存在的 deadline); - 选择回调类型:内置 Notifier 优先配
AsyncCallback(运行在 Triggerer);需要最高优先级执行的自定义逻辑用SyncCallback(运行在 executor,可用executor=指定执行器); - 验证环境:确认
[scheduler] scheduler_heartbeat_sec满足你的告警时效要求;自定义异步回调记得重启 Triggerer,自定义 Reference 记得重启 API Server 并确认插件已注册(见 plugins_manager.py); - 从简单开始:先在单个 DAG 上用
DAGRUN_LOGICAL_DATE做等价迁移验证,再逐步引入多级告警与自定义 Reference。
迁移完成后的验证可以参考仓库自带的示例与测试:官方示例 DAG example_deadline_alert.py 演示了完整的DAGRUN_QUEUED_AT用法;核心侧单元测试 tests/unit/models/test_deadline.py(覆盖TestDeadline、TestCalculatedDeadlineDatabaseCalls、TestDeadlineReference、TestCustomDeadlineReference等)以及 UI API 测试 tests/unit/api_fastapi/core_api/routes/ui/test_deadlines.py 可作为行为边界的权威参考。关于 Deadline Alerts 的完整特性说明,请继续阅读 deadline-alerts 指南。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考