Apache Airflow Java SDK 贡献者指南:双栈架构、Bundle 扫描与 Java–Python 协调器实现
2026/9/6 17:59:58 网站建设 项目流程

Apache Airflow Java SDK 贡献者指南:双栈架构、Bundle 扫描与 Java–Python 协调器实现

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

本文基于 Airflow 仓库中的 Java SDK 贡献者技能文档(SKILL.md)展开,系统讲解 Airflow Java SDK(AIP-108)的工程结构:从org.apache.airflow.sdk公开 API 与execution内部实现包的可见性边界,到 Bundle 目录扫描与 JAR 清单属性发现机制,再到 Python 协调器(JavaCoordinator)如何组装命令行、与 JVM 子进程完成双 Socket 握手,以及 Supervisor 线协议(4 字节长度前缀 + MessagePack)的升级流程。读完本文,你将掌握在java-sdk/与 Python 协调器两侧进行功能开发、测试与协议升级所需的完整知识链路。

Java SDK 的两个代码位置

Java SDK 让 Airflow 任务以 JVM 语言(Java、Kotlin 或任意 JVM 语言)执行:DAG 与调度仍留在 Python 侧,单个任务实例由JavaCoordinator派生的 JVM 子进程执行。贡献工作只发生在两个位置:

  • java-sdk/— JVM 侧库(Kotlin 源码实现,发布到 Maven)。SDK 与运行时逻辑用 Kotlin 编写,Java 是公开 API 的目标语言而非实现语言,示例(AnnotationExample.java)展示了 Java 侧的用法。
  • task-sdk/src/airflow/sdk/coordinators/java/— 负责启动 JVM 子进程的 Python 协调器(coordinator.py)。

开发前应尽早通读两份权威参考文档:

  • airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst — 面向用户的指南:注解式(@Builder.Dag/@Builder.Task)与接口式(BundleBuilder)API、XCom 类型映射、Gradle/Maven 步骤、协调器配置。
  • java-sdk/README.md — 贡献者指南:仓库布局、执行流程详解、Gradle + Breeze 测试命令、编码规范、常见任务与 PR 清单。

SDK 包架构:公开 API 与内部实现的可见性边界

JVM 侧库拆分为两个包,遵循严格的可见性规则:

  • org.apache.airflow.sdk— 公开的、面向用户的 API。这里的类(如ClientBundleBundleBuilderServer)是 DAG 作者和任务实现者直接导入的稳定契约,对该包的任何变更都视为破坏性变更
  • org.apache.airflow.sdk.execution— 内部实现细节。该包下的一切(CoordinatorCommLogSenderLogexecution下的Client、由 schema 生成的模型等)都不应被用户导入,可以在版本间无通知地变更

审查或编写代码时要守住这条边界:用户任务代码与BundleBuilder子类只能从org.apache.airflow.sdk导入;任何在用户可见 API 表面出现org.apache.airflow.sdk.execution.*导入的行为都是危险信号(red flag)。

从源码可以看到这条边界的实际落地:公开 API 的 Client.kt 是一个internal constructor的薄封装,持有StartupDetails和内部实现execution.Client,公开方法getConnectiongetVariablegetXComsetXCom全部委托给内部impl完成 Supervisor 线调用。例如setXCom会自动带上当前任务实例的dagId/taskId/runId/mapIndex,XCom 默认键为return_value(Client.kt)。

Bundle 组成与协调器发现机制

一个bundlejars_root目录(通常为build/bundle/)下的一个 JAR 文件集合。协调器在任务派发时扫描该目录,从中找到两个关键信息:

  1. Main-Class(标准 JAR 清单属性)— 入口点的全限定类名,协调器将以java -classpath … <Main-Class> --comm … --logs …方式调用它。该入口类必须具备public static void main(String[] args)方法;Gradle 插件org.apache.airflow.sdk会依据airflowBundle { mainClass = "…" }自动写入该属性,并在构建时校验类存在且签名正确。
  2. Airflow-Supervisor-Schema-Version(Airflow 专属清单属性)— JVM 侧与 Python supervisor 通信所期望的线协议版本。在 fat-JAR 模式(默认)下,Gradle 插件从runtimeClasspath中的airflow-sdkJAR 读取该值并复制进 shadow JAR 的清单;在 thin-JAR 模式(fatJar = false)下,该值留在与 bundle JAR 一同部署的airflow-sdkJAR 中。

