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-sdkpackage (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 |
|---|---|
|
Marks a method as the Java body of a task the Python Dag file declares with |
|
Marks a class as a Dag that Java itself owns. Attributes ( |
|
Marks a method as a task of a Java-owned Dag. If |
|
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. |
|
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. |
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
staticnested 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-Classis parsed. Each Dag its main class declares, throughBundle.registerof aDagDefor an@Builder.Dagclass, is stored with that JAR as its file. Task handlers for a Python Dag are not Dags. A JAR withoutMain-Classis skipped, but many dependency JARs set one (the PostgreSQL JDBC driver and H2 do, for example), so a thin bundle should setmain_classor list its dependency JARs in.airflowignore.Parsing needs a JAR built with a Java SDK whose supervisor schema version (the
Airflow-Supervisor-Schema-Versionmanifest attribute) is2026-10-30or 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 thatMain-Classare parsed. List JARs that should not be parsed in.airflowignore.Two JARs in the bundle that set the same
Main-Classfail to parse, because the JVM would load the classes of only one of them. Keep one in the bundle.Set
queueon every task, with@Builder.Task(queue = "java-native")orTaskDef.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
JavaCoordinatorthat 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, becausemain_classdoes not match the JAR’sMain-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 reserializedoes not store the Dags of a JAR, which only the Dag processor stores.airflow dags test,tasks testandtasks renderrefuse a native Java Dag.airflow tasks listlists 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-classicwithairflow-sdk-slf4jand no changes are needed in your task code.Apache Commons Logging (JCL) can be bridged to SLF4J via
org.slf4j:jcl-over-slf4jor to Log4j 2 viaorg.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 |
|---|---|---|
|
number (integer) |
|
|
number (decimal) |
|
|
string |
|
|
boolean |
|
|
null |
|
|
array |
|
|
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’s own Dag bundle) |
Name of the Dag bundle scanned recursively for |
|
|
Path to the |
|
|
Extra JVM arguments such as |
|
(auto-detect) |
Explicit entry-point class. If omitted, the coordinator scans the Dag bundle for a JAR
whose manifest sets |
|
|
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.pyfiles. The task uses the version that Dag bundle is on when it starts, pinned for the whole task.If
task_handler_bundle_nameis 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.