ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Airflow Java SDK 能力兼容性矩阵深度解读:TaskInstance 状态、运行时能力与 Native-Dag 支持全景

Apache Airflow Java SDK 能力兼容性矩阵深度解读:TaskInstance 状态、运行时能力与 Native-Dag 支持全景 Apache Airflow Java SDK 能力兼容性矩阵深度解读TaskInstance 状态、运行时能力与 Native-Dag 支持全景【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow Java SDK 允许开发者用 Java 及其他 JVM 语言编写 Airflow 任务实现并在标准 Airflow Worker 上由 Python 侧协调器拉起 JVM 子进程执行。本文以 java-sdk/sdk/module.md 这份官方兼容性矩阵文档为骨架逐行拆解其中 TaskInstance 状态、运行时能力与 Native-Dag 编写三个维度共 30 余项声明的真实含义并结合仓库源码Client.kt、Task.kt、Comm.kt、JavaCoordinator等验证每一项能力在代码层面的落地情况帮助读者准确判断当前 Java SDK 能做什么、不能做什么以及如何据此规划自己的任务实现。Java SDK 与 Language SDK 一致性规范的关系module.md开篇即给出定位The Apache Airflow Java SDK — author and run Airflow task implementations in JVM languages.它并不是一份产品宣传页而是一份机器可读的能力声明 人工可读的兼容性清单其语义基准来自仓库内的 Language SDK conformance specification即contributing-docs/30_new_language_sdk.rst。这份一致性规范是理解整个矩阵的前提它定义了三个关键分层关键词遵循 RFC 2119 语义MUSTSDK 若要被认为符合规范必须支持的能力SHOULD强烈建议支持不支持的 SDK 依然可用但会缺失某些体验MAY可选能力SDK 按自身成熟度声明是否支持。从源码结构看Airflow 3.3 起的标准 Worker 通过Python 协调器 目标语言 SDK双组件模式执行非 Python 任务协调器Python决定如何启动外来运行时、如何与它通信语言 SDK本仓库为 Kotlin 实现、面向 Java 公开 API实现通信的另一端。Java SDK 采用的是SubprocessCoordinator子类 TCP 双 socket 传输方案具体见下文运行时能力部分的源码验证。兼容性矩阵从哪来capabilities.yaml 是唯一事实源矩阵不是手写的表格而是由一个 prek 钩子自动生成的产物。规范文档中明确给出了生成链路java-sdk/capabilities.yaml - 唯一需要人工编辑的文件 | | hook: update-java-sdk-readme-matrix | -- java-sdk/README.md (面向贡献者) -- java-sdk/sdk/module.md (Dokka - 对外发布的 API 参考)也就是说java-sdk/capabilities.yaml 这份 YAML 才是 Java SDK 能力声明的唯一事实源。module.md中!-- BEGIN AUTO-GENERATED LANG-SDK COMPAT MATRIX --与!-- END ... --之间的全部内容都由该钩子重写一旦表格与 manifest 漂移构建就会失败退出。因此本文讲解的每一项能力读者都可以在 java-sdk/capabilities.yaml 中找到对应的supported: true/false声明。manifest 中还有一个关键元信息sdk: java min_airflow_version: 3.3 supervisor_schema_version: 2026-06-16min_airflow_version: 3.3说明该 SDK 要求的最低 Airflow 版本是 3.3Airflow 3.3 起才引入标准 Worker 运行外来语言任务的机制supervisor_schema_version: 2026-06-16是 SDK 所理解的监督者线协议Supervisor Schema版本它以YYYY-MM-DD日期串标识用于与 supervisor 协商消息格式。该值同时被刻印进 JAR 的 manifestAirflow-Supervisor-Schema-Version属性capabilities.yaml中的注释明确要求二者保持同步否则渲染钩子会报错。维度一TaskInstance 状态支持任务运行结束后子进程通过--commsocket 向 supervisor 上报终结消息对应状态转换。规范中强调只有子进程能发出的状态才属于 SDK 一致性范畴queued、scheduled、running、restarting、upstream_failed等由调度器拥有的状态永远不由运行时发送因此不在矩阵内。已支持success、failed、up_for_retry、removed自 3.3 起矩阵中这四项均为✓且起始版本都是 3.3。源码层面的证据集中在 java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Task.kt 的TaskResult工厂中fun success(...) SucceedTask().also { it.state success; ... } fun retry(...) RetryTask().also { ... } fun of(state: TaskState.State, ...) TaskState().also { it.state state; ... } fun failure(shouldRetry: Boolean) if (shouldRetry) retry() else of(TaskState.State.FAILED)success任务方法正常返回后发送SucceedTaskstate successfailed任务抛出未捕获异常且无重试机会时发送TaskState且state failedup_for_retry任务失败但仍有重试次数时发送RetryTask消息supervisor 据此将任务置为up_for_retry。规范特别提醒失败详情字段名必须是retry_reason而非reasonremoved当StartupDetails中携带的dag_id task_id在 Bundle 中查不到对应任务定义时TaskRunner.runTask直接返回TaskResult.of(TaskState.State.REMOVED)——这对应任务被移除的场景。TaskRunner的异常分类也值得注意构造器抛出的InvocationTargetException/ReflectiveOperationException一律shouldRetry false每次重试必然同样失败而Throwable类型静态初始化器、链接错误等则按request.tiContext.shouldRetry决定是否重试——因为换一个全新 JVM 可能就能成功。未支持skipped、deferred、up_for_reschedule、awaiting_input矩阵中这四项均为✗注释统一说明runtime does not emit ... yet。对照规范skippedSHOULD需要TaskState携带state skipped用于分支/跳过语义——运行时尚未实现deferredMAY需要DeferTask消息并桥接 triggerer——未实现up_for_rescheduleMAY需要RescheduleTask用于 reschedule 模式的传感器——未实现awaiting_inputMAY需要AwaitInputTask用于人工介入human-in-the-loop任务——未实现。从capabilities.yaml的注释可以确认这一现状The runtime terminates a task with SucceedTask, RetryTask, or TaskState (failed/removed); it does not yet emit skipped, DeferTask, RescheduleTask, or AwaitInputTask.对开发者的直接影响当前 Java SDK 只能覆盖普通任务跑到成功/失败并支持重试这一 MUST 层级分支跳过、延迟、reschedule 传感器等高级语义尚不可用在 Java 任务中不要依赖这些状态转换。维度二运行时能力Runtime capabilities运行时能力描述任务体在执行过程中能做什么与任务声明位置Python DAG 的task.stub还是原生 DAG无关。这一维度最能体现 SDK 的实际功能边界。mixed-lang-stub-targetMUST✓3.3SDK 能执行 Python DAG 中通过task.stub装饰器声明的任务——这是所有 Language SDK 的主要执行路径。Java 侧通过Builder.Dag/Builder.Task注解声明任务方法见 java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Builder.kt由BuilderProcessor注解处理器生成*Builder类最终装配成DagDef注册进Bundle供 stub 引用。task-loggingMUST✓3.3任务日志stdout/stderr 与结构化日志记录经--logssocket 转发给 supervisor最终汇入 Airflow 任务日志。矩阵注释点明实现方式SLF4J JPL 桥接到任务日志。仓库中对应四个独立子项目把主流 JVM 日志体系全部接到 Airflow 日志通道java-sdk/slf4jAirflowSlf4jProviderjava-sdk/juljava.util.logging的AirflowJulHandlerjava-sdk/jplSystem.LoggerJEP 264 的AirflowSystemLoggerFinderjava-sdk/log4j2Log4j 2 的AirflowLog4jAppender--logs通道只承载基础设施日志SDK 自身产生的记录而非用户代码的 stdout/stderr 输出。规范要求日志按newline-delimited JSON发送structlog 风格事件含event、level、timestamp字段且子进程启动早期、--logssocket 尚未建立时产生的记录必须先缓冲、连接建立后再按序冲刷——Java SDK 实现了该缓冲机制。xcom-read-writeMUST✓3.3任务可跨 Python 边界读写 XCom。公共 API 在 java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Client.kt 中fun getXCom(key XCOM_RETURN_KEY, dagId details.ti.dagId, taskId: String, runId details.ti.runId, mapIndex: Int? null, includePriorDates: Boolean false): Any? fun setXCom(key XCOM_RETURN_KEY, value: Any)默认 XCom key 是return_valueXCOM_RETURN_KEY常量任务返回值自动以此键推送getXCom默认作用域为当前 Dag run 与 task instance也支持跨 DAG/跨 run 读取mapIndex null时对映射任务返回集体结果按 map index 升序聚合的列表非映射任务等价于-1参数类型为原始类型如int却读到 null 时会抛出MissingXComException其消息建议改用装箱类型Integer——这是 Java 侧特有的空值语义细节。connection-readMUST✓3.3任务可按 ID 解析 Airflow Connection。Client.getConnection(id)返回包含id、type、host、schema、login、password、port、extra的Connection数据类。值得注意的实现细节wire 上端口等整数字段经 MessagePack 解码后是Long因此源码中做了(port as Number?)?.toInt()的数值转换而非类型强转。variable-read-writeMUST✗这是矩阵中声明为 MUST 但未实现的特例capabilities.yaml注明getVariable only; no write over the comm socket yet。当前Client只暴露getVariable(key): Any?读取尚未提供通过通信通道写入/删除 Variable 的 API。如果你需要在任务中设置变量现阶段只能借助其他通道如 Python 侧任务或等待后续版本补齐。self-contained-bundleMUST✓3.3SDK 的构建产物将 Airflow 元数据dag_id、task_id及完整任务描述符嵌入与任务代码相同的 artifact而非随附 sidecar 文件使部署单元自描述。Java 侧的落地方式Gradle 插件java-sdk/plugin 中的AirflowSdkPlugin负责 bundle 任务、manifest 属性注入与verifyBundleMainClassJAR 的META-INF/MANIFEST.MF中写入Main-Class与Airflow-Supervisor-Schema-Version属性Python 侧JavaCoordinator见 task-sdk/src/airflow/sdk/coordinators/java/coordinator.py通过zipfile读取 JAR manifest 的这两个属性来定位入口类并协商 schema 版本。未实现能力retry-policy 与状态/资产相关 API矩阵中retry-policy、task-state-store、asset-state-store、asset-event-emit、asset-event-read五项均为 MAY 层级且全部✗注释分别说明no task-facing retry-policy API yet、no task-facing state-store API yet、runtime does not emit asset events yet 等。注意retry-policy与上报up_for_retry是两回事SDK 支持普通重试RetryTask但尚未向任务代码开放检查失败并决定是否重试、自定义重试延迟的 API。维度三Native-Dag authoring条件性维度最后一个维度是用目标语言编写整个 DAG的能力由伞形能力native-dag-authoringSHOULD统辖。Java SDK 当前标注为✗注释明确 native Dag authoring not implemented yet。这里必须厘清一个容易误解的点Builder.Dag/Builder.Task注解并不等于 native-dag-authoring。注解驱动的*Builder模式Builder.kt用于在 Java 侧声明任务方法并注册到 Bundle供 Python DAG 的 stub 引用而 native-dag-authoring 指的是脱离 Python DAG 文件、直接在目标语言中定义整个可被 Airflow 调度的 DAG后者尚未实现。因此矩阵中用†标记的 10 项能力——task-args、dag-params、taskflow-dependenciesMUST †branching、dag-testSHOULD †task-group、dynamic-task-mapping、asset-inlets-outlets、asset-scheduling、object-storeMAY †——其层级仅在 native-dag-authoring 被支持时才适用。对尚未实现原生 DAG 的 SDK 而言它们标注为n/anot applicable而非不支持。当前 Java SDK 阶段用 Python DAG task.stub声明、Java 实现任务体依然是唯一受支持的任务执行方式。完整兼容性矩阵自动生成表原文以下为module.md与 java-sdk/README.md 中完全一致、由capabilities.yaml自动生成的兼容性矩阵全文最小 Airflow 版本 3.3supervisor schema 2026-06-16DimensionTierSupportedSinceNotesTaskInstance statesstate:successMUST✓3.3state:failedMUST✓3.3state:up_for_retryMUST✓3.3RetryTaskstate:skippedSHOULD✗–runtime does not emit TaskState skipped yetstate:deferredMAY✗–runtime does not emit DeferTask yetstate:up_for_rescheduleMAY✗–runtime does not emit RescheduleTask yetstate:awaiting_inputMAY✗–runtime does not emit AwaitInputTask yetstate:removedMAY✓3.3Runtime capabilitiescapability:mixed-lang-stub-targetMUST✓3.3task.stubcapability:task-loggingMUST✓3.3SLF4J JPL bridged to the task logcapability:xcom-read-writeMUST✓3.3capability:connection-readMUST✓3.3capability:variable-read-writeMUST✗–getVariable only; no write over the comm socket yetcapability:self-contained-bundleMUST✓3.3Airflow metadata embedded in the jar artifactcapability:retry-policyMAY✗–no task-facing retry-policy API yetcapability:task-state-storeMAY✗–no task-facing state-store API yetcapability:asset-state-storeMAY✗–no task-facing state-store API yetcapability:asset-event-emitMAY✗–runtime does not emit asset events yetcapability:asset-event-readMAY✗–no task-facing asset-event API yetNative-Dag authoringcapability:native-dag-authoringSHOULD✗–native Dag authoring not implemented yetcapability:task-argsMUST †n/a–capability:dag-paramsMUST †n/a–capability:taskflow-dependenciesMUST †n/a–capability:branchingSHOULD †n/a–capability:dag-testSHOULD †n/a–capability:task-groupMAY †n/a–capability:dynamic-task-mappingMAY †n/a–capability:asset-inlets-outletsMAY †n/a–capability:asset-schedulingMAY †n/a–capability:object-storeMAY †n/a–标记说明✓ 支持 · ✗ 不支持 · n/a 不适用带 † 的层级仅在native-dag-authoring被支持时生效。修改矩阵请编辑 java-sdk/capabilities.yaml 后运行update-java-sdk-readme-matrix钩子不要手工改表。源码级验证JavaCoordinator 如何拉起 JVM 子进程矩阵中self-contained-bundle与mixed-lang-stub-target的落地依赖 Python 侧协调器。JavaCoordinatortask-sdk/src/airflow/sdk/coordinators/java/coordinator.py继承SubprocessCoordinator其_build_execute_task_command实现展示了完整链路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要点递归扫描jars_root通过(st_dev, st_ino)去重防止符号链接环把所有.jar收集进 classpath从 JAR manifest 读取Main-Class与Airflow-Supervisor-Schema-Version后者作为subprocess_schema_version返回供 supervisor 协商消息格式——这正是supervisor_schema_version: 2026-06-16被刻进 JAR 的原因基类SubprocessCoordinator负责在127.0.0.1上绑定两个临时 TCP socket把--commhost:port与--logshost:port追加到命令尾部后拉起子进程并校验连接方属于该进程树。JVM 侧的入口在 java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.ktServer.create(args)解析--comm/--logs两个参数由 Airflow 自动传入不期望手工构造serve(bundle)立即并发连接两个 socket随后dispatchTask读取首帧StartupDetails按dag_id task_id在 Bundle 中查找任务定义并调用用户任务方法。线协议帧结构在 java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Comm.kt 中实现4 字节大端无符号长度前缀 MessagePack 载荷SDK→supervisor 为[id, body]二元组supervisor→SDK 为[id, body, error]三元组请求 id 单调递增、响应回显同 id 以支持并发关联。CoordinatorComm还对超大帧做了防护——单帧上限取Frame.MAX_FRAME_LENGTH与maxMemory() / 8的较小值防止恶意或损坏的帧长度前缀导致 OOM。从矩阵出发如何规划你的 Java 任务把矩阵翻译成日常开发约束任务状态你的任务方法只需关心正常返回success与抛异常failed/up_for_retry两种结局重试由tiContext.shouldRetry自动决定不要在 Java 任务里实现分支跳过skipped、deferral 等语义当前运行时不会发出这些状态数据交互Client是唯一入口——getXCom/setXCom传值注意返回值默认 key 为return_value、getConnection取连接、getVariable读变量变量目前只读不写日志直接使用 SLF4J、System.LoggerJPL、JUL 或 Log4j 2 中你熟悉的任一体系SDK 会自动把记录桥接到 Airflow 任务日志无需额外配置部署形态产物是一个嵌入了 Airflow 元数据的自包含 JAR含Main-Class与Airflow-Supervisor-Schema-Versionmanifest 属性把 bundle 目录配进AIRFLOW__SDK__COORDINATORS的jars_rootJava 任务即可被调度执行DAG 编写现阶段仍以 Python DAG task.stub声明、Java 实现任务体为主native DAG authoring 及其带 † 的能力branching、task-group、动态任务映射等尚不可用规划时勿作前提假设。端到端的运行方式可参考 java-sdk/README.md 的 Running the example 章节以及面向用户的完整使用文档 airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst。结语module.md的价值在于把Java SDK 当前能干什么压缩成了一张可机器校验、可人工检索的精确清单。透过它可以看到一条清晰的演进路径MUST 层级的状态与能力成功/失败/重试、stub 任务、日志、XCom、Connection、自包含 bundle已经全部落地SHOULD/MAY 层级的 skipped、deferral、变量写入、原生 DAG 编写与资产相关 API 尚未实现。这张矩阵会随 java-sdk/capabilities.yaml 的每次更新而自动重生成——它是 SDK 能力现状最权威、最不易过时的观测窗口。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表