news 2026/9/6 19:00:46

Apache Airflow 新语言 SDK 开发实战:Coordinator 选型、Wire Protocol 与 AIP-108 贡献规范

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow 新语言 SDK 开发实战:Coordinator 选型、Wire Protocol 与 AIP-108 贡献规范

Apache Airflow 新语言 SDK 开发实战:Coordinator 选型、Wire Protocol 与 AIP-108 贡献规范

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

在 Apache Airflow 3.3 起,标准 Worker 可以通过"外部语言 SDK"(AIP-108)执行非 Python 编写的任务代码。本文以仓库中 airflow-new-sdk 技能文档 为核心骨架,结合其指认的权威贡献指南 Creating a new Language SDK 以及 Java/Go 两个生产级参考实现的源码,系统讲解:如何为 Airflow 引入一门新编程语言的任务执行能力——包括 Coordinator 基类选型、Supervisor Schema 版本协商、AFBNDL01 原生可执行包格式、MessagePack 线协议、日志通道规范、E2E 测试搭建以及 PR 提交清单,帮助贡献者完整走通从目录规划到 CI 接线的落地全流程。

一、两个组件:Coordinator 与目标语言 SDK

Airflow 要理解"如何执行一门外语写的任务",需要两个基本独立的组件:

  • Python 编写的 Coordinator(协调器):当匹配的 queue 上触发 stub 任务时,由 Airflow 调用。它唯一的必需方法是execute_task,职责是启动外部运行时、把任务交出去、并返回最终任务状态。基类为airflow.sdk.execution_time.coordinator.BaseCoordinator
  • 目标语言编写的 Language SDK:实现协调器所选协议的"另一端"。

两者没有强制统一的通信机制——可以是 TCP socket 上的子进程、gRPC 服务、共享内存、消息队列等,传输方式完全由 coordinator 与其 SDK 对应方自行约定。但实践中几乎所有 SDK 都走"子进程 + TCP"路线,因此仓库为这条路径提供了大量现成脚手架。任务执行的整体架构可参考 Task execution architecture。

二、仓库布局:新 SDK 的代码放哪里

技能文档给出了统一的目录约定。Coordinator(Python 侧)放入 task-sdk:

task-sdk/src/airflow/sdk/coordinators/<language>/ __init__.py # re-export + 模块 docstring coordinator.py # SubprocessCoordinator 或 BaseCoordinator 的子类 task-sdk/tests/coordinators/<language>/ test_coordinator.py task-sdk/tests/integration/coordinators/<language>/ test_integration.py # 需要 Breeze

而语言 SDK 本身放在仓库顶层的<language>-sdk/目录,与现有的 java-sdk/、go-sdk/、ts-sdk/ 并列。当前仓库中已存在 java、executable、node 三个 coordinator 实现(见 task-sdk/src/airflow/sdk/coordinators/),其单元测试实际位于task-sdk/tests/task_sdk/coordinators/<language>/(如 test_coordinator.py)。

三、选择 Coordinator 基类:决策树与源码印证

技能文档给出一棵简洁的选型决策树:

运行时是否编译为自包含的原生可执行文件? 是 → 直接用 ExecutableCoordinator(零 Python 代码), 用打包工具给可执行文件追加 AFBNDL01 footer(参考 go-sdk)。 否 → 能否通过单条 shell 命令启动(node、ruby、dotnet…)? 是 → 继承 SubprocessCoordinator,只需实现 _build_execute_task_command。 否 → 直接继承 BaseCoordinator,从零实现 execute_task (少见:gRPC 守护进程、共享内存、常驻进程等场景)。

SubprocessCoordinator:只需实现一个方法

SubprocessCoordinator 接管了完整的子进程生命周期。任务被触发时它会:

  1. 127.0.0.1上绑定两个临时 TCP 服务 socket;
  2. 启动子进程,并在命令行末尾追加--comm=<host>:<port>--logs=<host>:<port>
  3. 等待子进程连接这两个 socket;
  4. 向子进程发信号,开始执行用户代码;
  5. --logs通道收到的任务日志行转发到 Airflow 日志基础设施;
  6. 子进程退出或启动超时时拆除一切资源。

