Apache Airflow Java SDK 贡献者指南:双栈架构、Bundle 扫描与 Java–Python 协调器实现
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 SDKAIP-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与接口式BundleBuilderAPI、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/mapIndexXCom 默认键为return_valueClient.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-VersionAirflow 专属清单属性— 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继承自SubprocessCoordinatorcoordinator.py其配置来自[sdk] coordinators条目支持以下参数源码字段与文档说明参数默认值说明java_executablejava依赖$PATHjava可执行文件路径如/usr/lib/jvm/java-17-openjdk/bin/javajvm_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上绑定两个临时 socketcomm 与 logs再以追加--commhost:port和--logshost: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.ktdispatchTask()读取第一帧期望StartupDetails进入runTask收到ErrorResponse则抛出ApiErrorServer.kt。结合 java-sdk/README.md 的执行流程描述完整链路为JavaCoordinator.execute_task()Python扫描jars_root、组装 classpath派生java -cp jars MainClass --commhost:port --logshost:portServer.kt启动后立即连接两个 socketsupervisor 发送StartupDetailsMessagePack 消息JVM 读取后按dag_idtask_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 降级为可捕获的FrameProcessingExceptionComm.kt写入侧通过互斥锁保证并发 client 调用的帧完整性——示例中的concurrentClientCalls任务正是用 8 线程 32 次并发getConnection来验证这一点AnnotationExample.java。新增消息类型需要同时改动两侧schema.jsonPython 侧与execution/Comm.ktexecution/Client.ktJVM 侧。升级 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公开 APIVariables、Connections、XComjava-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 安装的gradlecd java-sdk ./gradlew test # 运行单个测试类 ./gradlew :sdk:test --tests org.apache.airflow.sdk.execution.CommTestPython 协调器侧使用 Breeze不要在宿主机直接跑pytestbreeze testing task-sdk-tests -- task_sdk/coordinators/java端到端测试套件需要真实 Airflow 环境E2E_TEST_MODEjava_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 -PskipSigningtrue打包示例 bundle 到./example/build/bundle在example/目录下执行../gradlew bundle并把带 stub 任务的 DAG见 example/src/resources/dags放到 Airflow 可发现的位置。DAG 侧使用queuejava的 stub 任务dag def sales_pipeline(): task.stub(queuejava) def extract(): ... task.stub(queuejava) 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_VARIABLE123Java 侧实现通过Builder.Dag/Builder.Task注解声明任务参数用Builder.XCom(task ...)接收上游 XComClient参数在任务执行时由 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 testJVM与相应 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),仅供参考