Python 侧的扫描实现

Python 协调器JavaCoordinator通过_JarInfo.find()扫描jars_root下的每一个 JAR,从每个 ZIP 中读出META-INF/MANIFEST.MF,收集其中携带Main-ClassAirflow-Supervisor-Schema-Version的 JAR。相关实现在 coordinator.py:

  • _find_jars()递归遍历jars_root目录,并用(st_dev, st_ino)去重目录,防止符号链接环路或硬链接导致的无限递归;
  • _JarMetadata.from_jar()解析 JAR 清单,缺少清单或 JAR 损坏(BadZipFile)的文件会被记录日志后忽略;
  • _JarInfo.find()汇总扫描进度:找到匹配Main-ClassAirflow-Supervisor-Schema-Version的 JAR 后组合成_JarInfo返回。

发现规则与失败行为:

  • 若在JavaCoordinator实例上通过[sdk] coordinators的 kwargs 显式设置了main_class,扫描会以它作为过滤器;否则第一个Main-Class属性的 JAR 胜出(源码文档提示:存在多个可执行 JAR 时行为可能不确定)。
  • 无论如何,jars_root中至少要有一个JAR 携带Airflow-Supervisor-Schema-Version,否则启动失败,抛出FileNotFoundError,错误信息会区分"找不到带 Main-Class 的 JAR"与"找不到带 schema 版本元数据的 JAR"两种情形(coordinator.py)。
  • classpath 由_calculate_classpath()生成:所有 JAR 路径排序后以os.pathsep连接,保证输出确定性(coordinator.py)。

任务执行全链路:从 JavaCoordinator 到 JVM 子进程

JavaCoordinator继承自SubprocessCoordinator(coordinator.py),其配置来自[sdk] coordinators条目,支持以下参数(源码字段与文档说明):