子类唯一要实现的是:

def _build_execute_task_command(self, *, what: TaskInstanceDTO) -> tuple[list[str], str]: ...

返回(command, subprocess_schema_version)二元组:

  • command:子进程的 argv 列表。不要包含--comm/--logs——基类在绑定 socket 之后才会追加这两个参数;
  • subprocess_schema_version:子进程所理解的 wire-schema 版本(YYYY-MM-DD日期串),supervisor 用它跨 SDK 版本协商消息格式。

从源码结构看,_subprocess.py中还内建了两个容易被忽视的健壮性设计:_socket_address()会把双栈 JVM 的 loopback 连接(::ffff:127.0.0.1::127.0.0.1两种形式)都归一化为127.0.0.1,否则 Java 任务会因归属校验失败而被拒绝;_connection_owned_by_process_tree()则用psutil枚举整个子进程树,确认回连的对端确实属于被启动的进程树(允许 JVM launcher、shell 包装器 fork 出的后代进程回连),而非同机抢端口的无关进程。这正是文档强调"SDK 必须从同一进程或其子进程连接"这一约束的底层实现。

参考实现对照表

技能文档汇总了必须研读的参考实现:

研究对象位置
SubprocessCoordinator 基类task-sdk/src/airflow/sdk/coordinators/_subprocess.py
Java coordinator(SubprocessCoordinator 子类)task-sdk/src/airflow/sdk/coordinators/java/coordinator.py
ExecutableCoordinator(原生 bundle)task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py
Kotlin 侧 wire protocoljava-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/
Go 侧 wire protocolgo-sdk/pkg/execution/
AFBNDL01 footer(Go 参考)go-sdk/internal/bundlefooter/、task-sdk/docs/executable-bundle-spec.rst
全部消息类型与字段定义task-sdk/src/airflow/sdk/execution_time/schema/schema.json

Java 与 Go 是两个生产级参考点,按目标语言运行时模型就近学习:JVM/解释型看 Java,原生/编译型看 Go。以 Java coordinator 为例,它从 JAR 的META-INF/MANIFEST.MF中解析Main-ClassAirflow-Supervisor-Schema-Version两项元数据来定位可执行包——_build_execute_task_command的实现要点就藏在这里。

四、Supervisor Schema:协议契约与版本协商

Supervisor Schema 是 supervisor 与语言 SDK 子进程之间的正式契约:它由 supervisor 中定义的 Pydantic 模型生成,发布为 schema.json。该文件描述了 comm socket 双向所有消息类型的名称、字段与约束,并携带YYYY-MM-DD格式的api_version字段标识 schema 修订版本。

由于 schema 会随版本演进,SDK 必须向 supervisor 声明自己构建时对应的 API 版本,协商才能成立。目标语言 SDK 可以用支持 JSON Schema 输入的代码生成器直接从schema.json生成消息类型。一个现成的版本声明实例是 java-sdk/capabilities.yaml,其中supervisor_schema_version: "2026-06-16"gradle.properties中打入 JAR manifest 的airflowSupervisorSchemaVersion保持同步——渲染 hook 会在两者不一致时使构建失败。

五、AFBNDL01:原生可执行包格式

若目标运行时编译为自包含原生可执行文件,ExecutableCoordinator可以零 Python 代码地发现并启动它,前提是在构建阶段由自定义打包步骤给可执行文件追加AFBNDL01元数据 trailer(规范见 task-sdk/docs/executable-bundle-spec.rst)。

