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의 문제 해결. 4단계 진단 프로세스를 안내합니다: 실패 식별, 오류 세부 정보 추출, 컨텍스트 정보 수집, 실행 가능한 수정 단계 제공. 실패를 네 가지 유형(데이터, 코드, 인프라, 종속성)으로 분류하여 조사에 집중하고 적절한 수정을 제안합니다. 로그 검색, 실행 비교, 작업 정리, DAG...을 위한 즉시 사용 가능한 CLI 명령을 제공합니다.
delegating-to-otto
astronomer
Drives Astronomer's Otto agent (`astro otto`) as a delegated sub-agent for Airflow, dbt, and data-engineering work. Use when the user explicitly asks to "use…
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의 직접적인 소비자를 식별합니다. 테이블에서 대시보드, ML 모델에 이르기까지 모든 다운스트림 영향을 매핑하는 전체 종속성 트리를 구축합니다. 종속성을 중요도(심각, 높음, 중간, 낮음)별로 분류하여 이해관계자 커뮤니케이션 및 테스트의 우선순위를 지정합니다. 위험 평가, 영향을 받는 항목이 포함된 영향 보고서를 생성합니다...