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。这里的类(如Client、Bundle、BundleBuilder、Server)是 DAG 作者和任务实现者直接导入的稳定契约,对该包的任何变更都视为破坏性变更。org.apache.airflow.sdk.execution— 内部实现细节。该包下的一切(CoordinatorComm、LogSender、Log、execution下的Client、由 schema 生成的模型等)都不应被用户导入,可以在版本间无通知地变更。
审查或编写代码时要守住这条边界:用户任务代码与BundleBuilder子类只能从org.apache.airflow.sdk导入;任何在用户可见 API 表面出现org.apache.airflow.sdk.execution.*导入的行为都是危险信号(red flag)。
从源码可以看到这条边界的实际落地:公开 API 的 Client.kt 是一个internal constructor的薄封装,持有StartupDetails和内部实现execution.Client,公开方法getConnection、getVariable、getXCom、setXCom全部委托给内部impl完成 Supervisor 线调用。例如setXCom会自动带上当前任务实例的dagId/taskId/runId/mapIndex,XCom 默认键为return_value(Client.kt)。
Bundle 组成与协调器发现机制
一个bundle是jars_root目录(通常为build/bundle/)下的一个 JAR 文件集合。协调器在任务派发时扫描该目录,从中找到两个关键信息:
Main-Class(标准 JAR 清单属性)— 入口点的全限定类名,协调器将以java -classpath … <Main-Class> --comm … --logs …方式调用它。该入口类必须具备public static void main(String[] args)方法;Gradle 插件org.apache.airflow.sdk会依据airflowBundle { mainClass = "…" }自动写入该属性,并在构建时校验类存在且签名正确。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-Class和Airflow-Supervisor-Schema-Version的 JAR。相关实现在 coordinator.py:
_find_jars()递归遍历jars_root目录,并用(st_dev, st_ino)去重目录,防止符号链接环路或硬链接导致的无限递归;_JarMetadata.from_jar()解析 JAR 清单,缺少清单或 JAR 损坏(BadZipFile)的文件会被记录日志后忽略;_JarInfo.find()汇总扫描进度:找到匹配Main-Class与Airflow-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"(依赖$PATH) | java可执行文件路径,如/usr/lib/jvm/java-17-openjdk/bin/java |
jvm_args | 空列表 | 传给 JVM 的额外参数,如["-Xmx1024m"] |
jars_root | 必填(至少 1 项) | 扫描 JAR bundle 的目录列表 |
main_class | ""(自动发现) | 显式指定入口类 |
task_startup_timeout | 10.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 的执行流程描述,完整链路为:
JavaCoordinator.execute_task()(Python)扫描jars_root、组装 classpath,派生java -cp <jars> <MainClass> --comm=<host>:<port> --logs=<host>:<port>;Server.kt启动后立即连接两个 socket;- supervisor 发送
StartupDetailsMessagePack 消息;JVM 读取后按dag_id+task_id查找对应任务并调用用户任务方法; - 执行期间 JVM 向 supervisor 发起请求(
GetVariable、GetConnection、GetXCom、SetXCom等),supervisor 响应;所有帧均为 4 字节大端长度前缀 + MessagePack 载荷; - 完成(或异常)后 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 版本时的标准步骤:
- 重新生成模型:
./gradlew generateJsonSchema2Pojo; - 修改
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.kt | Supervisor 线调用 |
| java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Comm.kt | 4 字节前缀 MessagePack 帧层 |
| java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt | 进程入口;驱动执行循环 |
| java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt | kapt 注解处理器 |
| java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt | Gradle bundle 插件 |
| task-sdk/src/airflow/sdk/coordinators/java/coordinator.py | Python 侧——派生 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_root、java_executable、jvm_args、main_class如何被组装进命令(上文已给出完整实现)。再次强调规范:不要越过这个方法从 Python 深入 JVM 进程——socket 生命周期、连接校验、资源清理均由基类统一管理。
实战:端到端跑通示例 bundle
以下操作摘自 java-sdk/README.md 的Running the example与用户文档,可作为验证环境是否配置正确的标准流程。前置要求:worker 节点上 JRE 11+ 可用,apache-airflow-task-sdk(随 Airflow 安装)已提供协调器,无需额外 Python 包。
构建并发布 SDK 到本地 Maven 仓库:
./gradlew publishToMavenLocal -PskipSigning=true打包示例 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): ...配置 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"}'确保示例 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),仅供参考