Apache Airflow Java SDK 开发与运行指南:从 JVM 工作流 Bundle 到发布全流程
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow Java SDK 是 Airflow 为 JVM 生态提供的语言 SDK,允许开发者用 Java 或其他 JVM 兼容语言编写工作流任务,并以可被 Airflow 消费的 Bundle 形式交付执行。本文以java-sdk/README.md为骨架,结合仓库中的源码、示例与配置,系统讲解 SDK 的模块构成、构建与依赖校验、示例运行方式、任务编写 API、执行期通信协议、兼容性矩阵,以及从本地发布到 Maven Central 的完整发布与投票流程。读完本文,你将掌握如何构建 SDK、用 Java 编写并打包工作流、将其接入 Airflow 的 Java 队列,以及如何独立核验一个 Java SDK 发布候选。
一、SDK 定位与运行环境
Java SDK 是一个面向JVM的 Apache Airflow 开发套件。你可以使用任意 JVM 兼容语言编写工作流 Bundle,再由 Airflow 消费执行结果。SDK 本身与执行期逻辑使用Kotlin实现,但公开 API 面向 Java,仓库中捆绑了一个使用 Java 调用 SDK 的端到端示例,并在scala_spark_example/中提供了 Scala + Spark 的使用示例(见 ScalaSparkExample.scala)。
运行环境要求:
- SDK 运行时要求Java 11 或更高版本;
- 可选组件与开发工具可能有更高要求,具体见 airflow-jvm-conventions.gradle.kts 中定义的 toolchain;
- 服务端侧要求Airflow 3.3 及以上版本,且 supervisor 通信 schema 版本为
2026-06-16(见 capabilities.yaml 与gradle.properties中的airflowSupervisorSchemaVersion)。
更多 Java SDK 的使用细节,可参见 Airflow 官方文档Authoring and Scheduling下的 Java SDK 章节。
二、仓库布局:一个 JVM 多模块 Gradle 工程
java-sdk/是一个多模块 Gradle 工程,各模块职责清晰(源码结构见 java-sdk/settings.gradle.kts):
| 模块 | 职责 |
|---|---|
sdk/ | 核心库:公开 API(org.apache.airflow.sdk)与内部执行层(org.apache.airflow.sdk.execution),含Server、Client、Bundle、DagDef、Builder等核心类型 |
processor/ | 注解处理器,为被@Builder.Dag标注的类生成*Builder类(BuilderProcessor.kt,基于 kapt) |
plugin/ | Gradle 插件(org.apache.airflow.sdk),提供bundle打包任务、Manifest 属性注入与verifyBundleMainClass校验 |
bom/ | Bill of Materials POM,便于消费者以统一版本导入全部 SDK 构件 |
slf4j/ | SLF4J 日志 Provider,将 SLF4J 调用路由到 Airflow 日志存储 |
jul/ | java.util.loggingHandler,将 JUL 记录路由到 Airflow 日志存储(对应AirflowJulHandler.kt) |
jpl/ | Java Platform Logging Provider(System.Logger,JEP 264),路由 JPL 调用(对应AirflowSystemLoggerFinder.kt) |
log4j2/ | Log4j 2 Appender,路由 Log4j 2 事件到 Airflow 日志存储(对应AirflowLog4jAppender.java) |
example/ | 端到端示例 Bundle(注解 API + 接口 API,Java 源码) |
buildSrc/ | 共享 Gradle 约定插件(Java 版本、lint、格式化等) |
Python 侧的 JVM 启动协调器不在本目录内,位于:
- task-sdk/src/airflow/sdk/coordinators/java/:
JavaCoordinator(SubprocessCoordinator子类); - task-sdk/tests/coordinators/java/:Python 侧单元与集成测试。
三、构建 SDK 与依赖校验
构建整个 SDK 只需在java-sdk/目录下执行:
./gradlew build依赖变更审查
gradle/verification-metadata.xml为依赖、插件及其元数据固定了 SHA-256 校验和。只要该文件存在,Gradle 会自动启用依赖验证并默认采用strict模式,因此普通的./gradlew build就已经在执行校验。strict模式会在以下两种情况下失败:
- 校验和不匹配;
- 某个构件在元数据中完全没有条目——这通常是升级依赖时最常见的情况,意味着构建解析到了元数据未描述的内容。
仓库层面还有两条约定:
- 仓库集中在
settings.gradle.kts中声明,且**动态版本(dynamic)与变化版本(changing)**对项目依赖配置是被禁止的(插件 marker 与 detached 配置除外,需要手工固定版本); buildSrc/自行声明仓库,但其依赖同样被该元数据覆盖。
更新依赖或插件时,应从可信网络重新生成元数据。注意任务列表必须覆盖 CI 运行的全部内容,因为只有被调用任务实际解析到的依赖才会被记录:
./gradlew --write-verification-metadata sha256 --refresh-dependencies \ build \ :sdk:dokkaGeneratePublicationHtml :sdk:dokkaGeneratePublicationJavadoc \ sourceTarball checksumSourceTarball \ publishToMavenLocal -PskipSigning=true几点说明:
- 不带
-PskipSigning=true时签名任务会失败,但 Gradle 仍会从部分运行中写出元数据; - 重新生成只会追加,版本升级后过期的条目需要手工删除;
- 生成文件只记录仓库当时提供的内容,并不代表这些字节可信,务必把新坐标与校验和与依赖官方发布信息交叉核对;
- 生成元数据不覆盖:
example/、scala_spark_example/、kubernetes-tests/lang_sdk/java_example/的构建、foojay resolver 自动供应的 JDK,以及:sdk:syncSupervisorSchema拉取的 Supervisor Schema。
构建文档
./gradlew dokkaGenerate该命令使用 Dokka 构建 Java SDK 文档,同时生成 HTML 表示与 Javadoc(含:sdk:dokkaGeneratePublicationHtml与:sdk:dokkaGeneratePublicationJavadoc两个变体)。
四、端到端运行示例
要让示例真正跑起来,需要完成 SDK 构件发布、Bundle 打包、Airflow 队列配置与连接/变量准备四个环节。
1. 发布 SDK 到本地 Maven 仓库
./gradlew publishToMavenLocal -PskipSigning=true构建成功后,~/.m2/repository/org/apache/airflow/下会出现airflow-sdk、airflow-sdk-processor、airflow-sdk-bom、airflow-sdk-gradle-plugin等目录。
2. 打包示例 Bundle
进入示例工程并执行 bundle 任务:
# 已进入 example 目录,gradlew 在父目录 cd example ../gradlew bundleBundle 会输出到example/build/bundle。打包由 Gradle 插件(AirflowSdkPlugin.kt)负责,包括 bundle 任务、Manifest 属性注入与verifyBundleMainClass校验。
3. 放置 DAG 文件
将带 stub 任务的 DAG 放到 Airflow 能找到的位置,仓库提供了现成示例:java_examples.py。
4. 配置 Airflow 的 Java 队列
确保 Airflow 任务 worker 所在环境能执行java命令,然后配置[sdk]段的协调器,把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"}'classpath指向 Python 侧JavaCoordinator的完整导入路径;kwargs.jars_root是扫描 JAR Bundle 的目录列表;QUEUE_TO_COORDINATOR建立"队列名 → 协调器名"的映射,将java队列路由到名为java的协调器条目。
5. 准备连接与变量
示例 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五、编写任务:注解 API 与接口 API
SDK 提供两套编写任务的 API:注解驱动的声明式 API,以及面向底层编排的接口 API。
1. 注解 API(推荐)
org.apache.airflow.sdk.Builder容器类(Builder.kt)定义了三枚注解,由processor模块的BuilderProcessor.kt在编译期生成*Builder类:
| 注解 | 目标 | 参数 | 说明 |
|---|---|---|---|
@Builder.Dag | 类 | id(默认取类名)、to(生成的 Builder 类名,默认类名 +Builder) | 标注一个 DAG 定义类,处理器为其生成FooBuilder.build(),返回装配好的DagDef |
@Builder.Task | 方法 | id(默认取方法名) | 将方法标注为任务定义 |
@Builder.XCom | 方法参数 | task(默认取参数名) | 将参数标记为来自指定任务的 XCom 输入 |
仓库中的 AnnotationExample.java 展示了完整用法:
@Builder.Dag(id = "java_annotation_example") public class AnnotationExample { @Builder.Task(id = "extract") public long extractValue(Client client) throws InterruptedException { var pythonXcom = client.getXCom("python_task_1"); // 读取 Python 任务推送的 XCom var connection = client.getConnection("test_http"); // ... 业务逻辑 return new Date().getTime(); // 返回值自动作为 return_value XCom 推送 } @Builder.Task(id = "transform") public long transformValue(Client client, @Builder.XCom(task = "extract") long extracted) { var variable = client.getVariable("my_variable"); return new Date().getTime(); } }要点:
- 任务方法参数中的
Client由运行时自动注入; - 用
@Builder.XCom标注的参数会在任务执行前自动拉取上游任务同名(或task指定)的 XCom; - 方法的返回值会自动以
return_value为键推送为 XCom,供下游任务消费(常量XCOM_RETURN_KEY定义在 Client.kt); - 若某参数声明为原始类型(如
long)而对应 XCom 从未被推送,会抛出MissingXComException——改用装箱类型(如Long)即可接收null(见 Client.kt 的MissingXComException构造逻辑)。
示例中还演示了重试感知:load任务在第一次执行(context.ti.tryNumber == 1)时故意抛异常,由于 Java SDK 会在ti_context.should_retry置位时返回RetryTask而非终态的FAILED,supervisor 会把任务标记为UP_FOR_RETRY,重试后任务正常完成。同时concurrent任务验证了单个 supervisor 通道可以承载跨线程并发客户端调用(8 线程 × 32 次getConnection)。
2. 接口 API(底层编排)
当需要底层控制时,可直接使用 DagDef.kt 中的类型:
DagDef:DAG 定义,ID 只能包含 ASCII 字母数字、横线、点或下划线,且在一个 Bundle 内唯一;通过addTask(id, Class)链式注册任务;TaskDef:单个任务定义(ID + 实现Task接口的类,类需有公共无参构造器);Task接口:实现execute(context, client)方法,抛出的任何异常都会将任务实例标记为失败。
var dag = new DagDef("java_etl") .addTask("extract", Extract.class) .addTask("load", Load.class);3. Bundle 与进程入口
Bundle.kt 定义Bundle为当前 JVM 进程可执行的全部DagDef的不可变快照,DAG ID 重复会抛出IllegalArgumentException;BundleBuilder接口负责收集 DAG 并提供build()。进程入口通过Server启动,见 ExampleBundleBuilder.java:
public class ExampleBundleBuilder implements BundleBuilder { @Override public Iterable<DagDef> getDags() { return List.of( InterfaceExampleBuilder.build(), AnnotationExampleBuilder.build(), XComCastingExampleBuilder.build()); } public static void main(String[] args) { var bundle = new ExampleBundleBuilder().build(); Server.create(args).serve(bundle); } }Server(Server.kt)是 SDK 的运行核心:
Server.create(args)解析两个由 Airflow 自动注入的命令行参数:--comm host:port(任务执行消息通道)与--logs host:port(日志转发通道),无需手工构造;serve(bundle)是阻塞入口(对serveAsync的包装),连接协调器后分发任务执行请求,当协调器关闭连接(通常在一次任务实例执行后)时进程退出;- 进程启动时会同时打开
--comm与--logs两个 socket,随后等待首帧消息:若收到StartupDetails则按dag_id+task_id查表调用用户任务方法;若收到ErrorResponse则抛出ApiError。
4. 运行时 Client API
任务方法中注入的Client(Client.kt)封装了与 Airflow supervisor 的通信,所有读写默认自动限定在当前 DAG run 与任务实例范围内:
| 方法 | 作用 |
|---|---|
getConnection(id) | 读取 Airflow 连接存储中的连接,返回Connection(含 id/type/host/schema/login/password/port/extra 字段),连接不存在或调用失败抛ApiError |
getVariable(key) | 读取 Airflow 变量,未设置时返回null |
getXCom(key, dagId, taskId, runId, mapIndex, includePriorDates) | 读取其他任务推送的 XCom;map_index为null时对 mapped 任务返回按 map index 升序聚合的"集体结果"列表 |
setXCom(key, value) | 推送 XCom 供下游读取,值必须可 JSON 序列化 |
六、执行架构与线协议
JavaCoordinator与 JVM 子进程的协作流程如下(详见 README "Contributing" 一节与 coordinator.py):
- 启动 JVM:Airflow supervisor 判定任务应在 Java 队列运行后,调用
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 通过 comm 通道向 supervisor 发送请求(GetVariable、GetConnection、GetXCom、SetXCom 等)并等待响应;
- 结束:任务完成或抛异常时,JVM 发送
TaskState消息并关闭 socket,进程随之退出。
关键实现细节(Comm.kt):
- 帧格式:所有帧均为4 字节大端长度前缀 + MessagePack 载荷;
- 请求/响应关联:每个请求获得递增的
id,通过pending表(ConcurrentHashMap<Int, CompletableDeferred<IncomingFrame>>)与读循环协程配对响应;communicate<T>()对ErrorResponse自动抛ApiError; - 帧大小防护:入站帧上限取协议规定的
Frame.MAX_FRAME_LENGTH与堆内存 1/8(MAX_HEAP_FRACTION_PER_FRAME = 8,随-Xmx自适应)的较小值,超限帧在分配内存前即被拒绝,将不可恢复的 OOM 转化为可捕获的FrameProcessingException; - 流式解码:
ChannelFrameInput以 64 KiB 分块喂给 MessageUnpacker,避免超大帧一次性大分配;长度前缀承诺的字节数多于实际载荷(under-run)会被视为流失步; - 日志转发:SDK 自身(而非用户代码)产生的日志消息通过
--logssocket 转发,由 supervisor 追加到 Airflow 日志存储。对应实现见execution/Logger.kt与各日志桥接模块(SLF4J/JUL/JPL/Log4j 2)。
线协议由 task-sdk/src/airflow/sdk/execution_time/schema/schema.json 统一定义:新增消息类型需要同时修改 Python 侧schema.json与 JVM 侧execution/Comm.kt+execution/Client.kt。关于协调器架构与各语言 SDK 共享的核心集成面,可进一步阅读 airflow-core/adr/lang-sdk 下的架构决策记录(ADR)。
从 Python 协调器源码看,JavaCoordinator支持以下配置项(见 coordinator.py 类 docstring 与字段定义):
| 参数 | 默认值 | 说明 |
|---|---|---|
java_executable | "java" | java命令路径,默认依赖$PATH |
jvm_args | [] | 额外 JVM 参数,如["-Xmx1024m"] |
jars_root | 必填(至少 1 个) | 扫描 JAR Bundle 的目录列表 |
main_class | "" | 显式入口类;未指定时扫描jars_root寻找带Main-Class元数据的可执行 JAR,若存在多个可能不确定 |
task_startup_timeout | 10 秒 | 等待任务进程启动的最长时间 |
JavaCoordinator._build_execute_task_command()会从 JAR 的META-INF/MANIFEST.MF中读取Main-Class与Airflow-Supervisor-Schema-Version两个条目,后者用于核对线协议 schema 版本——Java SDK 打包的 JAR 会自动写入该元数据,只要依赖 JAR 原样部署就无需额外处理;若重新打包依赖,则必须在一个 JAR 中复现该条目。
七、兼容性矩阵
README 中的兼容性矩阵由 capabilities.yaml 自动生成(update-java-sdk-readme-matrixpre-commit 钩子负责重生成,请勿手工编辑表格),符合 contributing-docs/30_new_language_sdk.rst 定义的 Language SDK 一致性规范。
当前 SDK 的能力概览(对应 Airflow 3.3、supervisor schema2026-06-16):
TaskInstance 状态
- 支持(MUST,3.3):
success、failed、up_for_retry(通过RetryTask)、removed; - 未支持:
skipped、deferred、up_for_reschedule、awaiting_input(运行时尚未发出对应消息)。
运行时能力
- 支持(MUST,3.3):
mixed-lang-stub-target(@task.stub)、task-logging(SLF4J + JPL 桥接到任务日志)、xcom-read-write、connection-read、self-contained-bundle(Airflow 元数据内嵌于 jar 构件); - 未支持:
variable-read-write(目前仅getVariable,尚无写通道)、retry-policy、task-state-store、asset-state-store、asset-event-emit、asset-event-read。
Native-Dag 编写
native-dag-authoring尚未实现,因此task-args、dag-params、taskflow-dependencies、branching、dag-test、task-group、dynamic-task-mapping、asset-inlets-outlets、asset-scheduling、object-store等原生能力均标注为 n/a(仅在 native-dag-authoring 支持后才适用)。
在 Airflow 侧,混合语言场景通过@task.stub声明 stub 任务并指定队列,由 supervisor 把对应队列的任务路由给 Java 协调器执行,Java 任务通过 XCom 与 Python 任务交换数据(示例 DAG java_examples.py 中可见python_task_1/python_task_2与 Java 任务的 XCom 互通)。
八、发布流程
SDK 通过 ASF Nexus staging 仓库发布到 Maven Central。凡发布到 Maven Central(而非 Snapshots)的版本均视为正式发布(含 alpha、beta 等),每次发布都需 PMC 投票;仅-SNAPSHOT构建可免投票发布。前置条件:具备访问 repository.apache.org 的 ASF committer 账号,以及已加入项目 KEYS 文件并上传到公共密钥服务器的 GPG 密钥。
1. 版本号
编辑gradle.properties设置projectVersion=<VERSION>,或对单次命令用-PprojectVersion=<VERSION>覆盖。版本字符串需符合 Maven 版本顺序规范,例如 beta 1 写作1.0.0-beta1。main分支在发布之间应保持在-SNAPSHOT版本(snapshot 排序在<VERSION>之后、GA 之前,因此 beta 之后无需额外 bump)。
2. 打 RC 标签
git tag -s java-sdk/<VERSION>-rc<N> -m "Java SDK <VERSION> RC <N>" git push upstream java-sdk/<VERSION>-rc<N>RC 编号保留在标签名中(投票失败只需递增到下一 RC),构件版本本身不带 RC 后缀。投票前需先推送标签,便于评审者检出确切的待投票源码。
3. 本地核验 POM
rm -rf ~/.m2/repository/org/apache/airflow/ # 清空旧构件 ./gradlew publishToMavenLocal -PskipSigning=true less ~/.m2/repository/org/apache/airflow/airflow-sdk/*/airflow-sdk-*.pom less ~/.m2/repository/org/apache/airflow/airflow-sdk-bom/*/*.pom less ~/.m2/repository/org/apache/airflow/airflow-sdk-processor/*/airflow-sdk-*.pom less ~/.m2/repository/org/apache/airflow/airflow-sdk-gradle-plugin/*/airflow-sdk-*.pom less ~/.m2/repository/org/apache/airflow/sdk/org.apache.airflow.sdk.gradle.plugin/*/*.pom核对各 POM 的坐标、描述、license、SCM 与 organization 字段。
4. 对本地仓库做发布演练
rm -rf /tmp/local-maven-repo ./gradlew publish -PmavenUrl=file:///tmp/local-maven-repo -PskipSigning=true ls /tmp/local-maven-repo/org/apache/airflow/输出应与上一步~/.m2中的构件一致。本地演练不需要签名;如需测试签名,按下一节配置 GPG 私钥与口令并去掉-PskipSigning=true。
5. 发布到 ASF Nexus staging
将凭据写入~/.gradle/gradle.properties(避免进入 shell 历史):
mavenUsername=your-asf-nexus-token-username mavenPassword=your-asf-nexus-token-password signing.password=your-gpg-key-passphrase然后暂存并关闭发布(根build.gradle.kts应用的 Gradle Nexus Publish Plugin 会把所有模块聚合进同一个staging 仓库,无需逐模块协调):
./gradlew publishToApache closeApacheStagingRepository \ --no-configuration-cache \ -P"signing.key=$(gpg --armor --export-secret-keys your-gpg-key-fingerprint)"注意三点:
- 签名密钥通过命令行传入,因其含换行符不适合放在 Gradle properties 文件;
- 也可用环境变量提供凭据:
ASF_NEXUS_USERNAME、ASF_NEXUS_PASSWORD、SIGNING_KEY、SIGNING_PASSWORD(适合 CI); - 项目全局启用 configuration cache,但 staging 任务(
publishToApache、closeApacheStagingRepository、releaseApacheStagingRepository)通过 Nexus REST API 工作且不兼容 configuration cache,因此发布命令必须加--no-configuration-cache。
发布后在 Nexus 的Staging Repositories中打开已关闭的仓库,核验所有模块均包含 jar、-sources.jar、-javadoc.jar(适用处)、.pom与.asc签名,并检查Updated by、Uploaded Date、Last Modified字段。
6. 上传源码包
投票正式针对的是已签名的源码包(Maven 构件是 convenience binaries)。sourceRelease任务基于已提交的java-sdk源码(含LICENSE与NOTICE)一次生成签名与校验和:
./gradlew sourceRelease -PgitRef=java-sdk/<VERSION>-rc<N>输出到build/distributions/的三个文件:
apache-airflow-java-sdk-<VERSION>-src.tar.gz apache-airflow-java-sdk-<VERSION>-src.tar.gz.asc apache-airflow-java-sdk-<VERSION>-src.tar.gz.sha512注意:源码包刻意省略 Gradle wrapper 脚本(gradlew、gradlew.bat)与gradle/wrapper/gradle-wrapper.jar——ASF 源码发布不得包含编译产物(见 LEGAL-570),而缺少 jar 的脚本没有用处;gradle/wrapper/gradle-wrapper.properties会保留,以固定 Gradle 版本与发行版校验和供重新生成 wrapper 时核验。
将三个文件复制到 ASF distdev仓库并提交:
svn checkout https://dist.apache.org/repos/dist/dev/airflow <dist-dev-checkout> cd <dist-dev-checkout> mkdir -p java-sdk/<VERSION>-rc<N> cp <path-to>/java-sdk/build/distributions/apache-airflow-java-sdk-<VERSION>-src.tar.gz* \ java-sdk/<VERSION>-rc<N>/ svn add --parents java-sdk/<VERSION>-rc<N> svn commit -m "Add Apache Airflow Java SDK <VERSION>-rc<N> source release candidate"7. 发起投票
向dev@airflow.apache.org发送[VOTE]邮件,链接 git 标签与提交、dist/dev中的源码包、已关闭的 Nexus staging 仓库及KEYS文件。投票至少开放 72 小时。README 提供了完整模板(subject 为[VOTE] Release Apache Airflow Java SDK <VERSION> based on <VERSION>-rc<N>,正文列出 Maven 构件清单、Git 信息、源码包与 convenience binaries URL、投票选项)。发送前检查:所有构件能通过 BOM 交叉解析、"Changes since rc "与截止时间已填写、对渲染后的邮件执行grep '<' email.txt确认无输出(任何匹配都说明占位符未填完)。
8. 独立核验发布候选
任何社区成员都应在投票前独立核验候选,核验清单如下:
- 校验和:
sha512sum -c apache-airflow-java-sdk-<VERSION>-src.tar.gz.sha512; - 签名:下载并导入
KEYS文件,用gpg --verify验证.asc签名; - 与 git 标签比对:解包 tarball,与标签的干净检出做
diff -rq,除.gitattributesexport-ignore排除的文件(gradlew、gradlew.bat、gradle-wrapper.jar、scripts)外应无差异;解包后的顶层目录应为apache-airflow-java-sdk-<version>(不带仅出现在压缩包文件名中的-src后缀); - 无二进制文件:ASF 源码发布不得含编译产物,用
file扫描非文本文件应无输出; - 从源码构建:用本地 Gradle 依
gradle-wrapper.properties中的distributionUrl与distributionSha256Sum重新生成 wrapper,再./gradlew build; - staged 二进制冒烟测试:从一个临时工程指向 staging 仓库 URL,声明
org.apache.airflow:airflow-sdk-bom:<VERSION>依赖,确认传递构件(含airflow-sdk-jpl)可解析、示例 Bundle 能构建(对应脚本见 scripts/ci/smoke-test-staged-binaries.sh 与 scripts/ci/verify-source-release.sh)。
9. 投票成功后的收尾
回复
[RESULT][VOTE]统计,然后releasestaging 仓库(同步到 Maven Central 需数小时,不重新构建或签名):./gradlew releaseApacheStagingRepository --no-configuration-cache将源码包从
dist/dev移至dist/release(svn mv);在同一被投票的提交上打最终版本标签
java-sdk/<VERSION>(保留 RC 标签以便追溯);发布 GitHub release,附带被投票且已签名的源码构件(
gh release create,预发布版本加--prerelease,核验标签用--verify-tag);等待约 1 小时(Maven Central 同步)后发送纯文本
[ANNOUNCE]邮件到users@airflow.apache.org(抄送dev@),并在 ASF Committee Report Helper 中记录发布;触发Publish Docs to S3workflow 发布 API 文档,确认
https://airflow.apache.org/docs/java-sdk/stable/可解析且/docs/java-sdk/重定向到它。
若投票失败:关闭投票、在 Nexus 中dropstaging 仓库、删除dist/dev候选、修复问题后切下一 RC(...-rc2)。发布版本号不变,仅标签中的 RC 计数递增。
九、测试与编码规范
运行测试
# 运行全部 JVM 测试 ./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编码规范
- 所有 SDK 与 processor 源码均为Kotlin;Java 是公开 API 目标而非实现语言;
- 保持
sdk/src/main/kotlin/(公开 API 面)不混入内部实现细节,内部实现属于execution/子包; - 注解处理器(
BuilderProcessor.kt)使用 kapt:新增注解时在Builder.kt定义、在BuilderProcessor.kt处理,并在processor/src/test/kotlin/添加 golden-output 测试; - Python 协调器继承
SubprocessCoordinator,除_build_execute_task_command提供的接口外,不要从 Python 深入 JVM 进程内部; - 提交前运行
./gradlew ktLintCheck spotlessCheck(或ktLintFormat spotlessApply),项目强制 Kotlin 与 Java 格式; - 所有新文件必须带 Apache License 头。
常见开发任务
新增一个Client方法(如新的 Airflow API 调用):
- 若消息类型是新的,从
schema.json重新生成 POJO; - 在
execution/Comm.kt(或新文件)添加 Kotlin 请求/响应数据类; - 在公开
Client.kt添加委托给execution/Client.kt做 wire 调用的方法; - 在
sdk/src/test/kotlin/.../ClientTest.kt编写 mock socket 层的单元测试; - 若变更用户可见,同步更新
airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst。
新增注解:在Builder.kt定义注解接口 → 在BuilderProcessor.kt生成对应代码 → 在BuilderTest.kt添加期望生成输出的测试 → 更新java.rst的注解表。
修复组帧或协议 bug:聚焦execution/Comm.kt与execution/Frame.kt;CommTest.kt覆盖编解码往返,修复前先添加复现该 bug 的回归测试。
PR 检查清单
- 运行
./gradlew build test(JVM)与对应 pytest 套件(Python 协调器); - 确认示例 Bundle 仍可编译(按"运行示例"一节做到打包步骤);
- 若
schema.json变更,确认 JVM 与 Python 两侧都能处理新/改字段; - 为每项行为变更添加或更新测试;
- 对
task-sdk/的用户可见变更,在airflow-core/newsfragments/下添加 newsfragment。
十、小结
Apache Airflow Java SDK 让 JVM 团队以惯用的 Java/Kotlin/Scala 方式编写 Airflow 任务:通过@Builder.Dag/@Builder.Task/@Builder.XCom注解声明式描述 DAG,用 Gradle 插件一键打成自包含 Bundle,由JavaCoordinator以子进程方式拉起 JVM 并通过 4 字节长度前缀 + MessagePack 的 comm 通道与 supervisor 协作,天然支持与 Python 任务的 XCom 互通和@task.stub混合语言调度。其发布流程遵循 ASF 规范,从本地 POM 核验、staging 发布、源码包签名、PMC 投票到最终同步 Maven Central 均有清晰的步骤与可执行的验证清单,为在 Java 生态中落地 Airflow 提供了完整链路。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考