从 executable/coordinator.py 源码可以看到 trailer 的具体形态:固定 64 字节(FOOTER_SIZE = 64),文件末尾 8 字节为魔数AFBNDL01;开头三个小端 uint32 分别记录 source 区长度、metadata 区长度和 footer 版本号(当前仅支持 version 1);随后 32 字节是二进制的 SHA-256 校验值,另有 12 字节必须全零的保留区。解析时会严格校验各区域偏移,若声明的区域越过文件头、或二进制区为空则报错。Go 侧的打包工具在 go-sdk/internal/bundlefooter/footer.go 中有对应实现,可作为写新语言打包器的直接参照。

六、Wire Protocol:长度前缀 MessagePack 帧

当 coordinator 是SubprocessCoordinator(含其子类ExecutableCoordinator)时,--commsocket 上的所有通信使用长度前缀 MessagePack 帧

[4-byte big-endian uint32: payload 长度][payload 字节]

payload 是 MessagePack 编码的数组,两种形态:

  • SDK → Supervisor:2 元素数组[id, body]id是整型,在该连接内唯一标识一次请求;响应会回显同一个id以便关联。
  • Supervisor → SDK:3 元素数组[id, body, error]body在出错时为nullerror在失败时是ErrorResponsemap,成功时为null

payload 本身是带"type"键的 MessagePack map,用于指明消息类型;单帧最大2³² - 1字节,超过即为错误。帧层的具体编解码可参考 Go 的 frames.go 与 Kotlin 的 MsgPack.kt。

请求/响应关联:SDK 发出的每个请求携带单调递增的整数id,supervisor 在响应帧中回显该id。当存在并发执行的任务、或单个任务同时发出多个请求时,SDK 必须按id关联响应。

启动时序:supervisor 用--comm=<host>:<port>--logs=<host>:<port>两个参数拉起子进程,SDK 必须解析参数并尽快连接两个 socket(supervisor 会校验回连方属于该进程树)。两条连接建立后,supervisor 在 comm socket 上发送StartupDetails消息以启动执行。

七、日志通道:不丢"早到"日志的 NDJSON 规范

子进程相对--logssocket 的完整生命周期是:

1. Supervisor 以 --comm=… --logs=… 启动子进程 └─► 2. 子进程启动。其 logger 已经激活,但 --logs socket 尚未连接 └─► 3. 子进程连接 --comm 与 --logs └─► 4. 此后的每条记录经 --logs socket 发送

关键在第 2 阶段的"空窗期":SDK 自身的启动代码(参数解析、连接调用)在--logssocket 存在之前就可能已经产生日志记录。这些记录绝不能丢弃——应在内存中缓冲,socket 连上后按序 flush,之后再发送后续记录。Go SDK 与 Java SDK 均已实现该缓冲。

日志消息为换行分隔 JSON(NDJSON):每条日志是一行 UTF-8 JSON 对象,以\n结尾,采用 structlog 风格事件格式:

{"event": "Starting extraction", "level": "info", "logger": "com.example.SalesPipeline", "timestamp": "2026-06-22T12:00:00", "rows": 42}
  • event:日志消息本体;
  • level:小写的级别名,支持criticalerrorwarninginfodebugnotset,其他级别名会被丢弃;
  • timestamp:ISO-8601 时间戳;
  • 其余键作为结构化字段随记录转发。

技能文档补充了若干跨语言通用细节(权威定义在 30_new_language_sdk.rst 的 Logging 小节):

  • 级别数值:沿用 Pythonlogging量表,CRITICAL=50ERROR=40WARNING=30INFO=20DEBUG=10NOTSET=0,与 Airflow 其余部分阈值对齐;各语言应实现相应转换逻辑。
  • NAMESPACE_LEVELS解析:先按[\s,]+切分,再把每项按=拆为(logger_name, level_name);仅当记录级别>=匹配 logger 的阈值(无匹配时用全局阈值)时才发送。
  • 不要丢"晚到"日志:尽早连接--logssocket,并保持打开直到--comm通道结束,否则拆除阶段的记录会丢失。
  • 额外配置:运行时读不到 Airflow 配置文件,若 SDK 还需要[logging]中的其他设置,应从 coordinator 的start中像上面两个变量一样以环境变量传递。

