权威指南:基于 airflow.sdk 的稳定集成与扩展体系)
Apache Airflow 3.0 公共接口Public Interface权威指南基于 airflow.sdk 的稳定集成与扩展体系【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本篇技术指南系统讲解 Apache Airflow 3.0 的 Public Interface公共接口体系它定义了哪些接口受语义化版本控制、DAG 作者如何通过新的airflow.sdk命名空间编写稳定且可跨版本迁移的工作流以及平台维护者如何通过 Plugins、Triggers、Timetables、Executors、Secrets Backends 等扩展点增强 Airflow 能力。读完本文你将掌握从 直接引用内部模块 迁移到 公共接口 的完整路径学会使用 Task Context、Stable REST API 与 Python Client 替代被禁止的元数据库直连并了解集成外部系统时的全部官方扩展方式。什么是 Airflow 的 Public InterfaceApache Airflow 的 Public Interface 是 Airflow 中行为受语义化版本Semantic Versioning约束的接口与行为的集合。用户通过以下方式与公共接口交互创建和管理 DAG、管理任务与依赖关系通过编写新的 Executor、Plugin、Operator 和 Provider 扩展 Airflow 能力构建自定义工具、与其他系统集成以及自动化 Airflow 工作流的某些环节。公共接口的承诺是在同一个 MAJOR 版本内签名与行为保持向后兼容。以_开头的方法protected与以__开头的方法private不属于公共接口可能随时变化。⚠️ 本文档面向Airflow 3.0。若仍在使用 Airflow 2.x请参考 2.11 的旧版公共接口文档。公共接口的核心约定任务代码禁止直连元数据库从任务代码直接访问元数据库查询 DAG 状态、任务历史或 Dag Run在 Airflow 3 中已不再被允许工作进程只能通过 Execution API 通信。替代方案是 Task Context、Stable REST API 与 Python Client。推荐使用airflow.sdk命名空间自 Airflow 3.0 起airflow.sdk是官方公共接口对应 AIP-72 Task Execution Interface aka Task SDK其目标是让 DAG 编写与 Airflow 内部实现Scheduler、API Server 等解耦提供跨版本稳定、与版本无关的 DAG 编写与维护接口。稳定的程序化访问走 Stable REST APIREST API 基于 OpenAPI 规范定义CLI 的输出格式与可用参数可能在细节上变化因此需要程序化依赖时应优先使用 Stable REST API。公共接口的典型使用场景为现有用例编写自定义 Operator 或 Hook编写扩展 Airflow 功能的 Plugin如 Secrets、Timetables、Triggers、Listeners通常由 Airflow 实例管理员完成将自定义 Operator、Hook、Plugin 打包成 Provider 发布供外部服务或应用复用使用 TaskFlow API 编写任务依赖 Airflow 对象的一致行为。airflow.sdkDAG 作者的主接口airflow.sdk命名空间是 DAG 作者的主要公共接口。它提供了一个稳定、定义良好的 DAG 与任务创建接口不受内部实现变更影响。查看仓库中 task-sdk/src/airflow/sdk/init.py 的__all__列表当前 Task SDK 版本为 1.4.0可以拿到该命名空间导出的完整符号清单其全部类、装饰器与函数均通过惰性导入__lazy_imports按需加载。核心类Classes类用途airflow.sdk.Asset数据资产定义用于基于资产的调度airflow.sdk.BaseHook所有 Hook 的基类airflow.sdk.BaseNotifier通知器基类airflow.sdk.BaseOperator所有 Operator 的基类airflow.sdk.BaseOperatorLink自定义任务 Extra Link 的基类airflow.sdk.BaseSensorOperator传感器 Operator 基类airflow.sdk.Connection外部服务凭据与配置访问airflow.sdk.Context任务执行上下文TypedDictairflow.sdk.DAGDAG 核心实体airflow.sdk.EdgeModifier依赖边修饰器airflow.sdk.Label依赖边标签airflow.sdk.ObjectStoragePath对象存储路径抽象airflow.sdk.ParamDAG 参数定义airflow.sdk.TaskGroup任务组airflow.sdk.VariableAirflow 配置变量访问装饰器与函数Decorators and Functions符号用途airflow.sdk.asset资产定义装饰器airflow.sdk.dagDAG 定义装饰器airflow.sdk.task任务定义装饰器airflow.sdk.task_group任务组装饰器airflow.sdk.setup/airflow.sdk.teardownsetup/teardown 任务装饰器airflow.sdk.result结果收集装饰器airflow.sdk.chain/airflow.sdk.chain_linear/airflow.sdk.cross_downstream依赖编排函数airflow.sdk.get_current_context获取当前任务执行上下文airflow.sdk.get_parsing_context获取 DAG 解析上下文从 Airflow 2.x 迁移所有 DAG 都应改用airflow.sdk导入不再直接引用内部模块。旧导入路径如airflow.models.dag.DAG、airflow.decorator.task已被弃用并将在未来版本中移除。例如# 旧写法已弃用 # from airflow.models.dag import DAG # from airflow.decorators import task # 新写法 from airflow.sdk import DAG, dag, task详细的迁移说明包括导入变化与破坏性变更参见 installation 目录下的升级指南。airflow.sdk的完整可用类、装饰器与函数清单可直接检查airflow.sdk.__all__。DAG工作流的核心实体DAG 是 Airflow 表示周期性工作流的核心实体。可以通过在 DAG 文件中实例化airflow.sdk.DAG类来创建并可通过airflow.sdk.Param指定参数。官方推荐的方式是使用airflow.sdk命名空间中的dag装饰器。需要特别注意的边界airflow.models.dagbag.DagBag类仅用于 Airflow 内部从文件与目录加载 DAGDAG 作者应改用airflow.sdk.DAGairflow.models.dagrun.DagRun类仅用于内部 Dag Run 管理DAG 作者应通过 Task Contextget_current_context或airflow.sdk.types.DagRunProtocol接口获取 Dag Run 信息。从源码可见task-sdk/src/airflow/sdk/types.py 中DagRunProtocol定义了运行期可用的最小 Dag Run 接口dag_id、run_id、logical_date、data_interval_start/end、state、conf、partition_key等字段这就是任务代码能安全读取的 Dag Run 视图。Operators编写自定义任务airflow.sdk.BaseOperator与airflow.sdk.BaseSensorOperator是公开基类可被继承以创建新 Operator。在 task-sdk/src/airflow/sdk/bases/operator.py 中可以看到BaseOperator的完整实现它内置了重试策略、优先级权重含数据库安全范围裁剪db_safe_priority、依赖规则、模板渲染等通用能力。关于 Operator 的公共性有一个重要区分Airflow 官方发布的BaseOperator子类在行为上是公共的但在结构上不是——即 Operator 的参数与行为受语义化版本约束但其内部方法随时可能变化不要继承或依赖这些方法。Task Instances 与 Task Instance Keys任务实例Task Instance是某个任务在某个 Dag Run 中的单次运行。任务实例通过 Task Contextget_current_context访问不允许直接访问数据库。任务实例键Task Instance Key是任务实例的唯一标识是一个由dag_id、task_id、run_id、try_number、map_index组成的元组。任务代码不再允许通过airflow.models.taskinstance.TaskInstance模型直接访问键airflow.models.taskinstancekey.TaskInstanceKey同样仅限内部使用。通过 Task Context 访问任务实例信息from airflow.sdk import get_current_context def my_task(): context get_current_context() ti context[ti] dag_id ti.dag_id task_id ti.task_id run_id ti.run_id try_number ti.try_number map_index ti.map_index print(fTask: {dag_id}.{task_id}, Run: {run_id}, Try: {try_number}, Map Index: {map_index})get_current_context的实现位于 task-sdk/src/airflow/sdk/definitions/context.py它返回一个ContextTypedDict其中ti的类型是RuntimeTaskInstanceProtocol——一个定义在 task-sdk/src/airflow/sdk/types.py 的最小协议接口包含dag_id、run_id、try_number、map_index、state等属性以及xcom_pull、xcom_push、get_template_context、get_dr_count、get_dagrun_state、get_task_states等方法。这意味着任务代码只能通过该协议暴露的能力与 Airflow 交互底层数据获取由 Task SDK 通过 Supervisor 与 Execution API 完成。Hooks外部平台的统一接口Hook 是访问外部平台与数据库的接口在可能的情况下实现统一接口并作为 Operator 的构建块。所有 Hook 都派生自airflow.sdk.bases.hook.BaseHook。Airflow 内置的公共 Hook 集合允许你自由继承扩展其功能见 airflow-core/docs 中 Hooks 相关 API 参考。公共工具类Connection、Variable 与 XCom编写或扩展 Hook 与 Operator 时DAG 作者与开发者可以使用以下公共工具类airflow.sdk.Connection访问外部服务凭据与配置airflow.sdk.Variable访问 Airflow 配置变量airflow.sdk.execution_time.xcom.XCom访问任务间通信XCom数据。关键约束任务代码不再允许直接访问数据库中的airflow.models.connection.Connection与airflow.models.variable.Variable模型。Connection 与 Variable 操作应通过 Task Contextget_current_context与任务实例方法或直接通过airflow.sdk命名空间完成。通过 Task Context 访问 Connection 与 Variablefrom airflow.sdk import get_current_context def my_task(): context get_current_context() conn context[conn] my_connection conn.get(my_connection_id) var context[var] my_variable var.value.get(my_variable_name)直接使用airflow.sdk命名空间from airflow.sdk import Connection, Variable conn Connection.get(my_connection_id) var Variable.get(my_variable_name)更多信息可参阅 connections 相关 howto、variables 文档 与 xcoms 文档。公共异常与工具类编写自定义 Operator 与 Hook 时可以抛出与捕获 Airflow 暴露的公共异常见airflow/exceptions相关 API 参考也可复用airflow/utils/state下的公共工具类如状态枚举与状态相关辅助函数。使用 Plugin 机制扩展 Airflow 平台能力Airflow 通过 Plugin 机制扩展平台能力Plugin 不仅用于扩展 UI也是暴露以下自定义能力的官方途径Triggers、Timetables、Listeners 等。Provider 也可以实现 Plugin 端点并定制 Airflow UI。阅读入口Plugin 总体机制administration-and-deployment/plugins.rst自定义 UI 视图插件howto 目录相关文档无需 Plugin 的简单 UI 定制howto 目录相关文档Triggers触发器Airflow 使用 Trigger 实现基于asyncio的 Deferrable Operators可延迟算子。所有 Trigger 派生自airflow.triggers.base.BaseTrigger。Airflow 内置的公共 Trigger 集合允许自由继承扩展。更多内容参见 authoring-and-scheduling 目录的 deferring 文档。Timetables时间表自定义 Timetable 实现为 Scheduler 提供内置调度表达式无法实现的新调度逻辑。所有 Timetable 派生自airflow.timetables.base.Timetable。内置公共 Timetable 集合可自由继承扩展相关教程见 howto 目录的 timetable 文档。Listeners监听器Listener 使你能够响应 Dag/Task 生命周期事件通过airflow.listeners.listener.ListenerManager类提供的钩子实现该接口自 2.5 版本加入。详细文档见 administration-and-deployment/listeners.rst。Extra Links额外链接Extra Links 是可以独立于自定义 Operator 动态添加到 Airflow 的链接。通常由 Operator 定义但 Plugin 允许在全局层面覆盖这些链接。参见 howto 目录的 define-extra-link 文档。与外部服务和应用的集成任务通过 Hook 与 Operator 编排外部服务Airflow 的核心功能如认证也可以扩展以对接外部服务。社区 Provider 及其核心扩展点详见仓库的 providers/ 目录。Executors执行器Executor 是任务实例得以运行的机制所有 Executor 派生自airflow.executors.base_executor.BaseExecutor。Executor 接口BaseExecutor类是公共的内置 Executor如 KubernetesExecutor、LocalExecutor 等不是公共的——Airflow 可能在 minor 或 patch 版本中修改内置 Executor导致继承它的派生 Executor 被破坏。因此若需修改或扩展内置 Executor应把完整 Executor 代码复制到自己的项目中避免上游变更影响派生实现。自 2.6 版本起Executor 已完全解耦Airflow 核心不再需要知道具体 Executor 的行为在 2.6 之前实现自定义 Executor 会受到一些硬编码行为的限制无法提供与内置 Executor 完全相同的功能。如何编写自己的 Executor 参见 core-concepts/executor/index.rst。Secrets Backends密钥后端Airflow 可配置为通过 Secrets Backend 获取airflow.sdk.Connection与airflow.sdk.Variable。所有 Secrets Backend 派生自airflow.secrets.base_secrets.BaseSecretsBackend且所有 Secrets Backend 实现都是公共的可以自由扩展。文档见 security 目录下的 secrets-backend 文档社区 Provider 实现的可用 Secrets Backend 见 providers/ 目录。Auth Managers认证管理器Auth Manager 负责 Airflow 中的用户认证与授权。所有 Auth Manager 派生自airflow.api_fastapi.auth.managers.base_auth_manager.BaseAuthManager。Auth Manager 接口BaseAuthManager类是公共的不同 Auth Manager 实现如 FabAuthManager不是公共的。如何编写自己的 Auth Manager 参见 core-concepts/auth-manager/。Connections 与 Extra Links创建 Hook 时可以添加自定义 Connection 与自定义 Extra Link任务运行时显示。社区 Provider 中可用的 Connection 与 Extra Link 实现参见 providers/ 目录。日志与监控可以扩展 Airflow 的日志写入方式相关文档见 administration-and-deployment/logging-monitoring/社区 Provider 实现的日志写入器见 providers/ 目录。更多公共扩展点DecoratorsDAG 作者可使用 TaskFlow 概念编写 DAG所有 Decorator 派生自airflow.sdk.bases.decorator.TaskDecorator。面向 DAG 作者的主要装饰器都位于airflow.sdk命名空间dag、task、asset、setup、task_group、teardown、chain、chain_linear、cross_downstream、get_current_context、get_parsing_context自定义 Decorator 教程见 howto 目录的 create-custom-decorator 文档。TaskFlow 概念详见 core-concepts/taskflow.rst。Email notifications内置邮件通知支持可添加自定义邮件通知类见 howto 目录的 email-config 文档。Notifications内置可扩展的通知发送机制基于各种on_*_callback见 howto 目录的 notifications 文档。Cluster Policies集群策略动态地对被解析的 DAG 或被执行的任务应用集群级策略见 administration-and-deployment/cluster-policies.rst。Lineage血缘Airflow 可帮助追踪数据来源、变化与流动路径见 administration-and-deployment/lineage.rst。什么不属于 Public Interface本文档未提及的一切内容都应视为非公共接口。在其他应用中这些组件或许可以依赖其向后兼容性但在 Airflow 中它们不是公共接口可能随时变化数据库结构被认为是内部实现细节不保证以向后兼容的方式维护Web UI持续演进HTML 元素不提供向后兼容保证本文档未明确提及的 Python 类被认为是内部实现细节不保证向后兼容。禁止直接访问元数据库由 DAG 作者编写的代码不再允许直接访问元数据库来查询 DAG 状态、任务历史或 Dag Run——工作进程完全通过 Execution API 通信。这一变化改善了架构隔离并支持远程执行能力。请使用以下替代方案Task Context使用airflow.sdk.get_current_context访问任务实例信息以及RuntimeTaskInstanceProtocol提供的方法get_dr_count、get_dagrun_state、get_task_statesStable REST API通过 stable-rest-api-ref.rst 以编程方式访问 Airflow 元数据Python Client使用 Python Client SDK 进行基于 Python 的交互。从源码可以印证这套设计在 task-sdk/src/airflow/sdk/execution_time/request_handlers.py 中handle_get_dr_count、handle_get_dagrun_state、handle_get_task_states等处理函数将请求转换为对 API Client 的调用并由 task-sdk/src/airflow/sdk/execution_time/supervisor.py 统一分发任务侧通过 task-sdk/src/airflow/sdk/execution_time/task_runner.py 中的get_dr_count、get_task_states、get_dagrun_state静态方法发起请求。可见任务代码对元数据的每次读取都经过任务进程 → Supervisor → Execution API → 元数据库的受控链路而不是直接 SQL 访问。用 Task Context 替代直接数据库访问的完整示例from airflow.sdk import dag, get_current_context, task, DagRunState from datetime import datetime dag(dag_idexample_dag, start_datedatetime(2025, 1, 1), schedulehourly, tags[misc], catchupFalse) def example_dag(): task(task_idcheck_dagrun_state) def check_state(): context get_current_context() ti context[ti] dag_run context[dag_run] # 使用 Task Context 方法替代直接数据库访问 dr_count ti.get_dr_count(dag_idexample_dag) dagrun_state ti.get_dagrun_state(dag_idexample_dag, run_iddag_run.run_id) return fDag run count: {dr_count}, current state: {dagrun_state} check_state() example_dag()实践建议总结写 DAG 永远从airflow.sdk导入DAG、dag、task、task_group、Param、Asset、chain、get_current_context等让工作流代码与 Airflow 内部实现解耦天然获得跨版本稳定性任务内读取元数据走 Task Contextcontext[ti]、context[dag_run]、context[conn]、context[var]是官方提供的数据访问入口配合RuntimeTaskInstanceProtocol的方法如get_dr_count、get_dagrun_state、get_task_states即可完成状态查询程序化集成优先 Stable REST API需要脚本或外部系统与 Airflow 交互时基于 OpenAPI 的 Stable REST API 比 CLI 更可靠扩展功能遵循公共扩展点自定义 Operator/Hook 继承airflow.sdk的基类平台级扩展Trigger、Timetable、Listener、Secrets Backend、Auth Manager继承对应公共基类修改内置 Executor 时复制完整代码而非子类化识别非公共边界数据库结构、Web UI HTML、未列出的 Python 类均不承诺向后兼容不要在其上构建长期依赖。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考