Java SDK

This is an experimental feature.

The Java SDK lets you implement Airflow task logic in Java, Kotlin, or any other JVM language. The Dag and its scheduling remain in Python; individual tasks delegate to a JVM subprocess that is spawned by JavaCoordinator for each task instance.

API reference

The generated API reference for the Java SDK is published with the Airflow documentation at Java SDK API Reference.

Prerequisites

  • JDK 11 or later is required on the machine that builds the Java project. A local Gradle installation is only needed to generate the Gradle Wrapper for a new project.

  • JRE 11 or later must be available on the Airflow worker nodes and the Dag processor.

  • The compiled task JAR(s) and JVM dependencies must be accessible from the worker and the Dag processor.

  • The apache-airflow-task-sdk package (installed with Airflow) provides the coordinator; no additional Python packages are needed.

Quick start

This example uses the annotation-based API to build a Java bundle for two stub tasks. The annotation-based and interface-based APIs are two ways to define Java tasks, not two bundle formats. They use different Java code and dependencies, but the standard Gradle source layout, entry-point requirement, bundle task, and deployment process are the same. See Interface-based API for the interface-based task code.

The Python Dag source and the Java Gradle project are independent. They do not need to be in the same repository or have any particular relative filesystem layout. The Dag follows the deployment’s normal Dag delivery process; only the compiled Java bundle is deployed from the Gradle project, into a separate Dag bundle that holds the JARs.

Define the Python Dag

For a local installation with the default [core] dags_folder, create ${AIRFLOW_HOME}/dags/sales_pipeline.py. More generally, create sales_pipeline.py in the source location used by the deployment’s normal Dag delivery process, such as its configured Dags folder, Dag repository, or Dag bundle. This path is not relative to the Java Gradle project.

from airflow.sdk import dag, task


@dag
def sales_pipeline():
    @task.stub(queue="java")
    def extract(): ...

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

    @task()
    def load(transformed):
        print(f"Loaded: {transformed}")

    load(transform(extract()))


sales_pipeline()

Create the Java project

The guide uses com.mycompany.airflow.sales as a placeholder Java package. Replace it with the package for your project, and update the source paths, package declarations, and mainClass together.

Start with this standard Gradle project layout. The gradlew scripts and gradle/wrapper/ directory are generated by the Gradle Wrapper.

sales-pipeline-java/
├── build.gradle
├── gradle.properties
├── settings.gradle
├── gradlew
├── gradlew.bat
├── gradle/
│   └── wrapper/
│       ├── gradle-wrapper.jar
│       └── gradle-wrapper.properties
└── src/
    └── main/
        └── java/
            └── com/
                └── mycompany/
                    └── airflow/
                        └── sales/
                            ├── Main.java
                            └── SalesPipeline.java

Choose a Java SDK version published to Maven Central, then set it in gradle.properties by replacing JAVA_SDK_VERSION:

airflowJavaSdkVersion=JAVA_SDK_VERSION

The Java SDK is experimental, so Maven Central may list only prerelease versions. Select a version deliberately rather than copying a version that may become stale in this guide.

Name the project in settings.gradle:

rootProject.name = "sales-pipeline"

Configure the build in build.gradle:

plugins {
    id("org.apache.airflow.sdk") version "${airflowJavaSdkVersion}"
}

repositories {
    mavenCentral()
}

dependencies {
    annotationProcessor("org.apache.airflow:airflow-sdk-processor:${airflowJavaSdkVersion}")
    implementation("org.apache.airflow:airflow-sdk:${airflowJavaSdkVersion}")
}

java {
    toolchain {
        languageVersion.set(JavaLanguageVersion.of(11))
    }
}

airflowBundle {
    mainClass = "com.mycompany.airflow.sales.Main"
}

The annotationProcessor dependency is required for this annotation-based example. Omit it when using the interface-based API. If the project does not have the Gradle Wrapper yet, generate it from the project root with gradle wrapper.

Implement the Java tasks

Save the task implementations as src/main/java/com/mycompany/airflow/sales/SalesPipeline.java:

package com.mycompany.airflow.sales;

import org.apache.airflow.sdk.Builder;

public class SalesPipeline {

  @Builder.TaskHandler(dag = "sales_pipeline", task = "extract")
  public long extract() {
    return 3;
  }

  @Builder.TaskHandler(dag = "sales_pipeline", task = "transform")
  public long transform(long recordCount) {
    return recordCount * 2;
  }
}

Note

The graph is declared once, in the Python Dag file: transform(extract()) feeds the upstream’s return value into the downstream’s parameter by calling tasks like functions. The supervisor sends the resulting argument bindings to the Java runtime, and each Java data parameter receives whatever the Python call site bound at its position — an upstream task’s XCom or an inline literal. See Binding stub arguments.

Add the Java entry point

Save the entry point as src/main/java/com/mycompany/airflow/sales/Main.java:

package com.mycompany.airflow.sales;

import org.apache.airflow.sdk.Bundle;
import org.apache.airflow.sdk.Server;

public class Main {
  public static void main(String[] args) {
    Server.create(args).serve(new Bundle().register(SalesPipeline.class));
  }
}

register takes the class the annotations are on, so there is no second name to keep in sync: the annotation processor generates the registrar it reads during compilation.

Build and deploy

Run the bundle task from the sales-pipeline-java/ project root:

./gradlew bundle

The deployable JAR or JARs are now in build/bundle/. Copy or mount that entire directory into a path available on every worker that can consume the java queue. For example, deploy it as /opt/airflow/jars/sales-pipeline/. Java source files do not need to be deployed to Airflow.