过滤责任在 SDK 侧:supervisor 不对--logssocket 做过滤,而是通过环境变量提供信息——SubprocessCoordinator(含ExecutableCoordinator)启动外部运行时时会设置AIRFLOW__LOGGING__LOGGING_LEVEL(来自[logging] logging_level,如INFO)与AIRFLOW__LOGGING__NAMESPACE_LEVELS(来自[logging] namespace_levels,如sqlalchemy=INFO, botocore=WARNING)。SDK必须在启动时读取它们,并在发送前丢弃低于适用阈值的记录,使过滤结果与同一部署中的 Python 任务一致。

八、错误处理与 TaskInstance 状态符合性

错误处理四条规则(源自 30_new_language_sdk.rst 的 Error handling 小节):

  • 任务抛出未处理错误时,SDK 必须在关闭 comm socket之前发送"state": "failed"TaskState消息;
  • 任务失败但仍有重试次数时,必须改发RetryTask让 supervisor 把任务转入up_for_retry。字段名必须与 Supervisor Schema 完全一致——失败详情键是retry_reason而非reason
  • 进程未发送终结消息就退出时,supervisor 依据异常退出把任务实例标记为failed,但任务日志可能不完整;
  • supervisor 对任务中途的请求返回ErrorResponse时,SDK 应把它作为错误传播到任务函数。

状态符合性采用 RFC 2119 的 MUST/SHOULD/MAY 分级,各 SDK 只声明自己实际支持的子集:

状态级别上报方式 / 备注
successMUSTSucceedTask,或进程干净退出(exit 0)
failedMUSTTaskState"state": "failed";无终结消息时由非零退出码推断
up_for_retryMUST失败且尚有重试时发RetryTask;详情键为retry_reason
skippedSHOULDTaskState"state": "skipped",用于分支/跳过语义
deferredMAYDeferTask,需 SDK 桥接到 triggerer
up_for_rescheduleMAYRescheduleTask,用于 reschedule 模式 sensor
awaiting_inputMAYAwaitInputTask,用于 human-in-the-loop
removedMAYTaskState"state": "removed"

只实现 MUST 层的 SDK 已能让普通任务带重试地跑通成功/失败;SHOULD/MAY 层解锁分支、延迟执行、reschedule sensor 与人工介入。调度器状态(queuedscheduledrunningrestartingupstream_failed)由 Airflow 侧设置,不属于 SDK 符合性范围。

九、运行时能力与 Native-Dag 能力声明

运行能力描述任务体在执行期间能做什么(无论任务声明在 Python Dag 的@task.stub里还是原生 Dag 中),逐项独立声明:

  • MUSTmixed-lang-stub-target(执行 Python Dag 中@task.stub声明的任务,是每门语言 SDK 的主执行路径)、task-logging(经--logs转发 stdout/stderr 与结构化日志;远端日志存储由 supervisor 统一处理)、xcom-read-write(跨 Python 边界读写 XCom)、connection-read(按 id 解析 Connection)、variable-read-write(读写删 Variable)、self-contained-bundle(构建产物把 Airflow 元数据dag_id/task_id等与任务代码嵌入同一交付物——Go 用AFBNDL01trailer、JVM 嵌入 jar、Node 嵌入 package);
  • MAYretry-policytask-state-storeasset-state-storeasset-event-emitasset-event-read

Native-Dag authoring(整个 Dag 用目标语言书写,无需 Python 文件)由总括能力native-dag-authoring(SHOULD)控制。项目目标是每门语言 SDK 都达到这一标准,但仅执行@task.stub的 SDK 仍然有用且符合 MUST 层。在其前提下还有一组条件能力(记作,未支持原生 Dag 时为 n/a 而非"未支持"):task-args(MUST†)、dag-params(MUST†)、taskflow-dependencies(MUST†)、branching(SHOULD†)、dag-test(SHOULD†,airflow dags test本地演练)、task-group(MAY†)、dynamic-task-mapping(MAY†)、asset-inlets-outlets(MAY†)、asset-scheduling(MAY†)、object-store(MAY†)。

