authoring-java-sdk-tasks

作者: astronomer

使用 Airflow Java SDK 以 Java、Kotlin 或任何 JVM 語言編寫 Airflow 任務邏輯。當使用者想要以 Java/JVM 實作 Airflow 任務時使用,詢問…

npx skills add https://github.com/astronomer/agents --skill authoring-java-sdk-tasks

Authoring Java SDK Tasks

The Airflow Java SDK implements the language-SDK model for the JVM: your DAG stays in Python, and each task instance runs in a short-lived JVM subprocess. This skill covers the Java-specific native API. The shared model — the Python @task.stub pattern, ID matching, and the XCom-as-JSON contract — lives in authoring-language-sdk-tasks; read that first if you're new to language SDKs.

Experimental. The Java SDK is in preview. Artifact coordinates and APIs may change.

Related skills: authoring-language-sdk-tasks (shared Python stub + concepts), configuring-airflow-language-sdks (route the queue to JavaCoordinator), deploying-java-sdk-bundles (compile and ship the JAR).


Recap: the Python side

Java tasks are paired with Python stubs that carry no logic — they declare the task, queue, dependency graph, and retries. IDs must match the Java annotations exactly, and an upstream argument on a stub only declares the dependency (the value is fetched in Java). Full rules are in authoring-language-sdk-tasks; the minimal shape:

from airflow.sdk import dag, task


@dag
def sales_pipeline():                     # dag_id "sales_pipeline" -> @Builder.Dag(id="sales_pipeline")
    @task.stub(queue="java")
    def extract(): ...                    # task_id "extract" -> @Builder.Task(id="extract")

    @task.stub(queue="java")
    def transform(extracted): ...

    transform(extract())


sales_pipeline()

Java side: two APIs

Both APIs produce identical runtime behavior; pick by style, and you can mix them in one bundle.

Annotation-based API (recommended)

Annotate a plain class; an annotation processor generates the wiring (<ClassName>Builder) at compile time.

import static java.lang.System.Logger.Level.INFO;
import org.apache.airflow.sdk.*;

@Builder.Dag(id = "sales_pipeline")          // must match the Python dag_id
public class SalesPipeline {
  private static final System.Logger log = System.getLogger(SalesPipeline.class.getName());

  @Builder.Task(id = "extract")              // must match the Python @task.stub name
  public long extract(Client client) {
    var conn = client.getConnection("sales_db");
    log.log(INFO, "connected to {0}", conn.host);
    return 42L;                              // return value is pushed as the return_value XCom
  }

  @Builder.Task(id = "transform")
  public long transform(
      Client client,
      @Builder.XCom(task = "extract") long recordCount) {  // pulls extract's return_value
    var threshold = (String) client.getVariable("transform_threshold");
    return recordCount * 2;
  }

  @Builder.Task   // id omitted -> the method name "load" is used
  public void load(Context context, @Builder.XCom(task = "transform") long transformed) {
    log.log(INFO, "attempt {0}, value {1}", context.ti.tryNumber, transformed);
  }
}

Annotation reference:

AnnotationPurpose
@Builder.Dag(id = "...")Marks the class as a task container. id must match the Python dag_id; if omitted, the class name is used. Optional to = "..." renames the generated builder (default <ClassName>Builder).
@Builder.Task(id = "...")Marks a method as a task. id must match the Python @task.stub function name; if omitted, the method name is used.
@Builder.XCom(task = "...", key = "...")Injects an upstream task's XCom as a parameter. task defaults to the parameter name; key defaults to the producing task's return_value. The parameter type must be compatible with the stored JSON value.

A task method's return value is automatically pushed as that task's return_value XCom. A method may declare throws Exception; any uncaught exception fails the task instance (which triggers retries if the stub configured them).

Interface-based API

Implement Task directly when you want full control over registration and XCom handling.

import org.apache.airflow.sdk.*;

public class ExtractTask implements Task {
  @Override
  public void execute(Context context, Client client) throws Exception {
    var conn = client.getConnection("sales_db");
    // ... do work ...
    client.setXCom(42L);   // push return_value explicitly
  }
}

Register tasks manually in a Dag and expose it through a BundleBuilder:

public class MyBundle implements BundleBuilder {
  @Override
  public Iterable<Dag> getDags() {
    var dag = new Dag("sales_pipeline");      // DAG ID matches Python
    dag.addTask("extract", ExtractTask.class);
    dag.addTask("transform", TransformTask.class);
    return java.util.List.of(dag);
  }
}