For a local deployment where those paths are writable, the copy commands could be:

mkdir -p /opt/airflow/jars/sales-pipeline
cp build/bundle/*.jar /opt/airflow/jars/sales-pipeline/

Deploy sales_pipeline.py separately through the deployment’s normal Dag delivery process. For example, that process might sync it to ${AIRFLOW_HOME}/dags/ or package it in a Dag bundle; neither location is inside or relative to sales-pipeline-java/.

Configure Airflow to register the JAR directory as a Dag bundle, point the coordinator at it, and route the java queue to the coordinator. Add the following sections to the file selected by AIRFLOW_CONFIG (by default, ${AIRFLOW_HOME}/airflow.cfg), or set the equivalent AIRFLOW__* environment variables:

[dag_processor]
dag_bundle_config_list = [
    {"name": "dags-folder", "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", "kwargs": {}},
    {
      "name": "java-task-handlers",
      "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle",
      "kwargs": {"path": "/opt/airflow/jars/sales-pipeline"}
    }
  ]

[sdk]
coordinators = {
  "java": {
    "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
    "kwargs": {"task_handler_bundle_name": "java-task-handlers"}
  }
}
queue_to_coordinator = {"java": "java"}

java is a user-chosen coordinator name, not a reserved value. The value assigned to the queue in queue_to_coordinator must match a key in coordinators, and task_handler_bundle_name must match a Dag bundle name in dag_bundle_config_list. See JavaCoordinator configuration for how JARs are located.

Restart the affected Airflow components after changing this configuration. The coordinator config, the JARs and a JRE must be available wherever tasks execute and on the Dag processor. With CeleryExecutor, tasks execute on the Celery workers; with LocalExecutor, they run in subprocesses on the scheduler’s host. The Dag processor checks the stub tasks of sales_pipeline.py against the task handlers the JARs register, so it runs them too. The API server does not need any of it. Register the Dag bundle in [dag_processor] dag_bundle_config_list on every component, like your other Dag bundles: the worker and the Dag processor resolve task_handler_bundle_name through it, and wherever the [sdk] config is read it is rejected if the name is missing there. The Dag processor still receives sales_pipeline.py through the separate Dag delivery process.

A Dag processor with this [sdk] configuration also parses the executable JARs of every Dag bundle, and needs a JRE to do so (see Parsing native Java Dags). Keep it from parsing a bundle’s JARs with that bundle’s .airflowignore, whose pattern syntax follows [core] dag_ignore_file_syntax.

With task_handler_bundle_name set to its own bundle, as above, java-task-handlers holds only the JARs the Python Dag’s tasks run, so ignoring everything in it is safe:

# [core] dag_ignore_file_syntax = glob (the default)
echo '*' > /opt/airflow/jars/sales-pipeline/.airflowignore

# [core] dag_ignore_file_syntax = regexp; a bare * is not a valid pattern and is dropped
echo '.' > /opt/airflow/jars/sales-pipeline/.airflowignore

With task_handler_bundle_name unset, the JAR sits in the same Dag bundle as sales_pipeline.py, so ignore only the JAR, not the whole bundle:

# [core] dag_ignore_file_syntax = glob (the default)
echo '*.jar' >> /opt/airflow/dags/.airflowignore

# [core] dag_ignore_file_syntax = regexp
echo '\.jar$' >> /opt/airflow/dags/.airflowignore

After Airflow has parsed the Dag, trigger it from the UI or command line:

airflow dags trigger sales_pipeline

The Java extract task returns 3, transform returns 6 through XCom, and the Python load task logs Loaded: 6.

See JavaCoordinator configuration for the full list of accepted kwargs. For a larger example that exercises connections, variables, logging, and both task APIs, see the Java SDK example in the Airflow repository.

Writing tasks

The Java SDK offers two APIs for implementing tasks. Both produce the same runtime behavior; the choice is a matter of style.

Annotation-based API

Annotate a plain Java class and let the SDK generate the boilerplate at compile time.

Annotation

Purpose

@Builder.TaskHandler(dag = "...", task = "...")

Marks a method as the Java body of a task the Python Dag file declares with @task.stub. dag must match the dag_id and task the stub function name; omitting task derives it from the method name. There is no class-level annotation on this surface — the Dag is the Python file’s, so the handler names the pair it binds to.

@Builder.Dag(id = "...")

Marks a class as a Dag that Java itself owns. Attributes (schedule, description, tags, catchup, …) are Airflow’s own Dag settings; only attributes written explicitly are applied. See Native Java Dags.

@Builder.Task(id = "...")

Marks a method as a task of a Java-owned Dag. If id is omitted the method name is used. Further attributes (retries, queue, retryDelay, …) are Airflow’s own task settings; only attributes written explicitly are applied.

@Builder.Deps

Marks the nested class that declares the task graph in Java, TaskFlow-style. Required for a Dag that Java owns end to end. See Native Java Dags.

TaskInput / @ArgName("...")

Marks a class as a task’s input, so keyword arguments bind by name instead of by position: each public field receives the argument whose name matches it, ignoring case and underscores. @ArgName pins a name the match cannot reach, or renames the argument deliberately. See Binding stub arguments.

Besides the annotations, a task method may declare a Client and a Context parameter in any position; the SDK injects both. Every other parameter is a data parameter and receives an argument bound by the Python @task.stub call site.

The annotation processor generates a <ClassName>Builder class that wires up the task registry and resolves data parameters and XCom pushes automatically.

public class MyDag {

  @Builder.TaskHandler(dag = "my_dag", task = "fetch")
  public String fetch(Client client) throws Exception {
    var conn = client.getConnection("my_api");
    // implement task logic
    return result;
  }

  @Builder.TaskHandler(dag = "my_dag", task = "process")
  public long process(Client client, String fetched) {
    var threshold = (String) client.getVariable("process_threshold");
    // implement task logic
    return count;
  }
}

A task method may declare throws Exception; any uncaught exception causes the task instance to be marked as failed in Airflow (triggering retries if configured on the stub).

Interface-based API

Implement the Task interface directly for full control over how tasks are registered and how XComs are read. Each task is registered as a TaskDef on a DagDef; both carry a fluent config(key, value) whose keys are Airflow’s own setting names.

The runner creates a fresh instance of the task class through reflection for every task-instance run, which puts four constraints on the class:

  • The task class itself must be public.

  • It must be concrete: not abstract and not an interface.

  • It must declare a public no-argument constructor.

  • If nested inside another class, it must be a static nested class.

A class that violates any of these fails at runtime with a Cannot instantiate task class error in the task log.

import org.apache.airflow.sdk.*;

public class FetchTask implements Task {
  @Override
  public void execute(Context context, Client client) throws Exception {
    var conn = client.getConnection("my_api");
    // implement task logic
    client.setXCom(result);
  }
}

Implement InputTask<I> instead when the Python Dag calls the stub with TaskFlow arguments: the SDK resolves them from the call site and passes them in. The type argument is a TaskInput whose public fields declare what the task expects. See Binding stub arguments.

Register each task against the Dag and task the Python file declares. register is one overloaded verb: these ids, or the class an annotated handler lives on.

public class MyBundle {
  public static class ProcessInput implements TaskInput {
    public String fetched;
  }

  public static class ProcessTask implements InputTask<ProcessInput> {
    @Override
    public void execute(Context context, Client client, ProcessInput input) throws Exception {
      // implement task logic
      client.setXCom(input.fetched);
    }
  }

  public static void main(String[] args) {
    var bundle = new Bundle()
        .register("my_dag", "fetch", FetchTask.class)
        .register("my_dag", "process", ProcessTask.class);
    Server.create(args).serve(bundle);
  }
}

A task class can be top-level like FetchTask, or a nested static class like ProcessTask. Place them under the standard src/main/java/<package>/ source tree, and set airflowBundle.mainClass to the class that provides main. From that point onward, both APIs use the same ./gradlew bundle command and deploy the resulting build/bundle/ directory in the same way.

See the Java SDK API Reference for more details.

Binding stub arguments

Calling a @task.stub TaskFlow-style in the Python Dag is what declares the graph, and the supervisor delivers the resulting argument bindings to the Java runtime with every task run. A binding carries either an upstream task’s return_value XCom or an inline literal written at the call site.

Positional binding

A task method’s data parameters bind by position, in declaration order — the injected Client and Context parameters do not take up a position. Java parameter names are not part of the API, so renaming one in an IDE never rebinds an input.

@task.stub(queue="java")
def score(rows, threshold): ...


score(load_rows(), 0.75)
@Builder.TaskHandler(dag = "etl", task = "score")
public long score(Client client, long rows, double threshold) {
  // rows      <- the load_rows XCom (position 0)
  // threshold <- the literal 0.75   (position 1)
}

A primitive parameter cannot hold null, so the task fails with MissingXComException when its binding resolves to nothing; declare a boxed type (Long, Double, …) to receive null instead. The method must declare as many data parameters as the call site bound: positions carry the whole meaning of a flat binding, so any other count has already shifted them, and the task fails rather than running on arguments it has mistaken for others. A parameter the Python call omitted does not count towards that. Its default still arrives, but a method that does not declare it is not reading shifted arguments, so the SDK drops it before comparing the two counts.

Generic parameters are decoded element by element. Declare List<Double> and values.get(0) really is a Double, even though the call site passed whole numbers and the wire carries them as integers. Without that the list would hold Long values while claiming to hold Double, and the ClassCastException would land on whichever line first read an element rather than on the binding that got it wrong.

Named binding with a TaskInput

To bind keyword arguments by name, declare a class implementing TaskInput. Each of its public non-final fields receives the argument whose name matches it, ignoring case and underscores, so the stub’s snake_case arguments reach camelCase Java fields with nothing declared. That is the same fold the Go and TypeScript SDKs apply, so one Python signature binds identically in every SDK.

The class needs a public no-argument constructor. A name mismatch in either direction is logged rather than failed, because a field binds by name: a field nothing supplies keeps its Java default, and an argument no field claims changes nothing the task reads. Once a field has claimed its argument, an argument that resolves to nothing is a value and not a mistake, so a boxed or reference field takes null and a primitive field fails.

No two fields may claim argument names that differ only in case or underscores; that fails when the bundle is built, because the fold cannot tell them apart. Two arguments that collide that way reach only a field naming one of them exactly, and a field that would match both is left unfilled and reported rather than handed the wrong value.

@task.stub(queue="java")
def score(region_code, threshold): ...


score(region_code="emea", threshold=load_threshold())
public static class ScoreInput implements TaskInput {
  // Pinned so the field can be called region. Or drop the annotation and
  // write: public String regionCode;
  @ArgName("region_code")
  public String region;

  public double threshold; // binds threshold
}

@Builder.TaskHandler(dag = "etl", task = "score")
public long score(Client client, ScoreInput input) { ... }

Reach for @ArgName when the argument name is not a legal or usable Java identifier — a Python keyword such as class, say — or when the field should read differently from the argument. A pinned name is matched as written, with no folding.

A task method declares flat data parameters or one TaskInput, never both, so field names and flat positions cannot shift each other. Mixing them, or declaring two, fails the build.

Binding in the interface-based API

A task written against the interface has no parameter list to bind, so it declares its input as the type argument of InputTask<I> instead — the same TaskInput an annotated task would declare, binding the same way:

public class ScoreTask implements InputTask<ScoreInput> {
  @Override
  public void execute(Context context, Client client, ScoreInput input) throws Exception {
    // input.region, input.threshold
  }
}

A TaskInput is the only way an interface task receives bound values; there is no positional form. A positional read is safe in the code the annotation processor writes, because it type-checks each position against the method signature it serves, and is the wrong thing to ask of code a person writes and later reads. A call site with a single argument is worth the one-field class.

An InputTask whose type argument is not a concrete TaskInput fails when the bundle is built, rather than mid-run. Plain Task remains the right interface for a task the Dag file calls with no arguments.

Native Java Dags

A Dag can also be authored entirely in Java: the annotations (or the DagDef / TaskDef objects) carry the configuration, and Java declares the graph.

Building the Dag in Java

dag.task(...) registers a task as it creates it and hands back a handle. before and after draw every edge on this surface, and the task body moves the data itself, by reading the upstream’s XCom through Client:

var dag = new DagDef("java_etl");

var extract = dag.task("extract", Extract.class);
var transform = dag.task("transform", Transform.class);
var load = dag.task("load", Load.class);

transform.after(extract).before(load);

Both are variadic, so a.before(b, c) fans out and d.after(b, c) fans in, and both return their own receiver, so a chain reads from one task outwards. Flow.of(a, b).before(c, d), from org.apache.airflow.sdk.Deps.Flow, draws every edge between two sets in one call.

Edges are checked when the Dag is registered with a Bundle: an upstream that belongs to another Dag, or to no Dag, and a cycle anywhere in the graph both fail there rather than at the first task run.

Wiring the graph with @Builder.Deps

For a Dag written with annotations, the graph is declared by a nested @Builder.Deps class. The annotation processor generates a <ClassName>Deps interface, the wiring view, with one method per @Builder.Task method: the injected Client and Context parameters are dropped, each data parameter becomes an Arg<T>, and the return value becomes a TaskRef<T>. Calling a view method registers its task, and passing the handle one returned into another feeds the upstream’s output into the downstream’s parameter and wires the data edge. The call graph is the task graph, and javac type-checks it.

Declare the wiring class as a static nested class of the Dag class that implements the generated view, with a no-argument depends() method:

@Builder.Dag(
    id = "java_etl",
    schedule = "@daily",
    description = "Pure-Java Dag built with annotations",
    tags = {"example", "java-sdk"})
public class EtlPipeline {

  @Builder.Task(id = "extract", retries = 2)
  public long extract() {
    return 42L;
  }

  @Builder.Task(id = "transform")
  public long transform(long extracted, double factor) {
    return (long) (extracted * factor);
  }

  @Builder.Task(id = "load")
  public void load(long transformed) {
    // implement task logic
  }

  @Builder.Task(id = "audit")
  public void audit() {
    // side effect only, no data in or out
  }

  @Builder.Deps
  static class Wiring implements EtlPipelineDeps {
    void depends() {
      var rows = extract();
      load(transform(rows, lit(0.9)));
      rows.before(audit()); // ordering-only edge: audit waits for extract
    }
  }
}

Every @Builder.Task method must be called in the wiring class; a task the wiring missed fails at Dag-parse time. lit(...) wires an inline constant where no upstream feeds a parameter. A bare double cannot be an Arg, so a constant is wrapped. Airflow records what each task is called with, and that record travels as JSON, so a constant has to be a string, number, boolean, list or map. A view method that takes no arguments returns the same handle every time, so it names one node wherever it appears; one that takes arguments is called once, and the wiring fails if it is called again with arguments, so hold its handle in a local and reuse that.

before, after and Flow.of work here exactly as they do on the interface surface; inside the wiring class Flow is inherited by simple name, so it needs no import and never collides with java.util.concurrent.Flow.

Every @Builder.Dag class declares a wiring class, because the graph is what the Dag owns. A class that supplies only task bodies, for a Dag a Python file declares, carries @Builder.TaskHandler instead and contributes no Dag.

Note

A native Java Dag binds its task arguments from its own wiring, and the _arg_bindings it serializes are what Airflow records and shows. Runtime bindings (see Binding stub arguments) are what a @Builder.TaskHandler class reads, for a task whose Dag a Python file declares.

Task groups

A task group gathers tasks that the Airflow UI shows as one node, as Python’s TaskGroup does. Everything declared in a group carries the group’s ID as a prefix, so task stage in group staging is the task staging.stage. On the interface surface, taskGroup declares a group on the Dag or inside another group, and the group declares its tasks:

var staging = dag.taskGroup("staging");
var stage = staging.task("stage", Stage.class);              // "staging.stage"
staging.taskGroup("checks").task("nulls", Nulls.class).after(stage); // "staging.checks.nulls"
extract.before(staging);

With annotations, a @Builder.TaskGroup class holds the tasks of one group, and nesting one in another nests the groups:

@Builder.TaskGroup                      // the group "Staging", after the class
static class Staging {
  @Builder.Task
  public long stage(long rows) { ... }  // the task "Staging.stage"

  @Builder.TaskGroup(id = "checks")
  static class Checks {
    @Builder.Task
    public void nulls(long staged) { ... }   // "Staging.checks.nulls"
  }
}

@Builder.Deps
static class Wiring implements EtlPipelineDeps {
  void depends() {
    var rows = extract();
    load(transform(rows, lit(0.9)));
    rows.before(audit());
    var staged = staging().stage(rows);
    staging().checks().nulls(staged);
    extract().before(staging());
  }
}

The generated view nests the same way, so a group is both the namespace of what it holds and a point in the flow: staging().checks().nulls(staged) reaches a task, and extract().before(staging()) orders the whole group. Task method names scope to their own group, so two groups can each declare run(). A group class is static, non-private, and needs a no-argument constructor, because the generated task bodies instantiate it.

A group stands at either end of before, after and Flow.of. As an upstream it stands for its leaves, the tasks nothing else in the group runs after; as a downstream, for its roots, the tasks that run after nothing else in the group. A group ID contains only ASCII letters, digits, underscores, or dashes, and no task or other group in the Dag can share it.

Note

What a group holds is read once, when the Dag is serialized, which is what lets the wiring class above order a whole group before any of its tasks are declared, as extract().before(staging()) does. Python instead reads it at each >>. Edges are still resolved in the order they were drawn, as Python resolves them, so drawing an inner edge before or after an outer one gives different upstreams.

Configuration attributes

The @Builder.Dag and @Builder.Task configuration attributes, and the keys accepted by DagDef.config and TaskDef.config, are Airflow’s own Dag and task settings, under the names Airflow uses. Annotation attributes are camelCase (retryDelay); config keys are those names as Airflow writes them ("retry_delay"). Only attributes written explicitly at the use site are applied, so Airflow’s own defaults still apply to everything left out.

Durations and date-times are ISO-8601 strings in annotations (retryDelay = "PT5M", startDate = "2026-01-01T00:00:00Z", validated at compile time) and java.time.Duration / java.time.OffsetDateTime values in config calls. An unknown key or a mismatched value type fails the build for an annotation, and the config call itself for an object.

A Dag with a cron schedule runs in the time zone of its startDate. With no startDate it is scheduled in UTC, so set startDate to pin the zone.

Task state store

client.getTaskStateStore() gives a task key-value state that is scoped to the task instance and survives retry attempts within the same Dag run (see Task State Store). Use it to remember things like an external job ID so a retried task can resume instead of starting over:

@Builder.Task(id = "submit")
public void submit(Client client) throws Exception {
  var store = client.getTaskStateStore();
  var jobId = (String) store.get("job_id");
  if (jobId == null) {
    jobId = submitJob();
    store.set("job_id", jobId, Duration.ofHours(6));
  }
  waitForJob(jobId);
  store.delete("job_id");
}

get returns null when the key is not set. set stores any JSON-serializable value. Pass a positive java.time.Duration to expire the key after that long, TaskStateStore.NEVER_EXPIRE for a key that garbage collection skips, or omit the retention to use the deployment’s [state_store] default_retention_days (0 means never expire). The coordinator passes that setting to the JVM as AIRFLOW__STATE_STORE__DEFAULT_RETENTION_DAYS. A zero or negative retention is rejected. delete removes one key and clear removes every key for the task instance. The Java SDK does not use a [workers] state_store_backend: values always go to the metadata database as-is, so keys written by Python tasks through a custom backend are returned to Java as the raw reference marker rather than the stored value.

Parsing native Java Dags

To have Airflow parse the Dags a bundle JAR declares, put the JAR in a Dag bundle and configure a JavaCoordinator. The Dag processor runs the JAR’s main class to list its Dags, so it needs a Java executable, as the workers do:

[sdk]
coordinators = {
  "java-native": {
    "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
    "kwargs": {"java_executable": "/usr/lib/jvm/java-17-openjdk/bin/java"}
  }
}
queue_to_coordinator = {"java-native": "java-native"}

Once a JavaCoordinator is configured, the Dag processor parses the executable JARs of every Dag bundle, so it needs this [sdk] configuration and a JRE. With one JavaCoordinator, it parses them all. With several, map each Dag bundle that holds native Java Dags to one of them in [sdk] dag_bundle_to_coordinator. A JAR in a bundle that has no entry, or an entry that names no JavaCoordinator, fails to parse with an import error:

[sdk]
dag_bundle_to_coordinator = {"dags-folder": "java-native"}

A Dag bundle that holds only the JARs that Python Dags’ tasks run should list * in its .airflowignore either way: with one JavaCoordinator its JARs are still parsed on every loop, and a JAR built with a Java SDK older than schema 2026-10-30, every released one today, fails to parse as an import error. With several Java coordinators, an unmapped bundle’s JARs fail to parse too.

  • Every JAR in the bundle whose manifest sets Main-Class is parsed. Each Dag its main class declares, through Bundle.register of a DagDef or an @Builder.Dag class, is stored with that JAR as its file. Task handlers for a Python Dag are not Dags. A JAR without Main-Class is skipped, but many dependency JARs set one (the PostgreSQL JDBC driver and H2 do, for example), so a thin bundle should set main_class or list its dependency JARs in .airflowignore.

  • Parsing needs a JAR built with a Java SDK whose supervisor schema version (the Airflow-Supervisor-Schema-Version manifest attribute) is 2026-10-30 or later. An older JAR fails to parse, so list it in .airflowignore; its tasks still find it, because tasks do not read .airflowignore.

  • Do not declare a Dag in Java that a Python file in the same bundle also defines.

  • Keep one executable JAR per bundle, or set main_class, so that only JARs with that Main-Class are parsed. List JARs that should not be parsed in .airflowignore.

  • Two JARs in the bundle that set the same Main-Class fail to parse, because the JVM would load the classes of only one of them. Keep one in the bundle.

  • Set queue on every task, with @Builder.Task(queue = "java-native") or TaskDef.config("queue", "java-native"), so it runs on the coordinator’s queue. There is no Dag-level queue yet.

  • A task runs the JAR of its Dag on the JavaCoordinator that its queue routes to, which need not be the one that parsed the JAR. For example, a queue can route to a coordinator that uses another JDK. A task whose queue routes to another kind of coordinator fails without retries. A task whose JAR is missing, or that the coordinator cannot run (for example, because main_class does not match the JAR’s Main-Class), fails and retries while it has retries left.

  • The Code view shows the source of the JAR’s main class, which the Gradle plugin packs into the JAR.

  • Cluster policies (dag_policy, task_policy) are not applied to a native Java Dag.

  • airflow dags reserialize does not store the Dags of a JAR, which only the Dag processor stores. airflow dags test, tasks test and tasks render refuse a native Java Dag. airflow tasks list lists its tasks by running the JAR’s main class, so it needs a JRE.

Logging

Task code can emit log records through any common Java logging framework. The SDK ships optional integration libraries that forward those records to Airflow’s task log store, where they appear alongside the standard task output in the Airflow UI.

Declare a logger as a static field on the task class, using the class’s own type as the name. This is the conventional pattern regardless of which logging framework you choose:

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

@Builder.TaskHandler(dag = "etl", task = "extract")
public long extract(Client client) {
    log.log(System.Logger.Level.INFO, "Starting extraction");
    return recordCount;
}

The Gradle snippets below show the dependency declarations; all Airflow artifact versions are managed by airflow-sdk-bom. Maven users apply the same artifact IDs following the pattern in Maven.

System.Logger (Java Platform Logging)

Java 9’s new logging facade java.lang.System.Logger (JEP 264), commonly abbreviated JPL, can be used by libraries without pulling in any third-party API. The airflow-sdk-jpl artifact registers an AirflowSystemLoggerFinder via ServiceLoader, which routes all System.Logger calls directly to Airflow’s task log store.

implementation("org.apache.airflow:airflow-sdk-jpl:${version}")

No configuration file or startup call is required. The ServiceLoader mechanism discovers the provider automatically as long as the JAR is on the classpath.

Note

Do not add a second System.LoggerFinder implementation alongside airflow-sdk-jpl. The JVM selects one finder via ServiceLoader; having multiple providers on the classpath leads to unpredictable behaviour.

SLF4J 2.x

The SLF4J binding is discovered automatically via ServiceLoader; no configuration file or startup call is required.

implementation("org.apache.airflow:airflow-sdk-slf4j:${version}")

The above automatically pulls in the SLF4J API, so you don’t need to add slf4j-api yourself.

Note

Do not add a second SLF4J binding (such as logback-classic or slf4j-simple) alongside airflow-sdk-slf4j. SLF4J 2.x warns about multiple bindings and selects one unpredictably.

Log4j 2

airflow-sdk-log4j2 declares log4j-api as a transitive dependency, so you do not need to add the latter separately. You must also place log4j-core on the runtime classpath to host the plugin loader that discovers the custom AirflowAppender supplied by airflow-sdk-log4j2 at startup:

implementation("org.apache.airflow:airflow-sdk-log4j2:${version}")
runtimeOnly("org.apache.logging.log4j:log4j-core:${log4jVersion}")

Declare AirflowAppender in your log4j2.xml:

<?xml version="1.0" encoding="UTF-8"?>
<Configuration>
  <Appenders>
    <AirflowAppender name="Airflow"/>
  </Appenders>
  <Loggers>
    <Root level="info">
      <AppenderRef ref="Airflow"/>
    </Root>
  </Loggers>
</Configuration>

java.util.logging

Add the artifact:

implementation("org.apache.airflow:airflow-sdk-jul:${version}")

and call AirflowJulHandler.setup() on startup, before any task runs. It clears the JUL root logger’s existing handlers (including the default ConsoleHandler, whose stderr output Airflow would otherwise capture as task.stderr at ERROR level, duplicating each record and mislabeling its level) and installs AirflowJulHandler in their place:

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

Alternatively, declare the handler in a logging.properties file and point JUL at it with the java.util.logging.config.file system property (set via jvm_args in the coordinator configuration):

handlers = org.apache.airflow.sdk.jul.AirflowJulHandler
[sdk]
coordinators = {
  "java-jdk17": {
    "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
    "kwargs": {
      "task_handler_bundle_name": "java-task-handlers",
      "jvm_args": ["-Djava.util.logging.config.file=/opt/airflow/logging.properties"]
    }
  }
}

Other frameworks

Several commonly used logging APIs are covered without a dedicated Airflow artifact:

  • Logback is itself an SLF4J binding. Replace logback-classic with airflow-sdk-slf4j and no changes are needed in your task code.

  • Apache Commons Logging (JCL) can be bridged to SLF4J via org.slf4j:jcl-over-slf4j or to Log4j 2 via org.apache.logging.log4j:log4j-jcl.

XCom type mapping

XCom values are stored as JSON in Airflow’s metadata database. The table below shows how JSON types are represented as Java objects when read back via getXCom.

Python type

JSON

Java type (from getXCom)

int

number (integer)

Long (for values that fit; BigInteger otherwise)

float

number (decimal)

Double

str

string

String

bool

boolean

Boolean

None

null

null

list

array

List<Object>

dict

object

Map<String, Object>

Note

Avoid char and Character. JSON has no single-character type, so the value arrives as a string or a number and the SDK narrows it: a one-character string or an integer code point binds, an empty string binds as null, and anything longer fails at runtime with an IllegalArgumentException. Whether a value binds therefore depends on its length rather than on the stub signature, which no compile-time check can catch. Use String instead.

Note

A data parameter whose binding resolves to a value that was never pushed receives null. A boxed parameter (Integer, Long, Boolean, …) receives null safely, but a primitive parameter (int, long, boolean, …) cannot represent null and the task fails with MissingXComException. Declare the parameter with a boxed type when the upstream XCom may be absent.

Variables

Client reads, writes, and deletes Airflow Variables. Values are stored as strings. Serialize structured data (for example to JSON) before storing it.

var threshold = (String) client.getVariable("process_threshold");
client.setVariable("process_threshold", "42", "Rows above this count take the slow path");
client.deleteVariable("legacy_threshold");

Note

A value supplied by a secrets backend (for example an AIRFLOW_VAR_* environment variable) still takes precedence over the stored value when the Variable is read back. Calling setVariable without a description clears any existing description.

Building and packaging

The Java SDK is distributed as a JAR. The sections below show how to build a bundle with Gradle or Maven.

Gradle

Apply the Airflow SDK Gradle plugin in your build.gradle:

plugins {
    id("org.apache.airflow.sdk") version "${version}"
}

dependencies {
    annotationProcessor("org.apache.airflow:airflow-sdk-processor:${version}")
    implementation("org.apache.airflow:airflow-sdk:${version}")
}

airflowBundle {
    mainClass = "com.example.Main"  // Point to your main class instead.
}

Then run:

./gradlew bundle

The build/bundle/ directory contains all required JAR(s). Copy or mount it into the Dag bundle named by task_handler_bundle_name in the coordinator configuration. JavaCoordinator scans that Dag bundle recursively and builds the classpath automatically.

The plugin also packs the source file of mainClass, and of each class that declares a Dag in Java, into the bundle JAR, so the Airflow UI can show the source of a native Java Dag (see Parsing native Java Dags). To find those classes, the plugin runs mainClass once at build time. If that run fails, the build logs a warning and packs only the mainClass source.

A Dag is recorded against the class that constructed it. A Dag built by a factory therefore maps to the factory’s class, and one built by a factory in a dependency JAR has no source in the project and falls back to the mainClass source.

Note

You only need the annotationProcessor entry if you use the annotation-based API. It is not needed for the interface-based API.

Note

The plugin generates a fat JAR with the Shadow plugin by default. This is generally a good idea since you only deploy one JAR file to avoid dependency issues between projects. If this does not suit you, set fatJar = false in airflowBundle to produce thin JARs instead. The rest of the process stays the same, but you will need to put all dependency JARs in the same Dag bundle.

Maven

Import the airflow-sdk-bom Bill of Materials so that artifact versions and the ${airflow.supervisor.schema.version} property are managed in one place:

<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.apache.airflow</groupId>
            <artifactId>airflow-sdk-bom</artifactId>
            <version>${version}</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>

Add the SDK as a dependency (version is managed by the BOM):

<dependencies>
    <dependency>
        <groupId>org.apache.airflow</groupId>
        <artifactId>airflow-sdk</artifactId>
    </dependency>
</dependencies>

Wire the annotation processor through maven-compiler-plugin so it stays off the runtime classpath:

<plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-compiler-plugin</artifactId>
    <configuration>
        <annotationProcessorPaths>
            <path>
                <groupId>org.apache.airflow</groupId>
                <artifactId>airflow-sdk-processor</artifactId>
                <version>${version}</version>
            </path>
        </annotationProcessorPaths>
    </configuration>
</plugin>

Option 1 (recommended): fat JAR

Use maven-shade-plugin to bundle your code and all dependencies into a single JAR. This is the simplest deployment: one file, no dependency management at runtime.

<plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-shade-plugin</artifactId>
    <version>3.6.0</version>
    <executions>
        <execution>
            <phase>package</phase>
            <goals><goal>shade</goal></goals>
            <configuration>
                <transformers>
                    <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                        <!-- Replace with the class that provides main. -->
                        <mainClass>com.example.Main</mainClass>
                        <manifestEntries>
                            <!-- Resolved from the BOM; do not hard-code this value. -->
                            <Airflow-Supervisor-Schema-Version>${airflow.supervisor.schema.version}</Airflow-Supervisor-Schema-Version>
                        </manifestEntries>
                    </transformer>
                </transformers>
            </configuration>
        </execution>
    </executions>
</plugin>

Then run:

mvn package

The fat JAR is written to target/<artifactId>-<version>.jar. Copy it into the Dag bundle named by task_handler_bundle_name in your coordinator.

Option 2: thin JAR with separate dependencies

If a fat JAR does not suit your project, use maven-jar-plugin to set Main-Class on the regular JAR and maven-dependency-plugin to collect all runtime dependencies alongside it. Note that Airflow-Supervisor-Schema-Version does not need to be set here since Airflow reads it directly from the airflow-sdk JAR on the classpath.

<plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-jar-plugin</artifactId>
    <configuration>
        <archive>
            <manifestEntries>
                <Main-Class>com.example.Main</Main-Class>
            </manifestEntries>
        </archive>
    </configuration>
</plugin>

<plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-dependency-plugin</artifactId>
    <executions>
        <execution>
            <id>copy-dependencies</id>
            <phase>package</phase>
            <goals><goal>copy-dependencies</goal></goals>
            <configuration>
                <outputDirectory>${project.build.directory}/bundle</outputDirectory>
                <includeScope>runtime</includeScope>
            </configuration>
        </execution>
        <execution>
            <id>copy-artifact</id>
            <phase>package</phase>
            <goals><goal>copy</goal></goals>
            <configuration>
                <artifactItems>
                    <artifactItem>
                        <groupId>${project.groupId}</groupId>
                        <artifactId>${project.artifactId}</artifactId>
                        <version>${project.version}</version>
                        <outputDirectory>${project.build.directory}/bundle</outputDirectory>
                    </artifactItem>
                </artifactItems>
            </configuration>
        </execution>
    </executions>
</plugin>

Then run:

mvn package

target/bundle/ will contain the thin JAR and all runtime dependency JARs. Copy or mount this directory into the Dag bundle named by task_handler_bundle_name.

Note

You only need the annotationProcessorPaths entry if you use the annotation-based API.

Note

Unlike the Gradle plugin, Maven has no equivalent of the verifyBundleMainClass validation step. A wrong <mainClass> value will not be caught until runtime.

To show the source of a native Java Dag in the Airflow UI, pack the main class’s source file under META-INF/airflow/sources/ with an index that names it, src/main/resources/META-INF/airflow/sources.json:

{"entrypoint_path": "com/example/Main.java"}
<build>
    <resources>
        <!-- Declaring resources replaces the default, so keep it. -->
        <resource>
            <directory>src/main/resources</directory>
        </resource>
        <resource>
            <directory>src/main/java/com/example</directory>
            <includes><include>Main.java</include></includes>
            <targetPath>META-INF/airflow/sources/com/example</targetPath>
        </resource>
    </resources>
</build>

Then add <Airflow-Java-SDK-Sources>META-INF/airflow/sources.json</Airflow-Java-SDK-Sources> to the manifestEntries of maven-shade-plugin or maven-jar-plugin shown above.

JavaCoordinator configuration

All kwargs in the coordinators config entry are passed to the JavaCoordinator constructor:

Parameter

Default

Description

task_handler_bundle_name

(task’s own Dag bundle)

Name of the Dag bundle scanned recursively for .jar files. It is used only by mixed-language Dags, to locate the task handlers for the @task.stub tasks of a Python Dag; Dags defined natively in a language SDK do not use it. It must be registered in [dag_processor] dag_bundle_config_list. It is checked when the [sdk] configuration is loaded, so a typo fails there rather than on the first task.

java_executable

"java"

Path to the java binary. Defaults to java on $PATH.

jvm_args

[]

Extra JVM arguments such as ["-Xmx1g", "-Dsome.property=value"].

main_class

(auto-detect)

Explicit entry-point class. If omitted, the coordinator scans the Dag bundle for a JAR whose manifest sets Main-Class. If more than one JAR in that Dag bundle sets it, the first by path is used, so set main_class explicitly in that case. A task of a native Java Dag runs the JAR the Dag was parsed from. When the coordinator parses native Java Dags, only JARs with this Main-Class are parsed.

task_startup_timeout

10.0

Seconds to wait for the JVM subprocess to connect after launch. Increase this if your JVM startup is slow (e.g. on constrained hardware or with a large classpath).

Note

Locating JARs. The JARs for the @task.stub tasks of a Python Dag are read from a Dag bundle, so they are delivered, refreshed and versioned by the same machinery as your Dags.

  • The expected layout is a separate Dag bundle for the JARs, named by task_handler_bundle_name, rather than the Dag bundle that holds your .py files. The task uses the version that Dag bundle is on when it starts, pinned for the whole task.

  • If task_handler_bundle_name is unset, JARs are read from the task’s own Dag bundle, pinned to the version the run was created with.

  • Every JAR in the Dag bundle goes on one classpath, so all handlers in it share one set of dependencies. To isolate conflicting dependency versions, put the handlers in a second Dag bundle served by a second coordinator on its own queue.

A task of a native Java Dag ignores task_handler_bundle_name: it runs the JAR of its Dag from the Dag’s own bundle, at the version the run was created with. See Parsing native Java Dags.

Note

The [sdk] configuration is read at startup, so changes to coordinators or queue_to_coordinator (for example adding jvm_args) only take effect after you restart the components that read it: the workers (the scheduler with LocalExecutor), the Dag processor, or airflow standalone. A rebuilt bundle JAR, by contrast, is picked up on the next task launch without a restart, because a fresh JVM is spawned per task instance.

Pinning the Java executable

As a general recommendation, set java_executable to an absolute path rather than relying on java resolving from $PATH. This pins tasks to a known JDK, which matters most in production or corporate environments where the Airflow admin may not control the system-wide java (the same reasoning behind pinning a Python version).

For example, if you install the JDK with Homebrew on macOS, its java is not on $PATH, so point java_executable at it explicitly:

[sdk]
coordinators = {
  "java-jdk17": {
    "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
    "kwargs": {
      "task_handler_bundle_name": "java-task-handlers",
      "java_executable": "/opt/homebrew/opt/openjdk@17/bin/java"
    }
  }
}
queue_to_coordinator = {"java": "java-jdk17"}

Limitations

  • One JVM subprocess per task instance. Each task instance spawns a fresh JVM. Tasks that need to share in-process state between instances should use XCom or an external store instead.

  • Limited support for assets, deferral, and other Airflow features. They may be implemented in the future based on user feedback and demand.

Was this entry helpful?