这些维度除了文档中的文字描述,还需以机器可读方式声明:每个 SDK 手写一份<sdk>/capabilities.yaml,由 prek hook 生成发布的兼容矩阵表格。以 Java 为例:

java-sdk/capabilities.yaml <- 唯一需要手改的文件 | | hook: update-java-sdk-readme-matrix | +--> java-sdk/README.md (面向贡献者) +--> java-sdk/sdk/module.md (Dokka -> 发布的 API 参考)

hook 会重写目标文件,发现表格过期即非零退出,漂移的矩阵会让构建失败。manifest 应放在 SDK 发布物之外(Java 即位于settings.gradle.kts之上、覆盖所有子项目)。新增 SDK 时,在 scripts/ci/prek/lang_sdk_compat_matrix.py 的LANG_SDKS中注册、按同 schema 编写capabilities.yaml并添加等价 hook;新增/重命名维度时须在同一 PR 中同时修改该文件的STATE_DIMENSIONS/CAPABILITY_DIMENSIONS与上文文字描述——渲染器会对每份 manifest 与该列表做校验。

十、E2E 测试套件:镜像 java/go 的参考结构

airflow-e2e-tests下新增两个文件,镜像现有 java_sdk_tests/ 或 go_sdk_tests/ 的布局:

airflow-e2e-tests/tests/airflow_e2e_tests/<language>_sdk_tests/ __init__.py test_<language>_sdk_dag.py

测试文件应完成以下断言链:

  1. 通过AirflowClient.trigger_dag触发 SDK 的示例 Dag(放在<language>-sdk/dags/或等价位置);
  2. AirflowClient.wait_for_dag_run等待运行结束;
  3. 断言每个 SDK 任务实例到达"success"
  4. 断言至少一个 XCom 值——确认从任务返回值经 supervisor 到 XCom 存储的完整往返;
  5. 若 SDK 产生结构化日志则断言其内容(Go 测试套件即日志内容断言的范例)。

本地运行方式:

E2E_TEST_MODE=<language>_sdk uv run --project airflow-e2e-tests pytest \ tests/airflow_e2e_tests/<language>_sdk_tests/ -xvs

Java(java_sdk_tests/test_java_sdk_dag.py)与 Go(go_sdk_tests/test_go_sdk_dag.py)套件是参考实现。此外按 30_new_language_sdk.rst 的 Testing 小节,实现还应包含:帧层单元测试(编解码往返、超大帧、损坏的长度前缀)、每个消息类型双向的单元测试,以及使用 Breeze 对接真实 supervisor 的集成测试。

十一、PR 清单:技能文档独有的六项

30_new_language_sdk.rst 覆盖了 coordinator 位置、wire protocol 实现与测试要求;技能文档要求同一 PR 内再补齐这些条目:

  1. task-sdk/src/airflow/sdk/coordinators/<language>/__init__.py——简短模块 docstring 与 coordinator 类的__all__re-export;
  2. airflow-core/docs/authoring-and-scheduling/language-sdks/<language>.rst——面向用户的文档,结构仿照现有的 java.rst 或 go.rst;
  3. airflow-core/docs/authoring-and-scheduling/language-sdks/index.rst——把新文档加入 toctree;
  4. airflow-core/newsfragments/<PR>.feature.rst——新语言始终对用户可见,必须加 newsfragment;
  5. CI 接线——检查 dev/breeze/src/airflow_breeze/utils/selective_checks.py,确认<language>-sdk/的变更能触发正确的测试组,缺失则补上;
  6. E2E 测试——在airflow-e2e-tests/tests/airflow_e2e_tests/下新增<language>_sdk_tests/套件(见上一节)。