Each Task class needs a public no-arg constructor. Task IDs must be unique within a DAG, and DAG IDs unique within a bundle.


The entry point

Every bundle has a main that hands your DAGs to the SDK server. The server connects to the coordinator, runs one task instance, and exits.

import java.util.List;
import org.apache.airflow.sdk.*;

public class Main implements BundleBuilder {
  @Override
  public Iterable<Dag> getDags() {
    // With the annotation API, the *Builder classes are generated at compile time.
    return List.of(SalesPipelineBuilder.build());
  }

  public static void main(String[] args) {
    Server.create(args).serve(new Main().build());
  }
}

Server.create(args) parses the connection details Airflow passes on the command line — don't construct them by hand. Record this main class as the bundle's main class when you build it (see deploying-java-sdk-bundles).


Talking to Airflow from a task: Client

A Client is passed into every task and is scoped to the current DAG run and task instance.

CallReturnsNotes
client.getConnection(id)ConnectionFields: id, type, host, schema, login, password, port, extra. Any unset field is null. Throws if the connection doesn't exist.
client.getVariable(key)Object (or null)Cast to the type you expect, e.g. (String) client.getVariable("threshold").
client.getXCom(taskId)Object (or null)Reads another task's return_value by default. Overloads accept key, dagId, runId, mapIndex, and includePriorDates for cross-DAG/run reads and mapped tasks.
client.setXCom(value)—Pushes the return_value XCom (interface API). Value must be JSON-serializable. With the annotation API, returning a value does this for you.

Context

The Context parameter exposes run metadata: context.dagRun (dagId, runId) and context.ti (dagId, runId, taskId, mapIndex, tryNumber). tryNumber is useful for retry-aware logic.


XCom: Java types

XComs cross the boundary as JSON (the shared contract is in authoring-language-sdk-tasks). When you read one back in Java you get:

Python typeJSONJava type from getXCom
intintegerLong (or BigInteger if too large)
floatdecimalDouble
strstringString
boolbooleanBoolean
Nonenullnull
listarrayList<Object>
dictobjectMap<String, Object>

Declare @Builder.XCom parameter types to match. A mismatch (e.g. declaring int when the value is a String) fails the task.


Logging

Declare a logger as a static field named after the class — the conventional pattern regardless of framework:

private static final System.Logger log = System.getLogger(SalesPipeline.class.getName());

For records to reach Airflow's task log store (and show in the UI), the bundle must include one of the SDK logging integration artifacts (airflow-sdk-jpl, airflow-sdk-slf4j, airflow-sdk-log4j2, or airflow-sdk-jul). The dependencies and per-framework setup are in the logging integration section of deploying-java-sdk-bundles. System.Logger (JPL) with airflow-sdk-jpl is the lightest option and needs no configuration.


A complete worked example ships with the SDK

The SDK repository includes a runnable example under java-sdk/example/:

  • src/resources/dags/java_examples.py — Python DAGs pairing Python tasks with Java stubs, including a load stub with retries=1.
  • src/java/.../AnnotationExample.java — annotation API, including a task that fails on tryNumber == 1 and succeeds on retry.
  • src/java/.../InterfaceExampleBuilder.java — the same tasks via the Task interface and Dag.addTask(...).
  • src/java/.../ExampleBundleBuilder.java — a BundleBuilder returning both DAGs plus the main entry point.

Point users there for an end-to-end reference.


Java-specific pitfalls

  • Cast Object returns deliberately. getVariable and getXCom return Object; match the cast to the JSON type (see the table above).
  • @Builder.XCom parameter types must match the stored JSON type, or the task fails at runtime.
  • The annotation processor must be on the build for the annotation API (generates <ClassName>Builder); it is not needed for the interface API. See deploying-java-sdk-bundles.
  • See authoring-language-sdk-tasks for the language-agnostic pitfalls (ID matching, one JVM per task instance, queue/retries on the stub).

Related Skills

  • authoring-language-sdk-tasks: Shared Python-stub pattern and concepts (read first).
  • configuring-airflow-language-sdks: Route the java queue to JavaCoordinator and set JRE/coordinator options.
  • deploying-java-sdk-bundles: Build the bundle (Gradle/Maven) and place the JAR where Airflow can find it.
  • authoring-dags: General Airflow DAG authoring.

來自 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的直接消費者。建立完整的相依性樹狀圖,繪製從資料表到儀表板再到機器學習模型的所有下游影響。依關鍵性(關鍵、高、中、低)分類相依性,以優先處理利害關係人溝通與測試。產出包含風險評估、受影響範圍的影響報告。