参数默认值说明
java_executable"java"(依赖$PATHjava可执行文件路径,如/usr/lib/jvm/java-17-openjdk/bin/java
jvm_args空列表传给 JVM 的额外参数,如["-Xmx1024m"]
jars_root必填(至少 1 项)扫描 JAR bundle 的目录列表
main_class""(自动发现)显式指定入口类
task_startup_timeout10.0等待子进程连接两个 socket 的最长时间

典型的[sdk] coordinators配置(引自源码 docstring):

{ "jdk-17": { "classpath": "airflow.sdk.coordinators.java.JavaCoordinator", "kwargs": { "jars_root": ["~/airflow/jars"], "java_executable": "/usr/lib/jvm/java-17-openjdk/bin/java", "jvm_args": ["-Xmx1024m"] } } }

命令行组装

子类唯一必须实现的方法是_build_execute_task_command,它返回(argv, schema_version)二元组。Java 实现(coordinator.py)如下:

def _build_execute_task_command(self, *, what: TaskInstance) -> tuple[list[str], str | None]: jar = _JarInfo.find(self.jars_root, self.main_class) command = [ self.java_executable, "-classpath", _calculate_classpath(self.jars_root), *self.jvm_args, jar.main_class, ] return command, jar.schema_version

解析出的schema_version作为返回值交给基类SubprocessCoordinator,用于协商 supervisor 线协议;贡献规范明确要求:不要通过这一方法之外的方式从 Python 深入 JVM 进程内部(Do not reach into the JVM process from Python beyond what this method provides)。

SubprocessCoordinator:双 Socket 握手

基类(_subprocess.py)处理其余全部 socket 生命周期:监听、派生子进程、接受连接、排空启动期输出、失败时拆除资源。关键机制包括:

  • _PopenActivitySubprocess.start()先在127.0.0.1上绑定两个临时 socket(comm 与 logs),再以追加--comm=<host:port>--logs=<host:port>参数的方式启动子进程(_subprocess.py);
  • _accept_connections()阻塞等待子进程连上两个端口,期间用 selectors 排空子进程 stdout/stderr,防止管道阻塞(_subprocess.py);
  • 每个被接受的连接都会经过进程树归属校验:通过 psutil 检查连接是否属于子进程或其后代(JVM 启动器可能 fork 出真正回连的 worker),并处理双栈 JVM 的 IPv4-mapped/IPv4-compatible 地址规范化,避免 Java 任务被误拒(_subprocess.py);
  • 超时(默认task_startup_timeout10 秒)未连接则抛TimeoutError,子进程提前退出则抛RuntimeError

JVM 侧:Server 驱动执行循环

JVM 侧入口是 Server.kt,典型用法:

public static void main(String[] args) { Server.create(args).serve(new MyBundleBuilder().build()); }

serveAsync()并发打开 comm 与 logs 两个 TCP 连接(Server.kt);dispatchTask()读取第一帧:期望StartupDetails进入runTask,收到ErrorResponse则抛出ApiError(Server.kt)。

结合 java-sdk/README.md 的执行流程描述,完整链路为:

  1. JavaCoordinator.execute_task()(Python)扫描jars_root、组装 classpath,派生java -cp <jars> <MainClass> --comm=<host>:<port> --logs=<host>:<port>
  2. Server.kt启动后立即连接两个 socket;
  3. supervisor 发送StartupDetailsMessagePack 消息;JVM 读取后按dag_id+task_id查找对应任务并调用用户任务方法;
  4. 执行期间 JVM 向 supervisor 发起请求(GetVariableGetConnectionGetXComSetXCom等),supervisor 响应;所有帧均为 4 字节大端长度前缀 + MessagePack 载荷;
  5. 完成(或异常)后 JVM 发送TaskState消息并关闭 socket,进程退出。

SDK 自身产生的日志(而非用户代码)通过--logssocket 转发,由 supervisor 追加进 Airflow 日志存储。

线协议:4 字节长度前缀 + MessagePack

线协议定义在 task-sdk/src/airflow/sdk/execution_time/schema/schema.json,是两侧共享的单一事实来源。JVM 侧的帧层由 execution/Comm.kt 与 execution/Frame.kt 实现。从CoordinatorComm的源码结构看,帧层还带有防御性设计:入站帧大小上限取Runtime.getRuntime().maxMemory()的 1/8(受-Xmx影响)并与协议层的Frame.MAX_FRAME_LENGTH取较小值,把潜在 OOM 降级为可捕获的FrameProcessingException(Comm.kt);写入侧通过互斥锁保证并发 client 调用的帧完整性——示例中的concurrentClientCalls任务正是用 8 线程 32 次并发getConnection来验证这一点(AnnotationExample.java)。

新增消息类型需要同时改动两侧:schema.json(Python 侧)与execution/Comm.kt+execution/Client.kt(JVM 侧)。

升级 Supervisor Schema 客户端

升级到更新的 Supervisor Schema 版本时的标准步骤:

  1. 重新生成模型:./gradlew generateJsonSchema2Pojo
  2. 修改execution/Client.kt以处理变更。

java-sdk/README.md 的Contributing一节逐步展开了"新增一个 Client 方法"的完整序列:重新生成 POJO → 在execution/Comm.kt或新文件中添加 Kotlin 请求/响应数据类 → 在公开Client.kt添加面向用户的方法并委托execution/Client.kt→ 在sdk/src/test/kotlin/.../ClientTest.kt编写 mock socket 层的单测 → 若用户可见则更新 java.rst。

关键文件速查

文件用途
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Client.kt公开 API(Variables、Connections、XCom)
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Client.ktSupervisor 线调用
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Comm.kt4 字节前缀 MessagePack 帧层
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt进程入口;驱动执行循环
java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.ktkapt 注解处理器
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.ktGradle bundle 插件
task-sdk/src/airflow/sdk/coordinators/java/coordinator.pyPython 侧——派生 JVM 子进程
task-sdk/src/airflow/sdk/execution_time/schema/schema.json线协议定义(两侧共用)

运行测试

JVM 侧:一律使用java-sdk/目录内的./gradlew,绝不使用 apt 安装的gradle

cd java-sdk ./gradlew test # 运行单个测试类 ./gradlew :sdk:test --tests "org.apache.airflow.sdk.execution.CommTest"

Python 协调器侧:使用 Breeze(不要在宿主机直接跑pytest):

breeze testing task-sdk-tests -- task_sdk/coordinators/java

端到端测试套件(需要真实 Airflow 环境):

E2E_TEST_MODE=java_sdk uv run --project airflow-e2e-tests pytest \ tests/airflow_e2e_tests/java_sdk_tests/ -xvs

更新 Python 协调器

coordinator.py继承SubprocessCoordinator,子类唯一必须实现的方法就是_build_execute_task_command,返回(argv, schema_version)。参考现有实现了解jars_rootjava_executablejvm_argsmain_class如何被组装进命令(上文已给出完整实现)。再次强调规范:不要越过这个方法从 Python 深入 JVM 进程——socket 生命周期、连接校验、资源清理均由基类统一管理。

实战:端到端跑通示例 bundle

以下操作摘自 java-sdk/README.md 的Running the example与用户文档,可作为验证环境是否配置正确的标准流程。前置要求:worker 节点上 JRE 11+ 可用,apache-airflow-task-sdk(随 Airflow 安装)已提供协调器,无需额外 Python 包。

  1. 构建并发布 SDK 到本地 Maven 仓库:

    ./gradlew publishToMavenLocal -PskipSigning=true
  2. 打包示例 bundle 到./example/build/bundle(在example/目录下执行../gradlew bundle),并把带 stub 任务的 DAG(见 example/src/resources/dags)放到 Airflow 可发现的位置。DAG 侧使用queue="java"的 stub 任务:

    @dag def sales_pipeline(): @task.stub(queue="java") def extract(): ... @task.stub(queue="java") def transform(extracted): ...
  3. 配置 Airflow 将java队列的任务路由给 Java 协调器:

    export AIRFLOW__SDK__COORDINATORS='{ "java": { "classpath": "airflow.sdk.coordinators.java.JavaCoordinator", "kwargs": {"jars_root": ["/opt/airflow/java-sdk/example/build/bundle"]} } }' export AIRFLOW__SDK__QUEUE_TO_COORDINATOR='{"java": "java"}'
  4. 确保示例 DAG 所需的 Connection 与 Variable 可用:

    export AIRFLOW_CONN_TEST_HTTP='{ "conn_type": "http", "login": "user", "password": "pass", "host": "example.com", "port": 1234, "extra": {"param1": "val1", "param2": "val2"} }' export AIRFLOW_VAR_MY_VARIABLE=123

Java 侧实现通过@Builder.Dag/@Builder.Task注解声明任务,参数用@Builder.XCom(task = "...")接收上游 XCom,Client参数在任务执行时由 SDK 注入(详见 AnnotationExample.java)。

贡献清单与约定

按 java-sdk/README.md 的编码规范与 PR 清单,提交前应注意:

  • SDK 与 processor 源码全部为Kotlin;公开 API 表面(sdk/src/main/kotlin/顶层)保持干净,内部实现放在execution/子包;
  • 注解处理器使用kapt:新增注解需在Builder.kt定义、在BuilderProcessor.kt处理,并在processor/src/test/kotlin/添加 golden-output 测试;
  • 提交前运行./gradlew ktLintCheck spotlessCheck(或ktLintFormat spotlessApply);
  • 所有新文件需要 Apache License 头;
  • PR 检查项:跑./gradlew build test(JVM)与相应 pytest 套件(Python 协调器);确认示例 bundle 仍可编译;若schema.json有改动,验证 JVM 与 Python 两侧都能处理新增/变更字段;每个行为变更都需测试覆盖;对task-sdk/的用户可见变更在airflow-core/newsfragments/下添加 newsfragment。

此外,Java SDK 的能力边界(支持哪些 TaskInstance 状态与运行时能力)由 java-sdk/capabilities.yaml 声明并自动生成到 README 的兼容矩阵,例如当前声明variable-read-write仅支持getVariable、暂不支持经 comm socket 写入——在评估功能实现范围时应以该矩阵为准。

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

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询