airflow-java-sdk

作者: astronomer

Guide for contributing to the Airflow Java SDK (AIP-108). Use this skill whenever a contributor is working in the `java-sdk/` directory or on the Java…

npx skills add https://github.com/astronomer/airflow --skill airflow-java-sdk

Airflow Java SDK contributor guide

The Java SDK lets Airflow tasks execute JVM code (Java, Kotlin, or any JVM language). You are helping a contributor work in one or both of these locations:

  • java-sdk/ — the JVM-side library (Kotlin source, published to Maven)
  • task-sdk/src/airflow/sdk/coordinators/java/ — the Python coordinator that launches the JVM subprocess

Read these two documents early in every session — they contain the authoritative reference material:

  • airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst — user-facing guide: annotation vs. interface API, XCom type mapping, Gradle/Maven steps, coordinator config.
  • java-sdk/README.md — contributor guide: repository layout, detailed execution walkthrough, Gradle + Breeze test commands, coding conventions, common tasks, and PR checklist.

SDK package architecture

The JVM-side library is split into two packages with distinct visibility rules:

  • org.apache.airflow.sdk — public, user-facing API. Classes here (e.g. Client, Bundle, BundleBuilder, Server) are stable contracts that DAG authors and task implementers import directly. Changes to this package are breaking changes.
  • org.apache.airflow.sdk.execution — internal implementation detail. Everything in this package (CoordinatorComm, LogSender, Log, Client in execution/, generated schema models, etc.) is not intended to be imported by users. It may change between releases without notice.

When reviewing or writing code, enforce this boundary: user task code and BundleBuilder subclasses must only import from org.apache.airflow.sdk; any import of org.apache.airflow.sdk.execution.* in user-facing API surface is a red flag.


Bundle composition and coordinator discovery

A bundle is a directory of JAR files (typically build/bundle/) placed on the coordinator's jars_root. The coordinator scans the directory at task-dispatch time to find:

  1. Main-Class (standard JAR manifest attribute) — the fully-qualified class name of the entry point that the coordinator invokes with java -classpath … <Main-Class> --comm … --logs …. This must be a class with a public static void main(String[] args) method; the Gradle plugin org.apache.airflow.sdk writes it automatically from airflowBundle { mainClass = "…" } and validates that the class exists and has the right signature at build time.

  2. Airflow-Supervisor-Schema-Version (Airflow-specific manifest attribute) — the wire protocol version the JVM side expects when talking to the Python supervisor. In fat-JAR mode (the default), the Gradle plugin reads this value from the airflow-sdk JAR in runtimeClasspath and copies it into the shadow JAR manifest. In thin-JAR mode (fatJar = false), the value stays in the airflow-sdk JAR deployed alongside the bundle JAR.

The Python coordinator (JavaCoordinator) scans every JAR under jars_root with _JarInfo.find(), reads META-INF/MANIFEST.MF out of each ZIP, and collects Main-Class and Airflow-Supervisor-Schema-Version from whichever JARs carry them. The resolved schema version is then passed as the schema_version return value from _build_execute_task_command, which the base SubprocessCoordinator uses to negotiate the supervisor wire protocol.

If main_class is set explicitly on the JavaCoordinator instance (via [sdk] coordinators kwargs), the scan uses it as a filter; otherwise the first JAR with a Main-Class attribute wins. Either way, Airflow-Supervisor-Schema-Version must be present in at least one JAR in jars_root or startup fails.


Key files to know

FilePurpose
java-sdk/sdk/.../Client.ktPublic API (Variables, Connections, XCom)
java-sdk/sdk/.../execution/Client.ktSupervisor wire calls
java-sdk/sdk/.../execution/Comm.kt4-byte-prefix MessagePack framing
java-sdk/sdk/.../Server.ktEntry-point; drives the execution loop
java-sdk/processor/.../BuilderProcessor.ktKapt annotation processor
java-sdk/plugin/.../AirflowSdkPlugin.ktGradle bundle plugin
task-sdk/.../coordinators/java/coordinator.pyPython side — spawns the JVM
task-sdk/.../schema/schema.jsonWire protocol definition (both sides)

Running tests

Always use ./gradlew from inside java-sdk/; never run Gradle via apt's gradle. See java-sdk/README.md#testing for the full list of Gradle commands.

For the Python coordinator, use Breeze (never pytest directly on the host):

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

End-to-end test suite:

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

Updating the Python coordinator

coordinator.py extends SubprocessCoordinator. The only method subclasses must implement is _build_execute_task_command, which returns (argv, schema_version). Look at the existing implementation for how jars_root, java_executable, jvm_args, and main_class are assembled into the command. Do not reach into the JVM process from Python beyond what this method provides.


Upgrading Supervisor Schema client

When upgrading to a newer Supervisor Schema version:

  • Regenerate models with ./gradlew generateJsonSchema2Pojo
  • Modify execution/Client.kt to handle changes

The java-sdk/README.md#contributing section walks through the full "adding a new Client method" sequence step by step.

來自 astronomer 的更多技能

airflow-state-store
astronomer
Persists task and asset state across retries and DAG runs using Airflow 3.3's AIP-103 key/value stores (`task_state_store`, `asset_state_store`) and the…
creating-openlineage-extractors
astronomer
為不支援的Airflow運算子及複雜血緣場景設計的自訂OpenLineage提取器。提供兩種方法:直接在你擁有的運算子中加入OpenLineage方法(建議做法),或為無法修改的第三方運算子建立自訂提取器。提取器在三個時間點攔截運算子執行:執行前取得靜態血緣、成功後取得執行階段決定的輸出、以及選擇性地在失敗後取得部分血緣。可透過airflow.cfg或環境變數註冊提取器...
debugging-dags
astronomer
針對失敗的 Airflow DAG 進行系統性根本原因分析與修復,並提供結構化的調查流程。引導完成四個階段的診斷步驟:識別失敗、提取錯誤細節、收集背景資訊,以及提供可行的修復步驟。將失敗分為四種類型(資料、程式碼、基礎設施、相依性),以聚焦調查並建議適當的修正方式。提供可直接使用的 CLI 指令,用於日誌擷取、執行比較、任務清除與 DAG...
delegating-to-otto
astronomer
驅動 Astronomer 的 Otto 代理
deploying-airflow
astronomer
部署 Airflow DAG 和專案。當使用者想要部署程式碼、推送 DAG、設定 CI/CD、部署到生產環境,或詢問部署策略時使用…
deploying-go-sdk-bundles
astronomer
建置、打包並部署已編譯的 Airflow Go SDK 套件,以便 ExecutableCoordinator 能執行它們。當使用者想要編譯 Go 任務套件、要求…時使用。
testing-dags
astronomer
針對Airflow DAG的反覆測試-除錯-修復循環,提供全面的失敗診斷。從af runs trigger-wait <dag_id>開始執行DAG並等待完成,無需預先檢查。失敗時,使用af runs diagnose獲取完整的失敗摘要,並透過af tasks logs檢查特定任務的錯誤細節。支援自訂配置、超時設定與重試機制;能處理成功、失敗及超時情境,並提供清晰的回應解讀。快速驗證功能亦已就緒...
tracing-downstream-lineage
astronomer
追蹤下游資料血緣,在修改資料表或DAG前評估變更影響。透過原始碼搜尋、檢視相依性及BI工具連線,識別目標資料表或DAG的直接消費者。建立完整的相依性樹狀圖,繪製從資料表到儀表板再到機器學習模型的所有下游影響。依關鍵性(關鍵、高、中、低)分類相依性,以優先處理利害關係人溝通與測試。產出包含風險評估、受影響範圍的影響報告。