
Apache Airflow 3.2 中同步 deadline 回调对 Connection 与 Variable 的元数据库访问支持【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读本文基于 Apache Airflow 仓库中 newsfragment 65269 的技术变更airflow-core/newsfragments/65269.significant.rst展开讲解 Airflow 3.2 中新增的一项能力同步 deadline 回调SyncCallback现在可以访问元数据数据库中的 Connections 和 Variables。文章将说明该能力出现前的限制、底层实现原理、源码级证据以及与AsyncCallback的差异帮助你在 Deadline Alerts 场景中正确选择回调类型并编写可访问元数据的同步回调。背景Deadline Alerts 与两类回调Deadline Alerts 是 Airflow 3.1 引入的当前仓库中仍标注为实验性功能Dag Run 时间阈值告警机制允许你为 Dag Run 设置必须完成时间超过该时间后自动执行回调。一个DeadlineAlert由三个要素构成reference参考点计时起点如DeadlineReference.DAGRUN_QUEUED_AT入队时刻、DeadlineReference.DAGRUN_LOGICAL_DATE逻辑日期、DeadlineReference.FIXED_DATETIME固定时间点或DeadlineReference.AVERAGE_RUNTIME历史平均运行时长interval间隔以timedelta表示正负皆可决定在参考点之后多久或之前多久触发callback回调deadline 被错过时执行的AsyncCallback或SyncCallback对象。计算方式为Deadline Reference Interval见 airflow-core/docs/howto/deadline-alerts.rst 中的示意图[Reference] ------ [Interval] ------ [Deadline] ^ ^ | | Start time Trigger point在 Airflow 3.2 之前SyncCallback同步回调在 deadline 场景中存在一个明显的功能缺口它无法访问元数据数据库中的 Connections 和 Variables。本 newsfragment 正是补齐了这一缺口。变更核心同步回调现在可以读取元数据库65269.significant.rst的完整内容只有一句话Synchronous deadline callbacks (SyncCallback) can now access Connections and Variables from the Airflow metadata database.翻译过来即同步 deadline 回调SyncCallback现在可以访问来自 Airflow 元数据数据库的 Connections 和 Variables。这是一条significant级别的 newsfragment意味着该变更对使用者来说是可见的行为变化值得在升级时关注。它直接扩展了SyncCallback的适用范围此前如果你需要在 deadline 回调里查询 Connection 或 Variable只能选用在 Triggerer 上异步执行的AsyncCallback而现在同步回调同样具备该能力。同步回调 vs 异步回调执行位置差异要理解这条变更的意义先要理解两类回调的执行位置差异这一设计在 Task SDK 中定义AsyncCallback异步回调运行在 Triggerer 上回调必须是一个可等待awaitable的顶层函数。通过airflow.sdk.AsyncCallback导入注意Airflow 3.2 中其导入路径从airflow.sdk.definitions.deadline变更到了airflow.sdk见 deadline-alerts.rst 中的说明。可选参数queue用于指定 Trigger 队列。SyncCallback同步回调发送给执行器executor并被当作最高优先级的 DAG 任务来执行见 deadline-alerts.rst 的 noteSync callbacks are sent to the executor and treated just like a Dag task with top priority。可选参数executor用于指定目标执行器未指定时使用默认执行器。两类回调的定义位于 task-sdk/src/airflow/sdk/definitions/callback.pyCallback基类负责把 callable 转成可导入的 dot-path 并序列化AsyncCallback增加queue字段并校验 callable 必须可 awaitSyncCallback增加executor字段。回调被序列化后pathkwargs由 Core 侧的反序列化逻辑还原成可执行对象。为什么之前同步回调访问不到元数据库从源码结构看AsyncCallback与SyncCallback的运行环境差异导致了两者的数据访问能力差异异步回调在 Triggerer 进程中通过CallbackTrigger见 airflow-core/src/airflow/triggers/callback.py触发触发逻辑运行在 Triggerer 内拥有与元数据库的交互通道同步回调被封装为ExecutorCallback发送给执行器见 airflow-core/src/airflow/models/callback.py 中create_from_sdk_def对SyncCallback的处理执行时以ExecuteCallbackworkload 的形式运行见 airflow-core/src/airflow/executors/workloads/callback.py。在旧版本中这条执行路径并未打通对元数据库 Connections/Variables 的访问因此同步回调内读取这些资源会失败或不可用。底层实现Task SDK 的执行期数据访问通道要理解现在可以访问在代码层面意味着什么需要看 Task SDK 中执行期execution-time的存取实现。连接与变量的读取入口在 task-sdk/src/airflow/sdk/execution_time/context.py_get_connection(conn_id)依次检查预置连接ContextVar 注入、SecretCache缓存然后遍历已加载的 secrets backends包括 worker 上下文中的SupervisorCommsSecretsBackend或 API server 上下文中的MetastoreBackend任一 backend 命中即返回Connection同时会对连接中的密钥做掩码处理_mask_connection_secrets全部未命中则抛出AirflowNotFoundException。_get_variable(key, deserialize_json)同样先查SecretCache再遍历 secrets backends命中后经_mask_and_deserialize_variable完成密钥掩码与 JSON 反序列化未命中抛出AirflowRuntimeErrorVARIABLE_NOT_FOUND。在 SDK 的面向用户 API 中对应为 task-sdk/src/airflow/sdk/definitions/connection.py 的Connection.get(conn_id)延迟导入_get_connection与 task-sdk/src/airflow/sdk/definitions/variable.py 的Variable.get(key, defaultNOTSET, deserialize_jsonFalse)。在执行器/worker 环境中元数据数据库的访问正是通过SupervisorCommsSecretsBackend经 supervisor 进程与 API server 通信由 API server 侧的MetastoreBackend直连元数据库来实现的。因此这条变更的实际效果是同步回调在执行器上运行时经由这套 secrets backend 通道即可读取元数据库中的 Connections 与 Variables无需回调自己额外建立数据库连接。如何在同步回调中使用 Connection 与 Variable下面给出一个可直接运行的示例。回调函数被定义为顶层函数并放入 plugins 目录例如$AIRFLOW_HOME/plugins/deadline_callbacks.py这是因为SyncCallback要求 callable 必须在执行该回调的 worker 上可导入task-sdk/src/airflow/sdk/definitions/callback.py 中Callback.__init__会校验 dot-path 合法性。# $AIRFLOW_HOME/plugins/deadline_callbacks.py from airflow.sdk import Connection, Variable def notify_on_deadline(**kwargs): Deadline 错过时的同步回调读取元数据库中的 Connection 与 Variable。 context kwargs.get(context, {}) dag_id context.get(dag_run, {}).get(dag_id) # 从元数据库读取 Connection在 3.2 之前同步回调中不可用 conn Connection.get(my_alert_api_conn) endpoint conn.host # 从元数据库读取 Variabledeserialize_jsonTrue 可解析 JSON 值 alert_channel Variable.get(alert_channel, defaultgeneral) retry Variable.get(alert_max_retries, default3, deserialize_jsonTrue) print(fDag {dag_id} missed deadline; notify via {alert_channel} at {endpoint}, retry{retry})然后在 DAG 文件中使用该回调from datetime import timedelta from deadline_callbacks import notify_on_deadline from airflow.providers.standard.operators.empty import EmptyOperator from airflow.sdk import DAG, DeadlineAlert, DeadlineReference, SyncCallback with DAG( dag_idsync_callback_with_metadata, deadlineDeadlineAlert( referenceDeadlineReference.DAGRUN_QUEUED_AT, intervaltimedelta(minutes15), callbackSyncCallback( notify_on_deadline, kwargs{alert_type: time_exceeded}, ), ), ): EmptyOperator(task_idexample_task)需要说明的几点约束均来自 deadline-alerts.rst 与 SDK 源码context 参数deadline 被错过时 Airflow 会自动向回调传入contextkwarg包含 Dag Run 与 deadline 信息context是保留关键字不能出现在Callback的kwargs中否则 DAG 解析阶段会抛出ValueError见 task-sdk/src/airflow/sdk/definitions/callback.py 第 61-62 行的检查。指定执行器SyncCallback支持executorcelery_executor之类的参数定向到指定执行器未指定则用默认执行器。回调必须可导入支持传 callable 对象或字符串 dot-path如my_plugins.escalation.escalate_to_oncall。传字符串时 SDK 会尽力验证其可导入性但仅记录 debug 日志而非直接失败——因为 callable 可能存在于 DAG processor 之外的机器上。验证路径测试与文档功能文档Deadline Alerts 的完整指南见 airflow-core/docs/howto/deadline-alerts.rst其中明确写着SyncCallback support added in 3.2即同步回调支持是在 3.2 版本中加入的与本 newsfragment 相互印证。迁移文档从 SLA 迁移到 Deadline 的指南见 airflow-core/docs/howto/sla-to-deadlines.rst。核心模型Callback 的 ORM 模型含ExecutorCallback、TriggererCallback的多态映射见 airflow-core/src/airflow/models/callback.pyDeadline 与 DeadlineAlert 的元数据库模型见 airflow-core/src/airflow/models/deadline.pydeadline表记录deadline_time、missed状态并通过callback_id外键关联回调。执行负载ExecuteCallbackworkload 的 schema 见 airflow-core/src/airflow/executors/workloads/callback.py其中CallbackFetchMethod.IMPORT_PATH的注释明确写着For deadline callbacks since they import callbacks through the import path即 deadline 回调通过导入路径获取。小结Airflow 3.2 起SyncCallback可以在执行器上直接访问元数据数据库中的 Connections 和 Variables与AsyncCallback在这一能力上对齐。底层实现依赖 Task SDK 执行期通过 secrets backendsworker 经 supervisor 与 API server 通信API server 侧MetastoreBackend直连元数据库读取资源回调内使用Connection.get()与Variable.get()即可无需自行连接数据库。如果你的 deadline 处理逻辑依赖外部系统凭据或动态配置同步回调现在是可行选项若回调本身是异步 I/O 密集型逻辑仍应选用在 Triggerer 上运行的AsyncCallback。注意Deadline Alerts 在 Airflow 3.1 中引入并被标注为实验性功能见 deadline-alerts.rst 的 warning特性可能在后续版本中根据用户反馈发生调整本文描述的能力以当前仓库Airflow 3.2 开发版为准。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考