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

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

2026-09-05 18:45:47作者:裘晴惠Vivianne

在 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 protocol java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/
Go 侧 wire protocol go-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)时,--comm socket 上的所有通信使用长度前缀 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 在失败时是 ErrorResponse map,成功时为 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 规范

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

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

关键在第 2 阶段的"空窗期":SDK 自身的启动代码(参数解析、连接调用)在 --logs socket 存在之前就可能已经产生日志记录。这些记录绝不能丢弃——应在内存中缓冲,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 小节):

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

过滤责任在 SDK 侧:supervisor 不对 --logs socket 做过滤,而是通过环境变量提供信息——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 只声明自己实际支持的子集:

状态 级别 上报方式 / 备注
success MUST SucceedTask,或进程干净退出(exit 0)
failed MUST TaskState"state": "failed";无终结消息时由非零退出码推断
up_for_retry MUST 失败且尚有重试时发 RetryTask;详情键为 retry_reason
skipped SHOULD TaskState"state": "skipped",用于分支/跳过语义
deferred MAY DeferTask,需 SDK 桥接到 triggerer
up_for_reschedule MAY RescheduleTask,用于 reschedule 模式 sensor
awaiting_input MAY AwaitInputTask,用于 human-in-the-loop
removed MAY TaskState"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 用 AFBNDL01 trailer、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.pyLANG_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.rstgo.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 过滤 → 若为原生可执行文件,用打包器追加 AFBNDL01 footer(参考 go-sdk/internal/bundlefooter/footer.go)→ 编写 capabilities.yaml 并注册兼容矩阵 → 搭建单测/集成/E2E 三层测试 → 按六项 PR 清单补齐文档、newsfragment 与 CI 接线。整个过程中,Java 与 Go 两个生产实现始终是最直接的对标物。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
898
5.82 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
921
1.84 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.8 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
531
596
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.02 K
519
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.36 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
391