
Apache Airflow 从 SLA 迁移到 Deadline Alerts 完整实战指南范式对比、迁移路径与源码级原理解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowSLAService 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的目标非常相似——都是规定一个最晚时间点超时即告警——但它们采用的是两条完全不同的技术路线。理解这两条路线的差异是做出迁移决策的前提。SLADag Run 结束后才检查SLA 的工作方式如下当 Dag Run结束finishes时检查当前时间如果当前时间大于(logical_date sla)则执行sla_miss_callback如果 Dag Run 永远没有结束SLA 永远不会被检查。也就是说SLA 是一种事后判定机制它把检查动作挂在了 Dag Run 的结束事件上靠的是任务完成后的一次性判断而不是持续的轮询。这意味着两个天然的局限运行中的 Dag 无论超时多久在它结束之前你都收不到任何提醒一个卡死既不失败也不结束的 Dag 会完全绕过 SLA 检查。Deadline AlertsDag 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, ofDeadline, sessionsession, skip_lockedTrue, key_shareFalse, ) ): deadline.handle_miss(session)这段代码中值得注意的实现细节是with_row_locks(..., skip_lockedTrue)在启用 HA高可用多调度器副本时通过行级锁FOR UPDATE SKIP LOCKED保证同一行 Deadline 只会被一个调度器处理避免重复创建回调。这正是 Deadline 范式能够在运行中就感知超时的底层保障。在 config.yml 中可以看到该配置项的定义配置项所属 Section类型默认值说明scheduler_heartbeat_sec[scheduler]integer5调度器尝试触发新任务 / 执行调度循环的频率秒也就是说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 分钟的那一刻就发出告警。如果团队依赖的是事后统计迁移后告警的到达时机将显著提前这通常是期望中的改进但也需要提前与告警接收方对齐预期。等价示例 DAG1 小时 SLA vs 1 小时 Deadline下面先给出一个使用 1 小时 SLA 的 Dag再给出一个功能等价的、使用 Deadline Alerts 的 Dag两者可以并排对照。SLA 版本with DAG( minimal_sla_example, default_args{sla: timedelta(hours1)}, sla_miss_callbackSlackWebhookNotifier( textSLA missed for {{ dag_run.dag_id }}, ), ): BashOperator(task_idlong_task, bash_commandsleep 3600)Deadline Alerts 版本with DAG( minimal_deadline_example, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_LOGICAL_DATE, intervaltimedelta(hours1), callbackAsyncCallback( SlackWebhookNotifier, kwargs{ text: Deadline missed for {{ dag_run.dag_id }}, }, ), ), ): BashOperator(task_idlong_task, bash_commandsleep 3600)两个版本在语义上的对照关系如下概念SLA 版本Deadline 版本超时阈值default_args{sla: timedelta(hours1)}intervaltimedelta(hours1)基准点隐含使用logical_date显式声明referenceDeadlineReference.DAGRUN_LOGICAL_DATE回调sla_miss_callbackcallbackAsyncCallback(...)包一层 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_iddeadline_alert_example, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes15), callbackAsyncCallback( SlackWebhookNotifier, kwargs{ text: Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }} }, ), ), ): EmptyOperator(task_idexample_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、referenceJSON、intervalJSON与callback_defJSON并通过外键关联到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(sessionsession, 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.pyDeadlineReference.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_idexample_deadline_alert, ... deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(seconds30), callbackAsyncCallback(notify_deadline_missed), nameexample_deadline, ), ) as dag: task def hello_deadline(): time.sleep(60) # 故意睡过 30 秒的 deadline触发告警DeadlineReference.DAGRUN_LOGICAL_DATE引用 Dag Run计划开始scheduled to start的时间即迁移文档中推荐的、与 SLA 语义最接近的基准点。例如设置intervaltimedelta(minutes15)则无论 Dag 实际何时开始甚至从未开始只要在计划开始时间后 15 分钟尚未完成就会触发告警。适用于确保定时 DAG 在下一轮调度前完成。DeadlineReference.FIXED_DATETIME指定一个固定的时间点。适用于业务上存在硬性完成时间要求的场景如每天 10:00 前必须出报表。其核心实现FixedDatetimeDeadline直接返回构造时传入的 datetime不依赖数据库查询。DeadlineReference.AVERAGE_RUNTIME基于历史成功运行的耗时平均值动态计算 deadline。它分析历史执行数据来预测当前运行应该在何时完成deadline 当前时间 平均运行时长 interval。如果历史数据不足则不创建 deadline也就不会误报。参数说明max_runsint可选纳入统计的最近成功运行次数上限默认 10min_runsint可选计算平均值所需的最少成功运行次数默认与max_runs相同。# 使用默认设置分析最近 10 次运行且要求至少 10 次 DeadlineReference.AVERAGE_RUNTIME() # 分析最近 20 次运行但只要有 5 次即可计算 DeadlineReference.AVERAGE_RUNTIME(max_runs20, min_runs5) # 严格模式必须恰好有 15 次运行才计算 DeadlineReference.AVERAGE_RUNTIME(max_runs15, min_runs15)从源码实现看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_idaverage_runtime_deadline, deadlineDeadlineAlert( referenceDeadlineReference.AVERAGE_RUNTIME(max_runs15, min_runs5), intervaltimedelta(minutes30), # 超过平均耗时 30 分钟即告警 callbackAsyncCallback( SlackWebhookNotifier, kwargs{text: Dag {{ dag_run.dag_id }} is running longer than expected!}, ), ), ): EmptyOperator(task_iddata_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(days1), time(10, 0)) with DAG( dag_idfixed_deadline_alert, deadlineDeadlineAlert( referenceDeadlineReference.FIXED_DATETIME(tomorrow_at_ten), intervaltimedelta(minutes-30), # 在基准点前 30 分钟告警 callbackAsyncCallback( SlackWebhookNotifier, kwargs{ text: Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }} }, ), ), ): EmptyOperator(task_idexample_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_idslack_deadline_alert_async, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes30), callbackAsyncCallback( SlackWebhookNotifier, kwargs{ text: Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }} }, ), ), ): EmptyOperator(task_idexample_task) # 同步版本回调运行在 executor with DAG( dag_idslack_deadline_alert_sync, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes30), callbackSyncCallback( SlackWebhookNotifier, kwargs{ text: Dag {{ dag_run.dag_id }} missed deadline at {{ deadline.deadline_time }}. DagRun: {{ dag_run }} }, ), ), ): EmptyOperator(task_idexample_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.pydef custom_sync_callback(**kwargs): Handle deadline violation with custom logic. context kwargs.get(context, {}) print(fDeadline exceeded for Dag {context.get(dag_run, {}).get(dag_id)}!) print(fContext: {context}) print(fAlert 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_idcustom_sync_deadline_alert, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes15), callbackSyncCallback( custom_sync_callback, kwargs{alert_type: time_exceeded}, ), ), ): EmptyOperator(task_idexample_task)自定义异步回调示例第 1 步放入 plugins 文件夹async def custom_async_callback(**kwargs): Handle deadline violation with custom logic. context kwargs.get(context, {}) print(fDeadline exceeded for Dag {context.get(dag_run, {}).get(dag_id)}!) print(fContext: {context}) print(fAlert 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_idcustom_deadline_alert, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes15), callbackAsyncCallback( custom_async_callback, kwargs{alert_type: time_exceeded}, ), ), ): EmptyOperator(task_idexample_task)进阶提示SyncCallback支持可选的executor参数用于把回调路由到指定 executor不指定则使用默认 executorSyncCallback( my_callback, kwargs{msg: deadline missed}, executorcelery_executor, )AsyncCallback支持可选的queue参数把由此产生的 trigger 分配到特定 trigger queue不指定则运行在任意未加--queues限制的 triggerer 上AsyncCallback( my_callback, kwargs{msg: deadline missed}, queuealerts, )模板化与简化 ContextDeadline 回调当前收到的是一份简化版的 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基准时间来源典型场景常用 intervalDAGRUN_QUEUED_ATDagRun 表的queued_at监控资源挤压、入队后迟迟不启动正数入队后 X 分钟DAGRUN_LOGICAL_DATEDagRun 表的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( referenceDeadlineReference.FIXED_DATETIME(next_meeting), intervaltimedelta(hours-2), callbacknotify_team, ) # 场景 B计划 1 小时内未完成即告警DAGRUN_LOGICAL_DATE 正 interval DeadlineAlert( referenceDeadlineReference.DAGRUN_LOGICAL_DATE, intervaltimedelta(hours1), callbacknotify_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.pyfrom 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_idcustom_reference_example, deadlineDeadlineAlert( referenceDeadlineReference.MyCustomDecoratedReference, intervaltimedelta(hours2), callbackAsyncCallback(my_callback), ), ): # Your tasks here ...自定义 Reference 的五条硬性约束官方文档对自定义 Reference 提出了明确要求违反任何一条都会在注册或求值时直接报错时区感知_evaluate_with必须返回 timezone-aware 的 datetime 对象无参构造自定义 Reference 在注册时会被实例化因此必须能用无参数构造。如果确实需要参数用dataclass装饰并给每个字段默认值插件注册必须列入某个AirflowPlugin的deadline_references属性。只调用register_custom_reference装饰器内部行为只会影响运行 Dag 文件的进程不等于注册了插件未注册插件会在反序列化时抛出DeadlineReferenceNotRegisteredAPI Server 重启新增或修改自定义 Reference 后需要重启 Airflow API Serverrequired_kwargs白名单required_kwargs声明 Airflow 应向_evaluate_with()转发的 Dag Run context 值目前只有dag_id和run_id可用声明其他键会在求值时抛出ValueError。若要给 Reference 自身传配置应使用构造字段或读取 Airflow Variable需要查库时使用session参数。从源码看装饰器deadline_referencesdk 定义既可以裸用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_idmultiple_deadline_alerts, deadline[ # 第一级入队 30 分钟后未完成Slack 异步提醒 DeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes30), callbackAsyncCallback( SlackWebhookNotifier, kwargs{text: Dag {{ dag_run.dag_id }} is approaching its deadline.}, ), ), # 第二级入队 60 分钟后仍未完成自定义同步回调升级告警 DeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes60), callbackSyncCallback( my_plugins.escalation.escalate_to_oncall, kwargs{severity: high}, ), ), ], ): EmptyOperator(task_idexample_task)这个模式的意义在于先用低门槛的异步通知做预警再用高门槛的同步回调做升级把告警从信息噪音中区分出来。注意这里SyncCallback的第一个参数也可以直接传点路径字符串如my_plugins.escalation.escalate_to_oncallCore 侧的序列化机制会按点路径解析可调用对象。从 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),仅供参考