十二、落地路径小结

把上面的要素串起来,一次完整的"新语言 SDK"贡献大致按如下顺序推进:先读 30_new_language_sdk.rst 确定 coordinator 选型与 socket 生命周期 → 按task-sdk/src/airflow/sdk/coordinators/<language>/落 coordinator 代码(多数场景只需实现_build_execute_task_command并返回YYYY-MM-DD的 schema 版本)→ 在目标语言中按 schema.json 生成消息类型,实现长度前缀 MessagePack 帧、日志缓冲与 NDJSON 过滤 → 若为原生可执行文件,用打包器追加AFBNDL01footer(参考 go-sdk/internal/bundlefooter/footer.go)→ 编写capabilities.yaml并注册兼容矩阵 → 搭建单测/集成/E2E 三层测试 → 按六项 PR 清单补齐文档、newsfragment 与 CI 接线。整个过程中,Java 与 Go 两个生产实现始终是最直接的对标物。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/6 19:00:02

基于小程序的专业必修课程在线学习系统设计与实现

1. 项目背景与意义随着移动互联网的普及和高校信息化建设的不断深入&#xff0c;传统课堂教学模式在时间、空间上存在一定局限&#xff0c;学生课后复习、自主学习的需求日益增长。微信小程序凭借其无需下载安装、即用即走、跨平台兼容等优势&#xff0c;成为高校在线学习平台的…

作者头像 李华
网站建设 2026/9/6 19:00:00

物流仿真系统实验指南:从业务建模到数据决策

简介&#xff1a;物流仿真系统实验PDF是物流管理专业实践课程的完整指导资料&#xff0c;主要面向需要掌握RaLC乐龙仿真软件的学生与教师。内容分三篇&#xff0c;按基础到高级安排8个实验&#xff1a;从通过型物流中心的分拣分流模拟、仓储型物流中心建模&#xff0c;到复合型…

作者头像 李华
网站建设 2026/9/6 18:58:56

云效+Kubernetes:构建自动化CI/CD流水线的DevOps实践指南

简介&#xff1a;基于Kubernetes的DevOps工作流专题PDF&#xff0c;面向云原生与DevOps工程师&#xff0c;尤其适合具备一定容器基础、正在规划或希望搭建持续交付链路的开发与运维人员。内容从DevOps体系演进切入&#xff0c;结合阿里云效平台实践&#xff0c;梳理了需求、开发…

作者头像 李华
网站建设 2026/9/6 18:58:27

品质意识培训:从认知到行动的质量管理第一课

简介&#xff1a;《品质管理讲座之一&#xff1a;品质意识培训》是一份面向企业培训师、质量管理人员和生产一线主管的PPT课件&#xff0c;聚焦品质意识这一质量管理基础环节&#xff0c;可用于内部培训、班前会宣讲或质量课程备课。课件共87页&#xff0c;以单个pptx文件提供&…

作者头像 李华
网站建设 2026/9/6 18:56:51

纵横公路软件实操指南:从新建项目到报表导出的完整流程

简介&#xff1a;《纵横公路软件操作手册》是一本面向公路建设领域造价人员的实用指南&#xff0c;系统讲解纵横公路造价软件从项目建立到报表输出的完整操作流程。内容涵盖新建建设项目与造价文件、按标准项目表搭建概预算结构、套用定额计算建筑安装工程费、录入工料机预算单…

作者头像 李华
网站建设 2026/9/6 18:55:40

IEC 61215中文版标准深度解读:从PDF到产线上的可靠性与测试实战

简介&#xff1a;这是一份IEC 61215标准中文版PDF资料&#xff0c;主要面向光伏组件研发、质量检测与认证工程师&#xff0c;用于指导光伏组件在电性能、热性能、绝缘性、湿漏电等方面的型式试验与可靠性验证。资源为单个PDF文档&#xff0c;大小238KB&#xff0c;内容涵盖标准…

作者头像 李华