authoring-java-sdk-tasks

Menulis logika tugas Airflow dalam Java, Kotlin, atau bahasa JVM lainnya menggunakan Airflow Java SDK. Gunakan saat pengguna ingin mengimplementasikan tugas Airflow dalam Java/JVM, meminta…

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.

Lebih banyak skill dari astronomer

airflow
astronomer
Kueri, kelola, dan pecahkan masalah DAG, proses, tugas, serta konfigurasi sistem Apache Airflow. Mendukung 30+ perintah untuk inspeksi DAG, manajemen proses, pencatatan tugas, kueri konfigurasi, dan akses langsung REST API. Kelola beberapa instance Airflow dengan konfigurasi persisten; temukan secara otomatis deployment lokal dan Astro. Jalankan proses DAG secara sinkron (tunggu hingga selesai) atau asinkron, diagnosis kegagalan, hapus proses untuk percobaan ulang, dan akses log tugas dengan filter percobaan ulang/indeks peta. Keluaran...
official
airflow-hitl
astronomer
Gerbang persetujuan manusia, input formulir, dan percabangan dalam DAG Airflow menggunakan operator yang dapat ditunda. Empat jenis operator: ApprovalOperator untuk keputusan setuju/tolak, HITLOperator untuk pemilihan multi-opsi dengan formulir, HITLBranchOperator untuk perutean tugas yang digerakkan manusia, dan HITLEntryOperator untuk pengumpulan data formulir. Semua operator dapat ditunda, membebaskan slot pekerja sambil menunggu respons manusia melalui tab Required Actions di UI Airflow atau REST API. Mendukung fitur opsional termasuk kustom...
official
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…
official
analyzing-data
astronomer
Kueri gudang data Anda untuk menjawab pertanyaan bisnis dengan pola yang di-cache dan pemetaan konsep. Mendukung pencarian pola dan caching untuk jenis pertanyaan berulang, dengan pencatatan hasil untuk meningkatkan kueri di masa mendatang. Menyertakan cache pemetaan konsep-ke-tabel dan penemuan skema tabel melalui INFORMATION_SCHEMA atau grep basis kode. Menyediakan fungsi kernel run_sql() dan run_sql_pandas() yang mengembalikan DataFrame Polars atau Pandas untuk analisis. Perintah CLI untuk mengelola cache konsep, pola, dan tabel, plus...
official
annotating-task-lineage
astronomer
Anotasi tugas Airflow dengan lineage data menggunakan inlet dan outlet. Mendukung objek Dataset OpenLineage, Aset Airflow, dan Dataset Airflow untuk mendefinisikan input dan output di seluruh basis data, gudang data, dan penyimpanan cloud. Digunakan sebagai cadangan ketika operator tidak memiliki ekstraktor OpenLineage bawaan; mengikuti sistem prioritas empat tingkat di mana ekstraktor kustom dan metode OpenLineage diutamakan. Menyertakan pembantu penamaan dataset untuk Snowflake, BigQuery, S3, dan PostgreSQL guna memastikan konsistensi...
official
authoring-dags
astronomer
Panduan kerja untuk membuat DAG Apache Airflow dengan integrasi validasi dan pengujian. Pendekatan enam fase terstruktur: temukan lingkungan dan pola yang ada, rencanakan struktur DAG, implementasikan sesuai praktik terbaik, validasi dengan perintah CLI af, uji dengan persetujuan pengguna, dan lakukan iterasi perbaikan. Perintah CLI untuk penemuan (af config connections, af config providers, af dags list) dan validasi (af dags errors, af dags get, af dags explore) memberikan umpan balik langsung pada DAG...
official
authoring-go-sdk-tasks
astronomer
Writes Airflow task logic in Go using the Airflow Go SDK. Use when the user wants to implement Airflow tasks in Go, asks about `BundleProvider`/`RegisterDags`,…
official
authoring-language-sdk-tasks
astronomer
Fondasi netral bahasa untuk SDK bahasa Airflow — implementasikan logika tugas dalam bahasa non-Python sementara DAG tetap dalam Python. Gunakan ketika pengguna…
official