
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 可以通过外部语言 SDKAIP-108执行非 Python 编写的任务代码。本文以仓库中 airflow-new-sdk 技能文档 为核心骨架结合其指认的权威贡献指南 Creating a new Language SDK 以及 Java/Go 两个生产级参考实现的源码系统讲解如何为 Airflow 引入一门新编程语言的任务执行能力——包括 Coordinator 基类选型、Supervisor Schema 版本协商、AFBNDL01 原生可执行包格式、MessagePack 线协议、日志通道规范、E2E 测试搭建以及 PR 提交清单帮助贡献者完整走通从目录规划到 CI 接线的落地全流程。一、两个组件Coordinator 与目标语言 SDKAirflow 要理解如何执行一门外语写的任务需要两个基本独立的组件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 的代码放哪里技能文档给出了统一的目录约定。CoordinatorPython 侧放入 task-sdktask-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 接管了完整的子进程生命周期。任务被触发时它会在127.0.0.1上绑定两个临时 TCP 服务 socket启动子进程并在命令行末尾追加--commhost:port与--logshost:port等待子进程连接这两个 socket向子进程发信号开始执行用户代码把--logs通道收到的任务日志行转发到 Airflow 日志基础设施子进程退出或启动超时时拆除一切资源。子类唯一要实现的是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.pyJava coordinatorSubprocessCoordinator 子类task-sdk/src/airflow/sdk/coordinators/java/coordinator.pyExecutableCoordinator原生 bundletask-sdk/src/airflow/sdk/coordinators/executable/coordinator.pyKotlin 侧 wire protocoljava-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Go 侧 wire protocolgo-sdk/pkg/execution/AFBNDL01 footerGo 参考go-sdk/internal/bundlefooter/、task-sdk/docs/executable-bundle-spec.rst全部消息类型与字段定义task-sdk/src/airflow/sdk/execution_time/schema/schema.jsonJava 与 Go 是两个生产级参考点按目标语言运行时模型就近学习JVM/解释型看 Java原生/编译型看 Go。以 Java coordinator 为例它从 JAR 的META-INF/MANIFEST.MF中解析Main-Class与Airflow-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 → Supervisor2 元素数组[id, body]。id是整型在该连接内唯一标识一次请求响应会回显同一个id以便关联。Supervisor → SDK3 元素数组[id, body, error]。body在出错时为nullerror在失败时是ErrorResponsemap成功时为null。payload 本身是带type键的 MessagePack map用于指明消息类型单帧最大2³² - 1字节超过即为错误。帧层的具体编解码可参考 Go 的 frames.go 与 Kotlin 的 MsgPack.kt。请求/响应关联SDK 发出的每个请求携带单调递增的整数idsupervisor 在响应帧中回显该id。当存在并发执行的任务、或单个任务同时发出多个请求时SDK 必须按id关联响应。启动时序supervisor 用--commhost:port与--logshost:port两个参数拉起子进程SDK 必须解析参数并尽快连接两个 socketsupervisor 会校验回连方属于该进程树。两条连接建立后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 均已实现该缓冲。日志消息为换行分隔 JSONNDJSON每条日志是一行 UTF-8 JSON 对象以\n结尾采用 structlog 风格事件格式{event: Starting extraction, level: info, logger: com.example.SalesPipeline, timestamp: 2026-06-22T12:00:00, rows: 42}event日志消息本体level小写的级别名支持critical、error、warning、info、debug、notset其他级别名会被丢弃timestampISO-8601 时间戳其余键作为结构化字段随记录转发。技能文档补充了若干跨语言通用细节权威定义在 30_new_language_sdk.rst 的 Logging 小节级别数值沿用 Pythonlogging量表CRITICAL50、ERROR40、WARNING30、INFO20、DEBUG10、NOTSET0与 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如sqlalchemyINFO, botocoreWARNING。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 0failedMUSTTaskState且state: failed无终结消息时由非零退出码推断up_for_retryMUST失败且尚有重试时发RetryTask详情键为retry_reasonskippedSHOULDTaskState且state: skipped用于分支/跳过语义deferredMAYDeferTask需 SDK 桥接到 triggererup_for_rescheduleMAYRescheduleTask用于 reschedule 模式 sensorawaiting_inputMAYAwaitInputTask用于 human-in-the-loopremovedMAYTaskState且state: removed只实现 MUST 层的 SDK 已能让普通任务带重试地跑通成功/失败SHOULD/MAY 层解锁分支、延迟执行、reschedule sensor 与人工介入。调度器状态queued、scheduled、running、restarting、upstream_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 嵌入 packageMAYretry-policy、task-state-store、asset-state-store、asset-event-emit、asset-event-read。Native-Dag authoring整个 Dag 用目标语言书写无需 Python 文件由总括能力native-dag-authoringSHOULD控制。项目目标是每门语言 SDK 都达到这一标准但仅执行task.stub的 SDK 仍然有用且符合 MUST 层。在其前提下还有一组条件能力记作†未支持原生 Dag 时为 n/a 而非未支持task-argsMUST†、dag-paramsMUST†、taskflow-dependenciesMUST†、branchingSHOULD†、dag-testSHOULD†airflow dags test本地演练、task-groupMAY†、dynamic-task-mappingMAY†、asset-inlets-outletsMAY†、asset-schedulingMAY†、object-storeMAY†。这些维度除了文档中的文字描述还需以机器可读方式声明每个 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测试文件应完成以下断言链通过AirflowClient.trigger_dag触发 SDK 的示例 Dag放在language-sdk/dags/或等价位置用AirflowClient.wait_for_dag_run等待运行结束断言每个 SDK 任务实例到达success断言至少一个 XCom 值——确认从任务返回值经 supervisor 到 XCom 存储的完整往返若 SDK 产生结构化日志则断言其内容Go 测试套件即日志内容断言的范例。本地运行方式E2E_TEST_MODElanguage_sdk uv run --project airflow-e2e-tests pytest \ tests/airflow_e2e_tests/language_sdk_tests/ -xvsJavajava_sdk_tests/test_java_sdk_dag.py与 Gogo_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 内再补齐这些条目task-sdk/src/airflow/sdk/coordinators/language/__init__.py——简短模块 docstring 与 coordinator 类的__all__re-exportairflow-core/docs/authoring-and-scheduling/language-sdks/language.rst——面向用户的文档结构仿照现有的 java.rst 或 go.rstairflow-core/docs/authoring-and-scheduling/language-sdks/index.rst——把新文档加入 toctreeairflow-core/newsfragments/PR.feature.rst——新语言始终对用户可见必须加 newsfragmentCI 接线——检查 dev/breeze/src/airflow_breeze/utils/selective_checks.py确认language-sdk/的变更能触发正确的测试组缺失则补上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),仅供参考