← astronomer-dataCONTENT HISTORY

Update to astronomer-data

Snapshot Sep 30, 2026 · 23:17 UTC · version 0.1.0

Collection source: not recorded for this historical snapshot.

WHAT CHANGED · RULE-BASED ANALYSIS

First saved snapshot

No earlier snapshot is available to establish a change.

Compare saved observations

Download comparison JSON
Full technical diff · 0 changed fields
Full snapshot data
{
  "name": "authoring-java-sdk-tasks",
  "description": "Writes Airflow task logic in Java, Kotlin, or any JVM language using the Airflow Java SDK. Use when the user wants to implement Airflow tasks in Java/JVM, asks about `@Builder.Dag`/`@Builder.Task`/`@Builder.XCom`, the `Task`/`BundleBuilder` interfaces, reading connections/variables/XComs from Java, the JSON-to-Java type mapping, or logging from Java tasks. This skill covers the Java-specific native API; the shared Python-stub pattern and conceptual model live in authoring-language-sdk-tasks. For building/shipping the bundle see deploying-java-sdk-bundles; for coordinator config see configuring-airflow-language-sdks.",
  "included_files": [],
  "skill_md_contents": "---\nname: authoring-java-sdk-tasks\ndescription: Writes Airflow task logic in Java, Kotlin, or any JVM language using the Airflow Java SDK. Use when the user wants to implement Airflow tasks in Java/JVM, asks about `@Builder.Dag`/`@Builder.Task`/`@Builder.XCom`, the `Task`/`BundleBuilder` interfaces, reading connections/variables/XComs from Java, the JSON-to-Java type mapping, or logging from Java tasks. This skill covers the Java-specific native API; the shared Python-stub pattern and conceptual model live in authoring-language-sdk-tasks. For building/shipping the bundle see deploying-java-sdk-bundles; for coordinator config see configuring-airflow-language-sdks.\n---\n\n# Authoring Java SDK Tasks\n\nThe 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.\n\n> **Experimental.** The Java SDK is in preview. Artifact coordinates and APIs may change.\n\n> **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).\n\n---\n\n## Recap: the Python side\n\nJava 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:\n\n```python\nfrom airflow.sdk import dag, task\n\n\n@dag\ndef sales_pipeline():                     # dag_id \"sales_pipeline\" -> @Builder.Dag(id=\"sales_pipeline\")\n    @task.stub(queue=\"java\")\n    def extract(): ...                    # task_id \"extract\" -> @Builder.Task(id=\"extract\")\n\n    @task.stub(queue=\"java\")\n    def transform(extracted): ...\n\n    transform(extract())\n\n\nsales_pipeline()\n```\n\n---\n\n## Java side: two APIs\n\nBoth APIs produce identical runtime behavior; pick by style, and you can mix them in one bundle.\n\n### Annotation-based API (recommended)\n\nAnnotate a plain class; an annotation processor generates the wiring (`<ClassName>Builder`) at compile time.\n\n```java\nimport static java.lang.System.Logger.Level.INFO;\nimport org.apache.airflow.sdk.*;\n\n@Builder.Dag(id = \"sales_pipeline\")          // must match the Python dag_id\npublic class SalesPipeline {\n  private static final System.Logger log = System.getLogger(SalesPipeline.class.getName());\n\n  @Builder.Task(id = \"extract\")              // must match the Python @task.stub name\n  public long extract(Client client) {\n    var conn = client.getConnection(\"sales_db\");\n    log.log(INFO, \"connected to {0}\", conn.host);\n    return 42L;                              // return value is pushed as the return_value XCom\n  }\n\n  @Builder.Task(id = \"transform\")\n  public long transform(\n      Client client,\n      @Builder.XCom(task = \"extract\") long recordCount) {  // pulls extract's return_value\n    var threshold = (String) client.getVariable(\"transform_threshold\");\n    return recordCount * 2;\n  }\n\n  @Builder.Task   // id omitted -> the method name \"load\" is used\n  public void load(Context context, @Builder.XCom(task = \"transform\") long transformed) {\n    log.log(INFO, \"attempt {0}, value {1}\", context.ti.tryNumber, transformed);\n  }\n}\n```\n\nAnnotation reference:\n\n| Annotation | Purpose |\n|------------|---------|\n| `@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`). |\n| `@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. |\n| `@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. |\n\nA 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).\n\n### Interface-based API\n\nImplement `Task` directly when you want full control over registration and XCom handling.\n\n```java\nimport org.apache.airflow.sdk.*;\n\npublic class ExtractTask implements Task {\n  @Override\n  public void execute(Context context, Client client) throws Exception {\n    var conn = client.getConnection(\"sales_db\");\n    // ... do work ...\n    client.setXCom(42L);   // push return_value explicitly\n  }\n}\n```\n\nRegister tasks manually in a `Dag` and expose it through a `BundleBuilder`:\n\n```java\npublic class MyBundle implements BundleBuilder {\n  @Override\n  public Iterable<Dag> getDags() {\n    var dag = new Dag(\"sales_pipeline\");      // DAG ID matches Python\n    dag.addTask(\"extract\", ExtractTask.class);\n    dag.addTask(\"transform\", TransformTask.class);\n    return java.util.List.of(dag);\n  }\n}\n```\n\nEach `Task` class needs a public no-arg constructor. Task IDs must be unique within a DAG, and DAG IDs unique within a bundle.\n\n---\n\n## The entry point\n\nEvery bundle has a `main` that hands your DAGs to the SDK server. The server connects to the coordinator, runs one task instance, and exits.\n\n```java\nimport java.util.List;\nimport org.apache.airflow.sdk.*;\n\npublic class Main implements BundleBuilder {\n  @Override\n  public Iterable<Dag> getDags() {\n    // With the annotation API, the *Builder classes are generated at compile time.\n    return List.of(SalesPipelineBuilder.build());\n  }\n\n  public static void main(String[] args) {\n    Server.create(args).serve(new Main().build());\n  }\n}\n```\n\n`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**).\n\n---\n\n## Talking to Airflow from a task: `Client`\n\nA `Client` is passed into every task and is scoped to the current DAG run and task instance.\n\n| Call | Returns | Notes |\n|------|---------|-------|\n| `client.getConnection(id)` | `Connection` | Fields: `id`, `type`, `host`, `schema`, `login`, `password`, `port`, `extra`. Any unset field is `null`. Throws if the connection doesn't exist. |\n| `client.getVariable(key)` | `Object` (or `null`) | Cast to the type you expect, e.g. `(String) client.getVariable(\"threshold\")`. |\n| `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. |\n| `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. |\n\n### `Context`\n\nThe `Context` parameter exposes run metadata: `context.dagRun` (`dagId`, `runId`) and `context.ti` (`dagId`, `runId`, `taskId`, `mapIndex`, `tryNumber`). `tryNumber` is useful for retry-aware logic.\n\n---\n\n## XCom: Java types\n\nXComs cross the boundary as JSON (the shared contract is in **authoring-language-sdk-tasks**). When you read one back in Java you get:\n\n| Python type | JSON | Java type from `getXCom` |\n|-------------|------|--------------------------|\n| `int` | integer | `Long` (or `BigInteger` if too large) |\n| `float` | decimal | `Double` |\n| `str` | string | `String` |\n| `bool` | boolean | `Boolean` |\n| `None` | null | `null` |\n| `list` | array | `List<Object>` |\n| `dict` | object | `Map<String, Object>` |\n\nDeclare `@Builder.XCom` parameter types to match. A mismatch (e.g. declaring `int` when the value is a `String`) fails the task.\n\n---\n\n## Logging\n\nDeclare a logger as a static field named after the class — the conventional pattern regardless of framework:\n\n```java\nprivate static final System.Logger log = System.getLogger(SalesPipeline.class.getName());\n```\n\nFor 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.\n\n---\n\n## A complete worked example ships with the SDK\n\nThe SDK repository includes a runnable example under `java-sdk/example/`:\n\n- `src/resources/dags/java_examples.py` — Python DAGs pairing Python tasks with Java stubs, including a `load` stub with `retries=1`.\n- `src/java/.../AnnotationExample.java` — annotation API, including a task that fails on `tryNumber == 1` and succeeds on retry.\n- `src/java/.../InterfaceExampleBuilder.java` — the same tasks via the `Task` interface and `Dag.addTask(...)`.\n- `src/java/.../ExampleBundleBuilder.java` — a `BundleBuilder` returning both DAGs plus the `main` entry point.\n\nPoint users there for an end-to-end reference.\n\n---\n\n## Java-specific pitfalls\n\n- **Cast `Object` returns deliberately.** `getVariable` and `getXCom` return `Object`; match the cast to the JSON type (see the table above).\n- **`@Builder.XCom` parameter types must match the stored JSON type**, or the task fails at runtime.\n- **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**.\n- See **authoring-language-sdk-tasks** for the language-agnostic pitfalls (ID matching, one JVM per task instance, queue/retries on the stub).\n\n---\n\n## Related Skills\n\n- **authoring-language-sdk-tasks**: Shared Python-stub pattern and concepts (read first).\n- **configuring-airflow-language-sdks**: Route the `java` queue to `JavaCoordinator` and set JRE/coordinator options.\n- **deploying-java-sdk-bundles**: Build the bundle (Gradle/Maven) and place the JAR where Airflow can find it.\n- **authoring-dags**: General Airflow DAG authoring.\n"
}

SHA-256: 6645181954fab36e65c03f1608af77619ad9697d7337580c7433e94857c142e6