This is an automated email from the ASF dual-hosted git repository.
henry3260 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 21bbf2f40a2 Java SDK: Generate Dag and task configuration from the
serialization schema (#73595)
21bbf2f40a2 is described below
commit 21bbf2f40a2e4799d82898ee2d0fff350447588a
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Wed Sep 30 15:00:35 2026 +0800
Java SDK: Generate Dag and task configuration from the serialization schema
(#73595)
* Java SDK: Describe native Java Dags without reference to Python
A Dag authored entirely in Java has no Python in the picture, so teaching
its edge verbs as "Java's spelling of >> and <<" asks the reader to know a
language the feature does not involve. The mixed-language surface keeps its
Python references, because there the Python Dag file is part of what the
reader has to understand.
* Java SDK: Generate Dag and task configuration from the serialization
schema
A Java-authored Dag could carry no schedule, retries, tags or any other
Airflow setting, so it could be defined but not configured. Hand-listing
those settings in Java would drift from the Python semantics they mirror,
so they are recorded and validated against Airflow's own Dag serialization
schema instead, and take effect once the runtime answers parse requests.
The `Builder` annotations are generated from a copy of that schema vendored
at java-sdk/sdk/schema/dag-schema.json, so a source release builds without
the monorepo; a prek hook keeps the copy in sync.
---
.pre-commit-config.yaml | 11 +
.../language-sdks/java.rst | 37 +-
java-sdk/README.md | 19 +-
.../example/nativedag/InterfaceExample.java | 18 +-
.../org/apache/airflow/sdk/BuilderProcessor.kt | 118 ++++-
.../kotlin/org/apache/airflow/sdk/BuilderTest.kt | 172 ++++++--
java-sdk/sdk/build.gradle.kts | 443 ++++++++++++++++++-
java-sdk/sdk/schema/dag-schema.json | 474 +++++++++++++++++++++
.../src/main/kotlin/org/apache/airflow/sdk/Arg.kt | 18 +
.../main/kotlin/org/apache/airflow/sdk/Builder.kt | 109 -----
.../main/kotlin/org/apache/airflow/sdk/DagDef.kt | 69 ++-
.../src/main/kotlin/org/apache/airflow/sdk/Deps.kt | 18 +-
.../org/apache/airflow/sdk/internal/Fields.kt | 107 +++++
.../kotlin/org/apache/airflow/sdk/DagDefTest.kt | 179 ++++++++
14 files changed, 1611 insertions(+), 181 deletions(-)
diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml
index a454338a410..7fc98f8f3ab 100644
--- a/.pre-commit-config.yaml
+++ b/.pre-commit-config.yaml
@@ -337,6 +337,17 @@ repos:
(?x)
^airflow-core/src/airflow/serialization/schema\.json$|
^ts-sdk/schema/dag-schema\.json$
+ - id: sync-java-sdk-dag-schema
+ name: Sync Java SDK Dag serialization schema with airflow-core
+ description: "Copy airflow-core's serialization schema when Java SDK's
vendored dag-schema.json drifts"
+ entry: ./java-sdk/gradlew -p ./java-sdk :sdk:syncDagSchema
+ language: system
+ pass_filenames: false
+ # Re-vendoring is a java-sdk maintainer's deliberate step after an
airflow-core
+ # schema change, not something every commit should do, so this is
manual only:
+ # prek run sync-java-sdk-dag-schema --hook-stage manual
+ stages: ['manual']
+ files: ^java-sdk/sdk/schema/dag-schema\.json$
- id: check-go-version-in-sync
name: Check Go toolchain version is consistent across build files
entry: ./scripts/ci/prek/check_go_version_in_sync.py
diff --git a/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
b/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
index 1561ecfffcf..dbd5d8996ce 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
@@ -306,6 +306,14 @@ Annotate a plain Java class and let the SDK generate the
boilerplate at compile
``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 :ref:`java-sdk/native-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.
* - ``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
@@ -348,7 +356,8 @@ 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``.
+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:
@@ -537,16 +546,15 @@ calls with no arguments.
Native Java Dags
----------------
-A Dag can also be authored entirely in Java, with no Python stub file: the
``DagDef`` and
-``TaskDef`` objects hold the tasks, and Java declares the graph.
+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,
so there is no second
-``addTask`` call to forget. ``before`` and ``after`` draw every edge on this
surface — Python's
-``a >> b`` and ``b << a`` — and the task body moves the data itself, by
reading the upstream's XCom
-through ``Client``:
+``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``:
.. code-block:: java
@@ -566,6 +574,21 @@ Edges are checked when the Dag is registered with a
``Bundle``: an upstream that
Dag, or to no Dag, and a cycle anywhere in the graph both fail there rather
than at the first task
run.
+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.
+
.. _java-sdk/logging:
Logging
diff --git a/java-sdk/README.md b/java-sdk/README.md
index 12a43cbf3c5..b8600ab5d44 100644
--- a/java-sdk/README.md
+++ b/java-sdk/README.md
@@ -726,9 +726,13 @@ E2E_TEST_MODE=java_sdk uv run --project airflow-e2e-tests
pytest \
not the implementation language.
- Keep `sdk/src/main/kotlin/` (the public API surface) free of internal
implementation details; those belong in the `execution/` sub-package.
-- The annotation processor (`BuilderProcessor.kt`) uses `kapt`. When adding a
- new annotation, define it in `Builder.kt`, handle it in
- `BuilderProcessor.kt`, and add a golden-output test in
+- The annotation processor (`BuilderProcessor.kt`) uses `kapt`. The `Builder`
+ class holding the `@Builder.Dag` / `@Builder.Task` annotations is generated
+ from the Dag serialization schema by `:sdk:generateDagDsl` (vendored at
+ `sdk/schema/dag-schema.json`). The `Arg`/`TaskRef` and `Deps`/`Flow` graph
+ types are hand-written next to the rest of the public surface in
+ `sdk/src/main/kotlin/org/apache/airflow/sdk/`. When adding annotation
+ behaviour, handle it in `BuilderProcessor.kt` and add a golden-output test in
`processor/src/test/kotlin/`.
- The Python coordinator subclasses `SubprocessCoordinator`. Do not reach into
the JVM process from Python beyond what `_build_execute_task_command`
@@ -750,9 +754,14 @@ E2E_TEST_MODE=java_sdk uv run --project airflow-e2e-tests
pytest \
5. Update `airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst`
if the change is user-visible.
-**Adding a new annotation**:
+**Adding a new annotation or configuration attribute**:
-1. Define the annotation interface in `Builder.kt`.
+1. Hand-written annotations live in
`sdk/src/main/kotlin/org/apache/airflow/sdk/`
+ next to the runtime types (one package, so user code needs a single
+ `import org.apache.airflow.sdk.*`); the configuration attributes of
+ `@Builder.Dag` / `@Builder.Task` come from the Dag serialization schema via
+ `:sdk:generateDagDsl` (adjust its allowlist/exclusion rules in
+ `sdk/build.gradle.kts` when the exposed field set should change).
2. Handle it in `BuilderProcessor.kt` — generate the appropriate code in the
`*Builder` class.
3. Add a test in `BuilderTest.kt` with expected generated output.
diff --git
a/java-sdk/example/src/java/org/apache/airflow/example/nativedag/InterfaceExample.java
b/java-sdk/example/src/java/org/apache/airflow/example/nativedag/InterfaceExample.java
index 58c33218254..6047aeedd8e 100644
---
a/java-sdk/example/src/java/org/apache/airflow/example/nativedag/InterfaceExample.java
+++
b/java-sdk/example/src/java/org/apache/airflow/example/nativedag/InterfaceExample.java
@@ -22,11 +22,11 @@ package org.apache.airflow.example.nativedag;
import static java.lang.System.Logger.Level.INFO;
+import java.util.List;
import org.apache.airflow.sdk.*;
-// A Dag defined entirely in Java, interface-style: no Python stub file
-// describes it. dag.task registers a task as it creates it and hands back the
-// handle, and `before`/`after` wire the graph -- Java's spelling of `>>` and
`<<`.
+// A Dag defined entirely in Java, interface-style. dag.task registers a task
as
+// it creates it and hands back the handle, and `before`/`after` wire the
graph.
public class InterfaceExample {
private static final System.Logger log =
System.getLogger(InterfaceExample.class.getName());
@@ -56,9 +56,17 @@ public class InterfaceExample {
}
public static DagDef build() {
- var dag = new DagDef("java_native_interface_example");
+ var dag =
+ new DagDef("java_native_interface_example")
+ .config("description", "Pure-Java Dag authored with the interface
API")
+ .config("schedule", "@daily")
+ .config("catchup", false)
+ .config("tags", List.of("example", "java-sdk"));
- var extract = dag.task("extract", Extract.class);
+ var extract =
+ dag.task("extract", Extract.class)
+ .config("retries", 2)
+ .config("doc_md", "Extracts a value and pushes it as an XCom.");
var transform = dag.task("transform", Transform.class);
var load = dag.task("load", Load.class);
diff --git
a/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
b/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
index dfee528a30d..049aba0d060 100644
---
a/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
+++
b/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
@@ -29,16 +29,24 @@ import com.squareup.javapoet.ParameterizedTypeName
import com.squareup.javapoet.TypeName
import com.squareup.javapoet.TypeSpec
import org.apache.airflow.sdk.internal.ArgValues
+import org.apache.airflow.sdk.internal.Field
+import org.apache.airflow.sdk.internal.FieldType
+import org.apache.airflow.sdk.internal.SchemaFields
import org.apache.airflow.sdk.internal.TaskArgs
import org.apache.airflow.sdk.internal.TypeRef
import org.apache.airflow.sdk.internal.foldArgName
import org.apache.airflow.sdk.internal.registrarName
+import java.time.Duration
+import java.time.OffsetDateTime
+import java.time.format.DateTimeParseException
import javax.annotation.processing.AbstractProcessor
import javax.annotation.processing.ProcessingEnvironment
import javax.annotation.processing.RoundEnvironment
import javax.annotation.processing.SupportedAnnotationTypes
import javax.annotation.processing.SupportedSourceVersion
import javax.lang.model.SourceVersion
+import javax.lang.model.element.AnnotationValue
+import javax.lang.model.element.Element
import javax.lang.model.element.ElementKind
import javax.lang.model.element.ExecutableElement
import javax.lang.model.element.Modifier
@@ -61,13 +69,15 @@ import javax.tools.Diagnostic
* containing:
*
* - One inner class per [Builder.Task]-annotated method, implementing [Task].
- * - A static `build()` method that constructs the [DagDef] and registers those
- * inner classes as [TaskDef]s.
+ * - A static `build()` method that constructs the [DagDef], lowers every
+ * explicitly-written `@Builder.Dag` attribute into a `DagDef.config` call,
+ * and registers those inner classes as [TaskDef]s, each carrying its
+ * explicitly-written `@Builder.Task` attributes the same way.
*
* In the generated `execute` body, a task's data parameters resolve against
the
* arg bindings the supervisor delivered for the run: flat parameters through
* [TaskArgs], by their position among the data parameters, and [TaskInput]
- * [TaskInput] fields through [ArgValues], by argument name. Non-`void` return
values are
+ * fields through [ArgValues], by argument name. Non-`void` return values are
* forwarded to `client.setXCom`.
*/
@SupportedAnnotationTypes(
@@ -177,6 +187,9 @@ class BuilderProcessor : AbstractProcessor() {
.addModifiers(Modifier.PUBLIC, Modifier.STATIC)
.returns(DAG_DEF_TYPE)
.addStatement($$"var dag = new $T($S)", DAG_DEF_TYPE, ann.id.ifBlank {
el.simpleName })
+ explicitConfig(el, DAG_ANNOTATION, DAG_STRUCTURAL_ATTRIBUTES,
SchemaFields.DAG).forEach { (key, value) ->
+ buildMethod.addStatement($$"dag.config($S, $L)", key, value)
+ }
for (inner in el.enclosedElements) {
if (inner !is ExecutableElement) continue
@@ -188,10 +201,8 @@ class BuilderProcessor : AbstractProcessor() {
builderClass.addType(buildTask(innerName, inner, el))
buildMethod.addStatement(
- $$"dag.addTask(new $T($S, $L.class))",
- TASK_DEF_TYPE,
- taskAnn.id.ifBlank { inner.simpleName },
- innerName,
+ $$"dag.addTask($L)",
+ taskDefCode(inner, taskAnn.id.ifBlank { inner.simpleName.toString() },
innerName),
)
}
@@ -200,6 +211,93 @@ class BuilderProcessor : AbstractProcessor() {
return builderClass.build()
}
+ /**
+ * Emits `new TaskDef(id, <className>.class)` with the explicitly-written
+ * `@Builder.Task` attributes lowered into chained `.config` calls.
+ */
+ private fun taskDefCode(
+ method: ExecutableElement,
+ id: String,
+ className: String,
+ ): CodeBlock {
+ val taskDef =
+ CodeBlock
+ .builder()
+ .add($$"new $T($S, $L.class)", TASK_DEF_TYPE, id, className)
+ explicitConfig(method, TASK_ANNOTATION, TASK_STRUCTURAL_ATTRIBUTES,
SchemaFields.TASK).forEach { (key, value) ->
+ taskDef.add($$".config($S, $L)", key, value)
+ }
+ return taskDef.build()
+ }
+
+ /**
+ * Lowers the explicitly-written configuration attributes of [element]'s
+ * [annotationName] annotation into (schema key, value code) pairs. Only
+ * attributes present at the use site are lowered, so annotation defaults
+ * never override the schema's own defaults.
+ */
+ private fun explicitConfig(
+ element: Element,
+ annotationName: String,
+ structural: Set<String>,
+ table: Map<String, Field>,
+ ): List<Pair<String, CodeBlock>> {
+ val mirror =
+ element.annotationMirrors.firstOrNull {
+ (it.annotationType.asElement() as
TypeElement).qualifiedName.contentEquals(annotationName)
+ } ?: return emptyList()
+ val byAttribute = table.values.associateBy { it.attribute }
+ return mirror.elementValues.mapNotNull { (attr, value) ->
+ val name = attr.simpleName.toString()
+ if (name in structural) return@mapNotNull null
+ val field =
+ requireNotNull(byAttribute[name]) {
+ "Annotation attribute '$name' has no Dag serialization schema key"
+ }
+ field.key to configValueCode(field, value)
+ }
+ }
+
+ private fun configValueCode(
+ field: Field,
+ value: AnnotationValue,
+ ): CodeBlock =
+ when (field.type) {
+ FieldType.STRING -> CodeBlock.of($$"$S", value.value)
+ FieldType.BOOLEAN, FieldType.INTEGER, FieldType.NUMBER ->
CodeBlock.of($$"$L", value.value)
+ FieldType.STRING_ARRAY -> {
+ @Suppress("UNCHECKED_CAST")
+ val items = value.value as List<AnnotationValue>
+ CodeBlock.of(
+ $$"$T.of($L)",
+ ClassName.get(List::class.java),
+ CodeBlock.join(items.map { CodeBlock.of($$"$S", it.value) }, ", "),
+ )
+ }
+ FieldType.TIMEDELTA -> {
+ val text = value.value as String
+ parseTemporal(field, text) { Duration.parse(text) }
+ CodeBlock.of($$"$T.parse($S)", ClassName.get(Duration::class.java),
text)
+ }
+ FieldType.DATETIME -> {
+ val text = value.value as String
+ parseTemporal(field, text) { OffsetDateTime.parse(text) }
+ CodeBlock.of($$"$T.parse($S)",
ClassName.get(OffsetDateTime::class.java), text)
+ }
+ }
+
+ private fun parseTemporal(
+ field: Field,
+ text: String,
+ parse: () -> Any,
+ ) {
+ try {
+ parse()
+ } catch (e: DateTimeParseException) {
+ throw IllegalArgumentException("Annotation attribute
'${field.attribute}' is not valid ISO-8601: '$text'")
+ }
+ }
+
private fun buildTask(
name: String,
inner: ExecutableElement,
@@ -404,6 +502,12 @@ private val TASK_ARGS_TYPE =
ClassName.get(TaskArgs::class.java)
private val TYPE_REF_TYPE = ClassName.get(TypeRef::class.java)
private val ARG_VALUES_TYPE = ClassName.get(ArgValues::class.java)
+private const val DAG_ANNOTATION = "org.apache.airflow.sdk.Builder.Dag"
+private const val TASK_ANNOTATION = "org.apache.airflow.sdk.Builder.Task"
+
+private val DAG_STRUCTURAL_ATTRIBUTES = setOf("id", "to")
+private val TASK_STRUCTURAL_ATTRIBUTES = setOf("id")
+
private fun ProcessingEnvironment.isType(
t: TypeMirror,
c: ClassName,
diff --git
a/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
b/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
index 18e4d1501a1..6fd52e01b75 100644
--- a/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
+++ b/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
@@ -58,7 +58,7 @@ class BuilderTest {
@Builder.Task
public int t2(Client client) {
- return (Integer) client.getXCom("t0");
+ return 7;
}
@Builder.Task
@@ -95,18 +95,21 @@ class BuilderTest {
dag.addTask(new TaskDef("t3", T3.class));
return dag;
}
+
public static final class T1 implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
new TestExample().t1();
}
}
+
public static final class T2 implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
client.setXCom(new TestExample().t2(client));
}
}
+
public static final class T3 implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
@@ -164,6 +167,7 @@ class BuilderTest {
dag.addTask(new TaskDef("t", T.class));
return dag;
}
+
public static final class T implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
@@ -227,6 +231,7 @@ class BuilderTest {
dag.addTask(new TaskDef("t", T.class));
return dag;
}
+
public static final class T implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
@@ -245,27 +250,18 @@ class BuilderTest {
}
@Test
- @DisplayName("bind a TaskInput through the shared populator")
- fun generateBuilderBindsTaskInputFields() {
+ @DisplayName("lower explicit annotation attributes into config calls")
+ fun generateBuilderLowersConfigAttributes() {
val compilation =
compile(
"""
package org.apache.airflow.example;
- import java.util.List;
- import org.apache.airflow.sdk.ArgName;
import org.apache.airflow.sdk.Builder;
- import org.apache.airflow.sdk.Client;
- import org.apache.airflow.sdk.TaskInput;
- @Builder.Dag
+ @Builder.Dag(id = "cfg", schedule = "@daily", tags = {"a", "b"},
catchup = true,
+ startDate = "2026-01-01T00:00:00Z")
public class TestExample {
- public static class ScoreInput implements TaskInput {
- @ArgName("region_code") public String region;
- public double threshold;
- public List<String> tags;
- }
-
- @Builder.Task
- public double score(Client client, ScoreInput input) { return
input.threshold; }
+ @Builder.Task(retries = 2, queue = "q", retryDelay = "PT5M",
retryExponentialBackoff = 1.5)
+ public void t1() {}
}
""",
)
@@ -280,24 +276,30 @@ class BuilderTest {
import java.lang.Exception;
import java.lang.Override;
+ import java.time.Duration;
+ import java.time.OffsetDateTime;
+ import java.util.List;
import org.apache.airflow.sdk.Client;
import org.apache.airflow.sdk.Context;
import org.apache.airflow.sdk.DagDef;
import org.apache.airflow.sdk.Task;
import org.apache.airflow.sdk.TaskDef;
- import org.apache.airflow.sdk.internal.ArgValues;
public final class TestExampleBuilder {
public static DagDef build() {
- var dag = new DagDef("TestExample");
- dag.addTask(new TaskDef("score", Score.class));
+ var dag = new DagDef("cfg");
+ dag.config("schedule", "@daily");
+ dag.config("tags", List.of("a", "b"));
+ dag.config("catchup", true);
+ dag.config("start_date",
OffsetDateTime.parse("2026-01-01T00:00:00Z"));
+ dag.addTask(new TaskDef("t1", T1.class).config("retries",
2).config("queue", "q").config("retry_delay",
Duration.parse("PT5M")).config("retry_exponential_backoff", 1.5));
return dag;
}
- public static final class Score implements Task {
+
+ public static final class T1 implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
- TestExample.ScoreInput input = ArgValues.bindInput(client,
TestExample.ScoreInput.class);
- client.setXCom(new TestExample().score(client, input));
+ new TestExample().t1();
}
}
}
@@ -411,6 +413,7 @@ class BuilderTest {
dag.addTask(new TaskDef("named", Named.class));
return dag;
}
+
public static final class Flat implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
@@ -420,6 +423,7 @@ class BuilderTest {
new TestExample().flat(client_, context_);
}
}
+
public static final class Named implements Task {
@Override
public void execute(Context context, Client client) throws
Exception {
@@ -617,7 +621,9 @@ class BuilderTest {
"""
package org.apache.airflow.example;
import org.apache.airflow.sdk.DagDef;
- public final class TestExampleBuilder { public static DagDef build()
{ var dag = new DagDef("foo"); return dag; } }
+ public final class TestExampleBuilder {
+ public static DagDef build() { var dag = new DagDef("foo"); return
dag; }
+ }
""",
)
}
@@ -640,7 +646,9 @@ class BuilderTest {
"""
package org.apache.airflow.example;
import org.apache.airflow.sdk.DagDef;
- public final class Foo { public static DagDef build() { var dag = new
DagDef("TestExample"); return dag; } }
+ public final class Foo {
+ public static DagDef build() { var dag = new DagDef("TestExample");
return dag; }
+ }
""",
)
}
@@ -654,7 +662,9 @@ class BuilderTest {
package org.apache.airflow.example;
import org.apache.airflow.sdk.Builder;
@Builder.Dag
- public class TestExample { @Builder.Task(id = "foo") public void t1()
{} }
+ public class TestExample {
+ @Builder.Task(id = "foo") public void t1() {}
+ }
""",
)
@@ -664,6 +674,7 @@ class BuilderTest {
"org.apache.airflow.example.TestExampleBuilder",
"""
package org.apache.airflow.example;
+
import java.lang.Exception;
import java.lang.Override;
import org.apache.airflow.sdk.Client;
@@ -671,14 +682,19 @@ class BuilderTest {
import org.apache.airflow.sdk.DagDef;
import org.apache.airflow.sdk.Task;
import org.apache.airflow.sdk.TaskDef;
+
public final class TestExampleBuilder {
public static DagDef build() {
var dag = new DagDef("TestExample");
dag.addTask(new TaskDef("foo", T1.class));
return dag;
}
+
public static final class T1 implements Task {
- @Override public void execute(Context context, Client client)
throws Exception { new TestExample().t1(); }
+ @Override
+ public void execute(Context context, Client client) throws
Exception {
+ new TestExample().t1();
+ }
}
}
""",
@@ -703,6 +719,110 @@ class BuilderTest {
)
}
+ @Test
+ @DisplayName("reject a duration attribute that is not ISO-8601")
+ fun rejectInvalidDurationAttribute() {
+ val compilation =
+ compile(
+ """
+ package org.apache.airflow.example;
+ import org.apache.airflow.sdk.Builder;
+ @Builder.Dag
+ public class TestExample {
+ @Builder.Task(retryDelay = "5 minutes") public void t() {}
+ }
+ """,
+ )
+ assertThat(compilation).failed()
+ assertThat(compilation).hadErrorContaining(
+ "Annotation attribute 'retryDelay' is not valid ISO-8601: '5 minutes'",
+ )
+ }
+
+ @Test
+ @DisplayName("escape quotes and backslashes in a string-array attribute")
+ fun generateBuilderEscapesStringArrayValues() {
+ val compilation =
+ compile(
+ """
+ package org.apache.airflow.example;
+ import org.apache.airflow.sdk.Builder;
+ @Builder.Dag(tags = {"say \"hi\"", "back\\slash"})
+ public class TestExample {
+ @Builder.Task public void t() {}
+ }
+ """,
+ )
+
+ assertThat(compilation).succeeded()
+ assertThat(compilation)
+ .generatedSourceFile("org.apache.airflow.example.TestExampleBuilder")
+ .contentsAsUtf8String()
+ .contains("""dag.config("tags", List.of("say \"hi\"",
"back\\slash"));""")
+ }
+
+ @Test
+ @DisplayName("bind a TaskInput through the shared populator")
+ fun generateBuilderBindsTaskInputFields() {
+ val compilation =
+ compile(
+ """
+ package org.apache.airflow.example;
+ import java.util.List;
+ import org.apache.airflow.sdk.ArgName;
+ import org.apache.airflow.sdk.Builder;
+ import org.apache.airflow.sdk.Client;
+ import org.apache.airflow.sdk.TaskInput;
+ @Builder.Dag
+ public class TestExample {
+ public static class ScoreInput implements TaskInput {
+ @ArgName("region_code") public String region;
+ public double threshold;
+ public List<String> tags;
+ }
+
+ @Builder.Task
+ public double score(Client client, ScoreInput input) { return
input.threshold; }
+ }
+ """,
+ )
+
+ assertThat(compilation).succeeded()
+ assertThat(compilation)
+ .generatedSourceFile("org.apache.airflow.example.TestExampleBuilder")
+ .hasSourceEquivalentTo(
+ "org.apache.airflow.example.TestExampleBuilder",
+ """
+ package org.apache.airflow.example;
+
+ import java.lang.Exception;
+ import java.lang.Override;
+ import org.apache.airflow.sdk.Client;
+ import org.apache.airflow.sdk.Context;
+ import org.apache.airflow.sdk.DagDef;
+ import org.apache.airflow.sdk.Task;
+ import org.apache.airflow.sdk.TaskDef;
+ import org.apache.airflow.sdk.internal.ArgValues;
+
+ public final class TestExampleBuilder {
+ public static DagDef build() {
+ var dag = new DagDef("TestExample");
+ dag.addTask(new TaskDef("score", Score.class));
+ return dag;
+ }
+
+ public static final class Score implements Task {
+ @Override
+ public void execute(Context context, Client client) throws
Exception {
+ TestExample.ScoreInput input = ArgValues.bindInput(client,
TestExample.ScoreInput.class);
+ client.setXCom(new TestExample().score(client, input));
+ }
+ }
+ }
+ """,
+ )
+ }
+
@Test
@DisplayName("generate a registrar binding each handler to the ids its
annotation names")
fun generateHandlerRegistrar() {
diff --git a/java-sdk/sdk/build.gradle.kts b/java-sdk/sdk/build.gradle.kts
index fb74fcf9c9d..5e9a74a7046 100644
--- a/java-sdk/sdk/build.gradle.kts
+++ b/java-sdk/sdk/build.gradle.kts
@@ -21,6 +21,7 @@ import java.io.File
import java.net.URI
import java.nio.file.Files
import java.nio.file.StandardCopyOption
+import java.time.Duration
val airflowSupervisorSchemaVersion: String by project
@@ -43,6 +44,8 @@ val pointersDir =
layout.buildDirectory.dir("schema-pointers/main")
val jsonSchemaPackage = "org.apache.airflow.sdk.execution.comm"
val schemaModelsDir =
layout.buildDirectory.dir("generate-resources/main/src/main/java")
val discriminatorDir =
layout.buildDirectory.dir("generated-resources/main/src/main/kotlin")
+val dagSchemaInput = layout.projectDirectory.file("schema/dag-schema.json")
+val dagDslDir =
layout.buildDirectory.dir("generated-resources/dsl/src/main/kotlin")
dependencies {
compileOnly("com.github.spotbugs:spotbugs-annotations:4.10.4")
@@ -212,6 +215,426 @@ abstract class SyncSupervisorSchemaTask : DefaultTask() {
}
}
+// Keep the vendored Dag serialization schema in sync with the monorepo copy.
+// The vendored file makes standalone (source-release) builds work; in-repo
+// builds refresh it from airflow-core, and a prek hook guards against drift.
+abstract class SyncDagSchemaTask : DefaultTask() {
+ @get:Internal
+ abstract val sourceFile: RegularFileProperty
+
+ @get:Internal
+ abstract val targetFile: RegularFileProperty
+
+ @TaskAction
+ fun sync() {
+ val src = sourceFile.get().asFile
+ if (!src.exists()) {
+ logger.lifecycle("Monorepo serialization schema not present;
keeping vendored dag-schema.json.")
+ return
+ }
+ val dst = targetFile.get().asFile
+ if (dst.exists() && dst.readText() == src.readText()) {
+ logger.lifecycle("Vendored dag-schema.json is up-to-date.")
+ return
+ }
+ logger.lifecycle("Refreshing vendored dag-schema.json from
${src.path}")
+ src.copyTo(dst, overwrite = true)
+ }
+}
+
+// Generate the Dag-authoring DSL surface from the Dag serialization schema:
+//
+// - org.apache.airflow.sdk.Builder and its nested Dag / Task annotations,
+// whose configuration attributes mirror the scalar keys of the schema's
+// "dag" and "operator" definitions (the annotation processor lowers
+// explicitly-set attributes into DagDef.config / TaskDef.config calls), and
+// - org.apache.airflow.sdk.internal.SchemaFields, the key -> type table that
+// DagDef.config / TaskDef.config validate against at registration time.
+//
+// Field selection mirrors the Go SDK's TaskSpec generator: scalar properties
+// only (string/integer/number/boolean plus timedelta/datetime refs),
+// serializer-owned keys skipped ("_"-prefixed, schema-required, "has_on_"
+// callbacks), and a documented exclusion list for Python-only concerns. An
+// exclusion entry that stops matching an eligible key fails generation, so
+// the list cannot go stale.
+abstract class GenerateDagDslTask : DefaultTask() {
+ @get:InputFile
+ abstract val schemaFile: RegularFileProperty
+
+ @get:OutputDirectory
+ abstract val targetDirectory: DirectoryProperty
+
+ private data class DslField(
+ val key: String,
+ val attribute: String,
+ val fieldType: String,
+ val attrType: String,
+ val attrDefault: String,
+ val defaultJson: String?,
+ val doc: String,
+ )
+
+ private fun camelCase(key: String): String =
+ key
+ .split('_')
+ .filter { it.isNotEmpty() }
+ .mapIndexed { i, seg ->
+ if (i == 0) seg else seg.replaceFirstChar(Char::uppercase)
+ }.joinToString("")
+
+ private fun quote(s: String): String = "\"" + s.replace("\\",
"\\\\").replace("\"", "\\\"") + "\""
+
+ private fun resolveField(
+ key: String,
+ prop: com.fasterxml.jackson.databind.JsonNode,
+ typeOverride: String?,
+ ): DslField? {
+ val ref = prop.path("\$ref").asText("").substringAfterLast('/')
+ val schemaType =
+ when {
+ ref == "timedelta" -> "timedelta"
+ ref == "datetime" -> "datetime"
+ ref.isNotEmpty() -> return null
+ prop.path("type").isTextual -> prop.path("type").asText()
+ else -> return null
+ }
+ val default = prop.path("default")
+ val defaultJson = if (default.isMissingNode || default.isNull) null
else default.toString()
+ val attribute = camelCase(key)
+ return when (schemaType) {
+ "string" ->
+ DslField(
+ key,
+ attribute,
+ "STRING",
+ "String",
+ quote(default.asText("")),
+ defaultJson,
+ "Airflow's `$key` setting.",
+ )
+ "boolean" ->
+ DslField(
+ key,
+ attribute,
+ "BOOLEAN",
+ "Boolean",
+ default.asBoolean(false).toString(),
+ defaultJson,
+ "Airflow's `$key` setting.",
+ )
+ "integer", "number" ->
+ if (typeOverride == "Double") {
+ DslField(
+ key,
+ attribute,
+ "NUMBER",
+ "Double",
+ if (defaultJson != null) default.asDouble().toString()
else "-1.0",
+ defaultJson,
+ "Airflow's `$key` setting.",
+ )
+ } else {
+ DslField(
+ key,
+ attribute,
+ "INTEGER",
+ "Int",
+ if (defaultJson != null) default.asInt().toString()
else "-1",
+ defaultJson,
+ "Airflow's `$key` setting.",
+ )
+ }
+ "timedelta" ->
+ DslField(
+ key,
+ attribute,
+ "TIMEDELTA",
+ "String",
+ if (defaultJson != null) {
+ quote(Duration.ofSeconds(default.asLong()).toString())
+ } else {
+ quote("")
+ },
+ defaultJson,
+ "Airflow's `$key` setting; an ISO-8601 duration such as
`\"PT5M\"`.",
+ )
+ "datetime" ->
+ DslField(
+ key,
+ attribute,
+ "DATETIME",
+ "String",
+ quote(""),
+ defaultJson,
+ "Airflow's `$key` setting; an ISO-8601 date-time such as
`\"2026-01-01T00:00:00Z\"`.",
+ )
+ "array" ->
+ if (prop.path("items").path("type").asText("") == "string" ||
key == "tags") {
+ DslField(
+ key,
+ attribute,
+ "STRING_ARRAY",
+ "Array<String>",
+ "[]",
+ defaultJson,
+ "Airflow's `$key` setting.",
+ )
+ } else {
+ null
+ }
+ else -> null
+ }
+ }
+
+ @TaskAction
+ fun generate() {
+ // Python-only "operator" keys deliberately not exposed, mirroring the
+ // Go SDK's TaskSpec generator exclusion list.
+ val excludedTaskKeys =
+ setOf(
+ "doc",
+ "doc_json",
+ "doc_yaml",
+ "doc_rst",
+ "allow_nested_operators",
+ "multiple_outputs",
+ "start_from_trigger",
+ "is_setup",
+ "is_teardown",
+ "on_failure_fail_dagrun",
+ // Edges come from before/after, never from configuration.
+ "downstream_task_ids",
+ )
+ // Dag-level keys exposed for configuration, mirroring the Go SDK's
+ // hand-curated DagSpec field list. "schedule" is virtual: the schema
+ // models it as the serializer-owned "timetable" object.
+ val dagAllowlist =
+ listOf(
+ "description",
+ "dag_display_name",
+ "doc_md",
+ "start_date",
+ "end_date",
+ "dagrun_timeout",
+ "tags",
+ "max_active_tasks",
+ "max_active_runs",
+ "max_consecutive_failed_dag_runs",
+ "catchup",
+ "fail_fast",
+ "render_template_as_native_obj",
+ "disable_bundle_versioning",
+ "is_paused_upon_creation",
+ )
+
+ val root =
+ com.fasterxml.jackson.databind
+ .ObjectMapper()
+ .readTree(schemaFile.get().asFile)
+ val dagProps = root.path("definitions").path("dag").path("properties")
+ val operator = root.path("definitions").path("operator")
+ val operatorRequired = operator.path("required").map { it.asText()
}.toSet()
+
+ val dagFields =
+ buildList {
+ add(
+ DslField(
+ "schedule",
+ "schedule",
+ "STRING",
+ "String",
+ quote(""),
+ null,
+ "`\"@once\"`, `\"@continuous\"`, a cron expression, or
empty for no schedule.",
+ ),
+ )
+ dagAllowlist.forEach { key ->
+ val prop = dagProps.path(key)
+ if (prop.isMissingNode) {
+ throw GradleException("Dag allowlist key '$key' is
missing from the schema; update the allowlist")
+ }
+ add(
+ resolveField(key, prop, null)
+ ?: throw GradleException("Dag allowlist key '$key'
is not a scalar the DSL can express"),
+ )
+ }
+ }
+
+ val excludedSeen = mutableSetOf<String>()
+ val taskFields =
+ buildList {
+ operator.path("properties").fields().forEach { (key, prop) ->
+ val serializerOwned =
+ key.startsWith("_") || key in operatorRequired ||
key.startsWith("has_on_")
+ if (serializerOwned) return@forEach
+ if (key in excludedTaskKeys) {
+ excludedSeen += key
+ return@forEach
+ }
+ // retry_exponential_backoff is "number" with an integral
+ // default, but Python declares it float (a backoff
+ // multiplier), so the mechanical mapping would pick Int.
+ val override = if (key == "retry_exponential_backoff")
"Double" else null
+ resolveField(key, prop, override)?.let { add(it) }
+ }
+ }
+ (excludedTaskKeys - excludedSeen).takeIf { it.isNotEmpty() }?.let {
+ throw GradleException("Excluded task keys match no eligible schema
property; remove or fix: $it")
+ }
+ // "id"/"to" name the annotations' structural attributes, so a schema
+ // key camel-casing to either would silently shadow them.
+ (dagFields + taskFields).firstOrNull { it.attribute == "id" ||
it.attribute == "to" }?.let {
+ throw GradleException("Schema key '${it.key}' collides with a
structural annotation attribute")
+ }
+
+ val outDir = targetDirectory.get().asFile.also {
it.deleteRecursively() }
+
+ fun attrLines(fields: List<DslField>): String =
+ fields.joinToString("\n") { f ->
+ " /** ${f.doc} */\n val ${f.attribute}: ${f.attrType} =
${f.attrDefault},"
+ }
+
+ outDir.resolve("org/apache/airflow/sdk").apply { mkdirs()
}.resolve("Builder.kt").writeText(
+ """
+ |package org.apache.airflow.sdk
+ |
+ |// Generated from the Dag serialization schema
(sdk/schema/dag-schema.json); do not edit by hand.
+ |
+ |/**
+ | * Container for the annotation-based Dag-authoring API.
+ | *
+ | * Annotating a class with [Dag] generates a `<Class>Builder`
whose static
+ | * `build()` returns the [DagDef] to add to a [Bundle].
+ | *
+ | * Example:
+ | *
+ | * ```java
+ | * @Builder.Dag(id = "my_pipeline", schedule = "@daily")
+ | * public class MyPipeline {
+ | *
+ | * @Builder.Task(id = "extract", retries = 2)
+ | * public long extract(Client client) { ... }
+ | *
+ | * @Builder.Task(id = "transform")
+ | * public long transform(Client client, long extracted) { ...
}
+ | * }
+ | * ```
+ | *
+ | * A task method's data parameters — everything other than the
injected
+ | * [Client] and [Context] — receive, by position, the arguments
the Python
+ | * `@task.stub` call site bound. Keyword arguments bind by name
instead
+ | * through a single [TaskInput] parameter.
+ | */
+ |class Builder internal constructor() {
+ | /**
+ | * Annotation to automate a Dag-builder pattern.
+ | *
+ | * When applied on a class Foo, this generates a FooBuilder
class with a
+ | * static build method to create the Dag structure
automatically.
+ | *
+ | * Configuration attributes are Airflow's own Dag settings; only
+ | * attributes written explicitly at the use site are applied, so
+ | * Airflow's own defaults apply to everything left out.
+ | */
+ | @Target(AnnotationTarget.CLASS)
+ | @MustBeDocumented
+ | annotation class Dag(
+ | /** Dag ID. Empty derives it from the annotated class's name.
*/
+ | val id: String = "",
+ | /** Name of the generated builder class. Empty derives
`<Class>Builder`. */
+ | val to: String = "",
+ |${attrLines(dagFields)}
+ | )
+ |
+ | /**
+ | * Annotation to automate task definition in a Dag-builder
pattern.
+ | *
+ | * Configuration attributes are Airflow's own task settings;
only
+ | * attributes written explicitly at the use site are applied.
+ | */
+ | @Target(AnnotationTarget.FUNCTION)
+ | @MustBeDocumented
+ | annotation class Task(
+ | /** Task ID. Empty derives it from the annotated function's
name. */
+ | val id: String = "",
+ |${attrLines(taskFields)}
+ | )
+ |
+ | /**
+ | * Marks a method as the Java body of a task the Python Dag file
+ | * declares with `@task.stub`.
+ | *
+ | * This is not [Task] under another name. Python declares the
task and
+ | * Java supplies only its body, so the handler names the pair
it binds
+ | * to rather than an id it owns:
+ | *
+ | * ```java
+ | * @Builder.TaskHandler(dag = "etl", task = "score")
+ | * public long score(Client client, long rows) { ... }
+ | * ```
+ | *
+ | * Register every handler a class holds with [Bundle.register];
there
+ | * is no [Dag] annotation on this surface, because the Dag is
the
+ | * Python file's.
+ | *
+ | * @param dag Dag ID as declared in the Python Dag file.
+ | * @param task Task ID as declared by the `@task.stub`
function. Empty
+ | * derives it from the annotated method's name.
+ | */
+ | @Target(AnnotationTarget.FUNCTION)
+ | @MustBeDocumented
+ | annotation class TaskHandler(
+ | val dag: String,
+ | val task: String = "",
+ | )
+ |}
+ |
+ """.trimMargin(),
+ )
+
+ fun tableLines(fields: List<DslField>): String =
+ fields.joinToString("\n") { f ->
+ val defaultRepr = f.defaultJson?.let { quote(it) } ?: "null"
+ listOf(
+ " \"${f.key}\" to",
+ " Field(",
+ " \"${f.key}\",",
+ " \"${f.attribute}\",",
+ " FieldType.${f.fieldType},",
+ " $defaultRepr,",
+ " ),",
+ ).joinToString("\n")
+ }
+
+ outDir.resolve("org/apache/airflow/sdk/internal").apply { mkdirs()
}.resolve("SchemaFields.kt").writeText(
+ """
+ |package org.apache.airflow.sdk.internal
+ |
+ |// Generated from the Dag serialization schema
(sdk/schema/dag-schema.json); do not edit by hand.
+ |
+ |/**
+ | * Configuration keys accepted by `DagDef.config` and
`TaskDef.config`,
+ | * keyed by Dag serialization schema property name. Public so
that the
+ | * annotation processor can lower `@Builder.Dag` / `@Builder.Task`
+ | * attributes onto the same tables; not user-facing API.
+ | */
+ |object SchemaFields {
+ | val DAG: Map<String, Field> =
+ | linkedMapOf(
+ |${tableLines(dagFields)}
+ | )
+ |
+ | val TASK: Map<String, Field> =
+ | linkedMapOf(
+ |${tableLines(taskFields)}
+ | )
+ |}
+ |
+ """.trimMargin(),
+ )
+ }
+}
+
val syncSupervisorSchema by tasks.registering(SyncSupervisorSchemaTask::class)
{
description = "Ensure the bundled Supervisor Schema is up-to-date with the
Gradle property."
schemaVersion = airflowSupervisorSchemaVersion
@@ -234,6 +657,19 @@ tasks.register<GeneratePointersTask>("generatePointers") {
targetDirectory = pointersDir
}
+val syncDagSchema by tasks.registering(SyncDagSchemaTask::class) {
+ description = "Refresh the vendored Dag serialization schema from the
monorepo copy when present."
+ sourceFile =
layout.projectDirectory.file("../../airflow-core/src/airflow/serialization/schema.json")
+ targetFile = dagSchemaInput
+}
+
+tasks.register<GenerateDagDslTask>("generateDagDsl") {
+ dependsOn(syncDagSchema)
+ description = "Generate the Builder.Dag/Builder.Task annotations and
SchemaFields from the Dag serialization schema"
+ schemaFile = dagSchemaInput
+ targetDirectory = dagDslDir
+}
+
val javadocJar by tasks.registering(Jar::class) {
description = "Assembles Javadoc JAR from Dokka output"
group = JavaBasePlugin.DOCUMENTATION_GROUP
@@ -262,6 +698,7 @@ sourceSets {
main {
java.srcDir(tasks.named("generateJsonSchema2Pojo").map {
schemaModelsDir })
kotlin.srcDir(tasks.named("generateDiscriminator").map {
discriminatorDir })
+ kotlin.srcDir(tasks.named("generateDagDsl").map { dagDslDir })
}
}
@@ -304,15 +741,15 @@ tasks.named("compileKotlin") {
}
tasks.named("runKtlintCheckOverMainSourceSet") {
- dependsOn("generateJsonSchema2Pojo", "generateDiscriminator")
+ dependsOn("generateJsonSchema2Pojo", "generateDiscriminator",
"generateDagDsl")
}
tasks.matching { it.name.startsWith("dokkaGenerate") }.configureEach {
- dependsOn("generateJsonSchema2Pojo", "generateDiscriminator")
+ dependsOn("generateJsonSchema2Pojo", "generateDiscriminator",
"generateDagDsl")
}
tasks.withType<Jar> {
- dependsOn("generateJsonSchema2Pojo", "generateDiscriminator")
+ dependsOn("generateJsonSchema2Pojo", "generateDiscriminator",
"generateDagDsl")
manifest {
attributes(
"Airflow-Supervisor-Schema-Version" to
airflowSupervisorSchemaVersion,
diff --git a/java-sdk/sdk/schema/dag-schema.json
b/java-sdk/sdk/schema/dag-schema.json
new file mode 100644
index 00000000000..086bccb34bb
--- /dev/null
+++ b/java-sdk/sdk/schema/dag-schema.json
@@ -0,0 +1,474 @@
+{
+ "$schema": "http://json-schema.org/draft-07/schema#",
+ "$id": "https://airflow.apache.com/schemas/serialized-dags.json",
+ "definitions": {
+ "datetime": {
+ "description": "A date time, stored as fractional seconds since the
epoch",
+ "type": "number"
+ },
+ "timedelta": {
+ "type": "number",
+ "minimum": 0
+ },
+ "typed_timedelta": {
+ "type": "object",
+ "properties": {
+ "__type": {
+ "type": "string",
+ "const": "timedelta"
+ },
+ "__var": { "$ref": "#/definitions/timedelta" }
+ },
+ "required": [
+ "__type",
+ "__var"
+ ],
+ "additionalProperties": false
+ },
+ "typed_relativedelta": {
+ "type": "object",
+ "description": "A dateutil.relativedelta.relativedelta object",
+ "properties": {
+ "__type": {
+ "type": "string",
+ "const": "relativedelta"
+ },
+ "__var": {
+ "type": "object",
+ "properties": {
+ "weekday": {
+ "type": "array",
+ "items": { "type": "integer" },
+ "minItems": 1,
+ "maxItems": 2
+ }
+ },
+ "additionalProperties": { "type": "integer" }
+ }
+ }
+ },
+ "timezone": {
+ "anyOf": [
+ { "type": "string" },
+ { "type": "integer" }
+ ]
+ },
+ "asset_definition": {
+ "type": "object",
+ "properties": {
+ "uri": { "type": "string" },
+ "name": { "type": "string" },
+ "group": { "type": "string" },
+ "extra": {
+ "anyOf": [
+ {"type": "null"},
+ { "$ref": "#/definitions/dict" }
+ ]
+ },
+ "watchers": {
+ "type": "array",
+ "items": { "$ref": "#/definitions/trigger" }
+ }
+ },
+ "required": [ "uri", "extra" ]
+ },
+ "asset": {
+ "type": "object",
+ "properties": {
+ "uri": { "type": "string" },
+ "extra": {
+ "anyOf": [
+ {"type": "null"},
+ { "$ref": "#/definitions/dict" }
+ ]
+ }
+ },
+ "required": [ "uri", "extra" ]
+ },
+ "typed_asset": {
+ "type": "object",
+ "properties": {
+ "__type": {
+ "type": "string",
+ "constant": "asset"
+ },
+ "__var": { "$ref": "#/definitions/asset" }
+ },
+ "required": [
+ "__type",
+ "__var"
+ ],
+ "additionalProperties": false
+ },
+ "typed_asset_cond": {
+ "type": "object",
+ "properties": {
+ "__type": {
+ "anyOf": [{
+ "type": "string",
+ "constant": "asset_or"
+ },
+ {
+ "type": "string",
+ "constant": "asset_and"
+ }
+ ]
+ },
+ "__var": {
+ "type": "array",
+ "items": {
+ "anyOf": [
+ {"$ref": "#/definitions/typed_asset"},
+ { "$ref": "#/definitions/typed_asset_cond"}
+ ]
+ }
+ }
+ },
+ "required": [
+ "__type",
+ "__var"
+ ],
+ "additionalProperties": false
+ },
+ "trigger": {
+ "type": "object",
+ "properties": {
+ "classpath": { "type": "string" },
+ "kwargs": { "$ref": "#/definitions/dict" }
+ },
+ "required": [ "classpath", "kwargs" ]
+ },
+ "dict": {
+ "description": "A python dictionary containing values of any type",
+ "type": "object"
+ },
+ "arg_binding": {
+ "$comment": "One captured TaskFlow call argument of a @task.stub task.
Materialized directly by _serialize_node, so it stays plain JSON with no
{__type, __var} encoding. The object stays open so future binding fields keep
validating on older cores",
+ "type": "object",
+ "properties": {
+ "name": { "type": "string" },
+ "kind": { "type": "string", "enum": [ "xcom", "literal" ] },
+ "value_schema": { "type": "object" },
+ "task_id": { "type": "string" },
+ "value": {},
+ "from_default": { "type": "boolean" }
+ },
+ "required": [ "name", "kind" ]
+ },
+ "color": {
+ "type": "string",
+ "pattern": "^#[a-fA-F0-9]{3,6}$"
+ },
+ "extra_links": {
+ "type": "array",
+ "items": {
+ "type": "object",
+ "minProperties": 1,
+ "maxProperties": 1
+ }
+ },
+ "dag_dependencies": {
+ "type": "array",
+ "items": {
+ "type": "object"
+ }
+ },
+ "dag": {
+ "type": "object",
+ "properties": {
+ "params": { "$ref": "#/definitions/params" },
+ "dag_id": { "type": "string" },
+ "tasks": { "$ref": "#/definitions/tasks" },
+ "timezone": { "$ref": "#/definitions/timezone" },
+ "owner_links": { "type": "object" },
+ "timetable": {
+ "type": "object",
+ "properties": {
+ "type": { "type": "string" },
+ "value": { "$ref": "#/definitions/dict" }
+ }
+ },
+ "catchup": { "type": "boolean" },
+ "allowed_run_types": {
+ "anyOf": [
+ { "type": "array", "items": { "type": "string" } },
+ { "type": "null" }
+ ]
+ },
+ "fail_fast": { "type": "boolean", "default": false },
+ "fileloc": { "type" : "string"},
+ "relative_fileloc": { "type" : "string"},
+ "bundle_name": { "anyOf": [{ "type": "null" }, { "type": "string" }] },
+ "_processor_dags_folder": {
+ "anyOf": [
+ { "type": "null" },
+ { "type": "string" }
+ ]
+ },
+ "dag_display_name": { "type" : "string"},
+ "description": { "type" : "string"},
+ "deadline": {
+ "anyOf": [
+ { "$ref": "#/definitions/dict" },
+ {
+ "type": "array",
+ "items": { "$ref": "#/definitions/dict" }
+ },
+ {
+ "$comment": "Once persisted, a Dag's deadline alerts live
as rows in the deadline_alert table and the serialized Dag keeps only a list of
UUID strings referencing them (see
SerializedDagModel._generate_deadline_uuids). This branch lets the stored form
validate at any lifecycle stage, not only before the dict->UUID rewrite.",
+ "type": "array",
+ "items": { "type": "string",
+ "pattern":
"^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$"
+ }
+ },
+ { "type": "null" }
+ ]
+ },
+ "_concurrency": { "type" : "number"},
+ "max_active_tasks": { "type" : "number" },
+ "max_active_runs": { "type" : "number" },
+ "max_consecutive_failed_dag_runs": { "type" : "number" },
+ "default_args": { "$ref": "#/definitions/dict" },
+ "start_date": { "$ref": "#/definitions/datetime" },
+ "end_date": { "$ref": "#/definitions/datetime" },
+ "dagrun_timeout": { "$ref": "#/definitions/timedelta" },
+ "doc_md": { "type" : "string"},
+ "access_control": {"$ref": "#/definitions/dict" },
+ "is_paused_upon_creation": { "type": "boolean" },
+ "has_on_success_callback": { "type": "boolean", "default": false },
+ "has_on_failure_callback": { "type": "boolean", "default": false },
+ "render_template_as_native_obj": { "type": "boolean", "default":
false },
+ "tags": { "type": "array" },
+ "task_group": {"anyOf": [
+ { "type": "null" },
+ { "$ref": "#/definitions/task_group" }
+ ]},
+ "edge_info": { "$ref": "#/definitions/edge_info" },
+ "dag_dependencies": { "$ref": "#/definitions/dag_dependencies" },
+ "disable_bundle_versioning": {"type": "boolean" },
+ "rerun_with_latest_version": {"type": ["boolean", "null"], "default":
null}
+ },
+ "required": [
+ "dag_id",
+ "fileloc",
+ "tasks"
+ ],
+ "additionalProperties": false
+ },
+ "tasks": {
+ "type": "array",
+ "additionalProperties": { "$ref": "#/definitions/operator" }
+ },
+ "params": {
+ "type": "array",
+ "prefixItems": [
+ { "type": "string" },
+ { "$ref": "#/definitions/param" }
+ ],
+ "unevaluatedItems": false
+ },
+ "param": {
+ "$comment": "A param for a dag / operator",
+ "type": "object",
+ "required": [
+ "__class",
+ "default"
+ ],
+ "properties": {
+ "__class": { "type": "string" },
+ "default": {},
+ "description": {"anyOf": [{"type":"string"}, {"type":"null"}]},
+ "schema": { "$ref": "#/definitions/dict" }
+ }
+ },
+ "operator": {
+ "$comment": "A task/operator in a DAG",
+ "type": "object",
+ "required": [
+ "task_type",
+ "_task_module",
+ "task_id",
+ "ui_color",
+ "ui_fgcolor",
+ "template_fields"
+ ],
+ "properties": {
+ "task_type": { "type": "string", "default": "BaseOperator"},
+ "_task_module": { "type": "string" },
+ "_operator_extra_links": { "$ref": "#/definitions/extra_links" },
+ "task_id": { "type": "string" },
+ "_task_display_name": { "type": "string" },
+ "owner": { "type": "string", "default": "airflow" },
+ "start_date": { "$ref": "#/definitions/datetime" },
+ "end_date": { "$ref": "#/definitions/datetime" },
+ "trigger_rule": { "type": "string", "default": "all_success" },
+ "depends_on_past": { "type": "boolean", "default": false },
+ "ignore_first_depends_on_past": { "type": "boolean", "default": false
},
+ "wait_for_past_depends_before_skipping": { "type": "boolean",
"default": false },
+ "wait_for_downstream": { "type": "boolean", "default": false },
+ "retries": { "type": "number", "default": 0 },
+ "queue": { "type": "string", "default": "default" },
+ "pool": { "type": "string", "default": "default_pool" },
+ "pool_slots": { "type": "number", "default": 1 },
+ "execution_timeout": { "$ref": "#/definitions/timedelta" },
+ "retry_delay": { "$ref": "#/definitions/timedelta", "default": 300.0 },
+ "retry_exponential_backoff": { "type": "number", "default": 0 },
+ "max_retry_delay": { "$ref": "#/definitions/timedelta" },
+ "params": { "$ref": "#/definitions/params" },
+ "priority_weight": { "type": "number", "default": 1 },
+ "weight_rule": { "type": "string", "default": "downstream" },
+ "executor": { "type": "string" },
+ "executor_config": { "$ref": "#/definitions/dict" },
+ "do_xcom_push": { "type": "boolean", "default": true },
+ "email_on_failure": { "type": "boolean", "default": true },
+ "email_on_retry": { "type": "boolean", "default": true },
+ "ui_color": { "type": "string", "default": "#fff" },
+ "ui_fgcolor": { "type": "string", "default": "#000" },
+ "template_fields": {
+ "type": "array",
+ "items": { "type": "string" },
+ "default": []
+ },
+ "template_ext": {"type": "array", "default": []},
+ "template_fields_renderers": {"$ref": "#/definitions/dict", "default":
{}},
+ "downstream_task_ids": {
+ "type": "array",
+ "items": { "type": "string" },
+ "default": []
+ },
+ "doc": { "type": "string" },
+ "doc_md": { "type": "string" },
+ "doc_json": { "type": "string" },
+ "doc_yaml": { "type": "string" },
+ "doc_rst": { "type": "string" },
+ "_logger_name": { "type": "string" },
+ "_needs_expansion": { "type": "boolean"},
+ "_is_mapped": { "const": true, "$comment": "only present when True",
"default": false },
+ "_is_sensor": { "const": true, "$comment": "only present when True",
"default": false },
+ "partial_kwargs": { "type": "object" },
+ "_disallow_kwargs_override": { "type": "boolean"},
+ "_expand_input_attr": { "type": "string" },
+ "map_index_template": { "type": "string" },
+ "allow_nested_operators": { "type": "boolean", "default": true },
+ "render_template_as_native_obj": { "anyOf": [{"type": "boolean"},
{"type": "null"}], "default": null },
+ "inlets": {"type": "array", "default": []},
+ "outlets": {"type": "array", "default": []},
+ "has_on_execute_callback": {"type": "boolean", "default": false},
+ "has_on_failure_callback": {"type": "boolean", "default": false},
+ "has_on_skipped_callback": {"type": "boolean", "default": false},
+ "has_on_success_callback": {"type": "boolean", "default": false},
+ "has_on_retry_callback": {"type": "boolean", "default": false},
+ "multiple_outputs": {"type": "boolean", "default": false},
+ "start_from_trigger": {"type": "boolean", "default": false},
+ "start_trigger_args": {"type": "object", "default": null},
+ "is_setup": {"type": "boolean", "default": false},
+ "is_teardown": {"type": "boolean", "default": false},
+ "on_failure_fail_dagrun": {"type": "boolean", "default": false},
+ "max_active_tis_per_dag": {"type": "integer"},
+ "max_active_tis_per_dagrun": {"type": "integer"},
+ "_arg_bindings": {
+ "$comment": "Used for mixed-language tasks (the @task.stub operator)
or tasks defined in a non-Python SDK Dag",
+ "type": "array",
+ "items": { "$ref": "#/definitions/arg_binding" }
+ }
+ },
+ "dependencies": {
+ "expand_input": ["partial_kwargs", "_is_mapped"],
+ "partial_kwargs": ["expand_input", "_is_mapped"],
+ "_is_mapped": ["expand_input", "partial_kwargs"]
+ },
+ "additionalProperties": true
+ },
+ "task_group": {
+ "$comment": "A TaskGroup containing tasks",
+ "type": "object",
+ "required": [
+ "_group_id",
+ "group_display_name",
+ "prefix_group_id",
+ "children",
+ "tooltip",
+ "ui_color",
+ "ui_fgcolor",
+ "upstream_group_ids",
+ "downstream_group_ids",
+ "upstream_task_ids",
+ "downstream_task_ids"
+ ],
+ "properties": {
+ "_group_id": {"anyOf": [{"type": "null"}, { "type": "string" }]},
+ "group_display_name": {"type": "string" },
+ "is_mapped": { "type": "boolean" },
+ "prefix_group_id": { "type": "boolean" },
+ "children": { "$ref": "#/definitions/dict" },
+ "tooltip": { "type": "string" },
+ "doc_md": {
+ "anyOf": [
+ { "type": "string" },
+ { "type": "null" }
+ ]},
+ "ui_color": { "type": "string" },
+ "ui_fgcolor": { "type": "string" },
+ "upstream_group_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ },
+ "downstream_group_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ },
+ "upstream_task_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ },
+ "downstream_task_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ }
+ },
+ "additionalProperties": false
+ },
+ "edge_info": {
+ "$comment": "Metadata about DAG edges",
+ "type": "object",
+ "additionalProperties": {
+ "type": "object",
+ "additionalProperties": {
+ "type": "object",
+ "properties": {
+ "label": { "type": "string" }
+ },
+ "required": ["label"],
+ "additionalProperties": false
+ }
+ }
+ }
+ },
+
+ "type": "object",
+ "allOf": [
+ {
+ "type": "object",
+ "properties": {
+ "__version": {
+ "type": "integer",
+ "exclusiveMinimum": 0
+ },
+ "dag": { "$ref": "#/definitions/dag" },
+ "client_defaults": {
+ "type": "object",
+ "description": "SDK-specific default values that differ from schema
defaults",
+ "properties": {
+ "tasks": {
+ "type": "object",
+ "description": "Task-level default overrides"
+ }
+ },
+ "additionalProperties": false
+ }
+ },
+ "additionalProperties": false,
+ "required": [ "__version", "dag" ]
+ }
+ ]
+}
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Arg.kt
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Arg.kt
index 1834a548e6a..8bedb4f4cb8 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Arg.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Arg.kt
@@ -50,4 +50,22 @@ class TaskRef<T> internal constructor(
super<Deps.Flow>.after(*previous)
return this
}
+
+ /**
+ * Sets one task-level configuration value, so a task built through
+ * [DagDef.task] is configured where it is created.
+ *
+ * @param key Airflow task setting name.
+ * @param value Value matching the key's schema type.
+ * @return This handle, for chaining.
+ * @throws IllegalArgumentException if the key is unknown or the value type
+ * does not match.
+ */
+ fun config(
+ key: String,
+ value: Any?,
+ ): TaskRef<T> {
+ def.config(key, value)
+ return this
+ }
}
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Builder.kt
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Builder.kt
deleted file mode 100644
index 8d00b7ca7a9..00000000000
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Builder.kt
+++ /dev/null
@@ -1,109 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied. See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-
-package org.apache.airflow.sdk
-
-/**
- * Container for the annotation-based Dag-authoring API.
- *
- * This class is not instantiated directly. Its nested annotations drive the
- * `BuilderProcessor` annotation processor in the :processor project,
- * which generates a `*Builder` class for each class annotated with
[Builder.Dag].
- *
- * Example:
- *
- * ```java
- * @Builder.Dag(id = "my_pipeline")
- * public class MyPipeline {
- *
- * @Builder.Task(id = "extract")
- * public long extract(Client client) { ... }
- *
- * @Builder.Task(id = "transform")
- * public long transform(Client client, long extracted) { ... }
- * }
- * ```
- *
- * A task method's data parameters — everything other than the injected
- * [Client] and [Context] — receive the arguments the Python `@task.stub` call
- * site bound, by position. Keyword arguments bind by name instead through a
- * single [TaskInput] parameter.
- *
- * The processor generates `MyPipelineBuilder.build()`, which returns a
- * fully wired-up [DagDef] ready to add to a [Bundle].
- */
-class Builder internal constructor() {
- /**
- * Annotation to automate a Dag-builder pattern.
- *
- * When applied on a class Foo, this generates a FooBuilder class with a
- * static build method to create the Dag structure automatically.
- *
- * @param id Override the Dag ID. If empty or not provided, the annotated
- * class's name is used by default.
- * @param to Name of the Dag-builder class. If empty or not provided, use the
- * annotated class name + "Builder".
- */
- @Target(AnnotationTarget.CLASS)
- @MustBeDocumented
- annotation class Dag(
- val id: String = "",
- val to: String = "",
- )
-
- /**
- * Annotation to automate task definition in a Dag-builder pattern.
- *
- * @param id Override the task ID. If empty or not provided, the annotated
- * function's name is used by default.
- */
- @Target(AnnotationTarget.FUNCTION)
- @MustBeDocumented
- annotation class Task(
- val id: String = "",
- )
-
- /**
- * Marks a method as the Java body of a task the Python Dag file declares
- * with `@task.stub`.
- *
- * This is not [Task] under another name. Python declares the task and Java
- * supplies only its body, so the handler names the pair it binds to rather
- * than an id it owns — Python owns both, and the processor generates the
- * registration from them:
- *
- * ```java
- * @Builder.TaskHandler(dag = "etl", task = "score")
- * public long score(Client client, long rows, double threshold) { ... }
- * ```
- *
- * Register every handler a class holds with [Bundle.register]; there is no
- * [Dag] annotation on this surface, because the Dag is the Python file's.
- *
- * @param dag Dag ID as declared in the Python Dag file.
- * @param task Task ID as declared by the `@task.stub` function. Empty
- * derives it from the annotated method's name.
- */
- @Target(AnnotationTarget.FUNCTION)
- @MustBeDocumented
- annotation class TaskHandler(
- val dag: String,
- val task: String = "",
- )
-}
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
index 90f77d86dee..6ba3ff7f165 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
@@ -19,6 +19,8 @@
package org.apache.airflow.sdk
+import org.apache.airflow.sdk.internal.SchemaFields
+import org.apache.airflow.sdk.internal.checkConfigValue
import org.apache.airflow.sdk.internal.validateTaskInput
import kotlin.Throws
@@ -32,8 +34,8 @@ import kotlin.Throws
* class directly if you need to do low-level plumbing:
*
* ```java
- * var dag = new DagDef("java_etl");
- * var extract = dag.task("extract", Extract.class);
+ * var dag = new DagDef("java_etl").config("schedule", "@daily");
+ * var extract = dag.task("extract", Extract.class).config("retries", 2);
* extract.before(dag.task("load", Load.class));
* ```
*
@@ -46,6 +48,31 @@ class DagDef(
val id: String, // TODO: charset check?
) {
internal val tasks = linkedMapOf<String, TaskDef>()
+ internal val dagConfig = linkedMapOf<String, Any>()
+
+ /**
+ * Sets one Dag-level configuration value.
+ *
+ * Keys are Airflow's own Dag setting names (for example `"schedule"`,
+ * `"description"`, `"tags"`, `"catchup"`); unknown keys and
+ * mismatched value types are rejected on the call, so mistakes surface where
+ * the Dag is defined.
+ *
+ * @param key Airflow Dag setting name.
+ * @param value Value matching the key's schema type. Durations take
+ * [java.time.Duration], date-times [java.time.OffsetDateTime] or
+ * [java.time.Instant], string arrays any `Iterable` of `String`.
+ * @return This Dag, for chaining.
+ * @throws IllegalArgumentException if the key is unknown or the value type
+ * does not match.
+ */
+ fun config(
+ key: String,
+ value: Any?,
+ ): DagDef {
+ dagConfig[key] = checkConfigValue("Dag", SchemaFields.DAG, key, value)
+ return this
+ }
/**
* Registers a task from its ID and implementation class.
@@ -63,12 +90,11 @@ class DagDef(
): DagDef = addTask(TaskDef(id, definition))
/**
- * Creates a task, registers it, and hands back its handle — so there is no
- * second `addTask` call to forget, and the handle is ready to wire edges
- * with [Deps.Flow.before].
+ * Creates a task, registers it, and hands back its handle, ready to carry
+ * configuration and to wire edges with [Deps.Flow.before].
*
* ```java
- * var extract = dag.task("extract", Extract.class);
+ * var extract = dag.task("extract", Extract.class).config("retries", 2);
* var load = dag.task("load", Load.class);
* extract.before(load);
* ```
@@ -116,14 +142,14 @@ class DagDef(
}
/**
- * One task definition: its ID, the class that implements it, and its upstream
- * dependencies.
+ * One task definition: its ID, the class that implements it, its upstream
+ * dependencies, and its task-level configuration.
*
* Edges are drawn on the handles that [DagDef.task] returns, not here:
*
* ```java
* var dag = new DagDef("java_etl");
- * dag.addTask(new TaskDef("extract", Extract.class));
+ * dag.addTask(new TaskDef("extract", Extract.class).config("retries", 2));
* ```
*
* @param id Task identifier, unique within a [DagDef].
@@ -143,9 +169,34 @@ class TaskDef(
validateTaskInput(definition)
}
+ internal val configValues = linkedMapOf<String, Any>()
internal val upstreams = linkedSetOf<TaskDef>()
internal var owner: DagDef? = null
+ /**
+ * Sets one task-level configuration value.
+ *
+ * Keys are Airflow's own task setting names (for example `"retries"`,
+ * `"queue"`, `"retry_delay"`); unknown keys and mismatched
+ * value types are rejected on the call, so mistakes surface where the task
is
+ * defined.
+ *
+ * @param key Airflow task setting name, e.g. `"retries"`.
+ * @param value Value matching the key's schema type. Durations take
+ * [java.time.Duration], date-times [java.time.OffsetDateTime] or
+ * [java.time.Instant].
+ * @return This task definition, for chaining.
+ * @throws IllegalArgumentException if the key is unknown or the value type
+ * does not match.
+ */
+ fun config(
+ key: String,
+ value: Any?,
+ ): TaskDef {
+ configValues[key] = checkConfigValue("task", SchemaFields.TASK, key, value)
+ return this
+ }
+
/** Records that this task runs after [upstreams], backing
[Deps.Flow.before] and [Deps.Flow.after]. */
internal fun dependsOn(vararg upstreams: TaskDef): TaskDef {
this.upstreams += upstreams
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
index 7860a69a0c8..38a14d00354 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
@@ -24,20 +24,19 @@ interface Deps {
/**
* A point in the task graph: one task, or a set of them.
*
- * [Flow] is Java's spelling of Python's `>>` and `<<`, for a dependency
- * where nothing flows but the ordering. An edge that carries a value is
- * declared by passing the upstream's handle instead.
+ * [Flow] declares a dependency where nothing flows but the ordering. An
+ * edge that carries a value is declared by passing the upstream's handle
+ * instead.
*/
interface Flow {
/** The tasks at this point in the flow. */
fun nodes(): List<TaskDef>
/**
- * Runs the tasks here before each of [next], carrying no value — Java's
- * spelling of Python's `>>`.
+ * Runs the tasks here before each of [next], carrying no value.
*
* ```java
- * loaded.before(cleaned, notified); // load >> [cleanup, notify]
+ * loaded.before(cleaned, notified); // cleanup and notify both wait for
load
* ```
*
* Variadic, so one call fans out, and it returns its own receiver: a
@@ -54,11 +53,10 @@ interface Deps {
}
/**
- * Runs the tasks here after each of [previous], carrying no value —
- * Python's `<<`.
+ * Runs the tasks here after each of [previous], carrying no value.
*
* ```java
- * cleaned.after(loaded, transformed); // [load, transform] >> cleanup
+ * cleaned.after(loaded, transformed); // cleanup waits for load and
transform
* ```
*
* @param previous Tasks that run before the ones here.
@@ -76,7 +74,7 @@ interface Deps {
* every edge between two sets:
*
* ```java
- * Flow.of(a, b).before(c, d); // [a, b] >> [c, d]
+ * Flow.of(a, b).before(c, d); // c and d both wait for a and b
* ```
*/
@JvmStatic
diff --git
a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Fields.kt
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Fields.kt
new file mode 100644
index 00000000000..cc43db31865
--- /dev/null
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Fields.kt
@@ -0,0 +1,107 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.airflow.sdk.internal
+
+import java.time.Duration
+import java.time.Instant
+import java.time.OffsetDateTime
+
+/**
+ * Value shape of one Dag serialization schema key. Public so that the
+ * annotation processor can lower `@Builder.Dag` / `@Builder.Task` attributes;
+ * not user-facing API.
+ */
+enum class FieldType {
+ STRING,
+ INTEGER,
+ NUMBER,
+ BOOLEAN,
+ STRING_ARRAY,
+ TIMEDELTA,
+ DATETIME,
+}
+
+/**
+ * One configuration key from the Dag serialization schema. Public so that the
+ * annotation processor can lower `@Builder.Dag` / `@Builder.Task` attributes;
+ * not user-facing API.
+ *
+ * @property key Schema property name, e.g. `retry_delay`.
+ * @property attribute Annotation attribute name, e.g. `retryDelay`.
+ * @property type Accepted value shape.
+ * @property defaultJson Schema default as raw JSON, or `null` when the schema
+ * declares no default.
+ */
+class Field(
+ val key: String,
+ val attribute: String,
+ val type: FieldType,
+ val defaultJson: String?,
+)
+
+/**
+ * Validates one `config(key, value)` call against a schema field table and
+ * returns the value to store.
+ *
+ * @throws IllegalArgumentException if the key is not a configurable schema
+ * key or the value does not match the key's type.
+ */
+internal fun checkConfigValue(
+ scope: String,
+ table: Map<String, Field>,
+ key: String,
+ value: Any?,
+): Any {
+ val field =
+ requireNotNull(table[key]) {
+ "Unknown $scope config key: '$key'"
+ }
+ requireNotNull(value) {
+ "Value for $scope config key '$key' must not be null"
+ }
+
+ fun mismatch(expected: String): Nothing =
+ throw IllegalArgumentException(
+ "Value for $scope config key '$key' must be $expected, got:
${value.javaClass.name}",
+ )
+ return when (field.type) {
+ FieldType.STRING -> value as? String ?: mismatch("a String")
+ FieldType.BOOLEAN -> value as? Boolean ?: mismatch("a Boolean")
+ FieldType.NUMBER -> value as? Number ?: mismatch("a Number")
+ FieldType.INTEGER ->
+ when (value) {
+ is Byte, is Short, is Int, is Long -> value
+ else -> mismatch("an integral Number")
+ }
+ FieldType.TIMEDELTA -> value as? Duration ?: mismatch("a
java.time.Duration")
+ FieldType.DATETIME ->
+ when (value) {
+ is OffsetDateTime -> value
+ is Instant -> value
+ else -> mismatch("a java.time.OffsetDateTime or java.time.Instant")
+ }
+ FieldType.STRING_ARRAY ->
+ when {
+ value is Iterable<*> && value.all { it is String } -> value.map { it
as String }
+ value is Array<*> && value.all { it is String } -> value.map { it as
String }
+ else -> mismatch("an Iterable of String")
+ }
+ }
+}
diff --git a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/DagDefTest.kt
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/DagDefTest.kt
index 5ad4125e985..70bb24552e0 100644
--- a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/DagDefTest.kt
+++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/DagDefTest.kt
@@ -23,6 +23,8 @@ import org.apache.airflow.sdk.Deps.Flow
import org.junit.jupiter.api.Assertions
import org.junit.jupiter.api.DisplayName
import org.junit.jupiter.api.Test
+import java.time.Duration
+import java.time.OffsetDateTime
internal class DagDefTest {
private class NoOp : Task {
@@ -161,6 +163,183 @@ internal class DagDefTest {
Assertions.assertEquals(setOf("left", "right"), upstreamsOf(dag, "join"))
}
+ @Test
+ @DisplayName("Should store validated dag config values keyed by schema name")
+ fun shouldStoreDagConfigValues() {
+ val dag =
+ DagDef("dag")
+ .config("schedule", "@daily")
+ .config("description", "demo")
+ .config("catchup", true)
+ .config("max_active_runs", 3)
+ .config("dagrun_timeout", Duration.ofMinutes(5))
+ .config("start_date", OffsetDateTime.parse("2026-01-01T00:00:00Z"))
+ .config("tags", listOf("a", "b"))
+
+ Assertions.assertEquals(
+ mapOf(
+ "schedule" to "@daily",
+ "description" to "demo",
+ "catchup" to true,
+ "max_active_runs" to 3,
+ "dagrun_timeout" to Duration.ofMinutes(5),
+ "start_date" to OffsetDateTime.parse("2026-01-01T00:00:00Z"),
+ "tags" to listOf("a", "b"),
+ ),
+ dag.dagConfig,
+ )
+ }
+
+ @Test
+ @DisplayName("Should reject unknown dag config keys")
+ fun shouldRejectUnknownDagConfigKey() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("scheduel", "@daily")
+ }
+
+ Assertions.assertEquals("Unknown Dag config key: 'scheduel'",
error.message)
+ }
+
+ @Test
+ @DisplayName("Should reject dag config values of the wrong type")
+ fun shouldRejectMismatchedDagConfigValue() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("catchup", "yes")
+ }
+
+ Assertions.assertEquals(
+ "Value for Dag config key 'catchup' must be a Boolean, got:
java.lang.String",
+ error.message,
+ )
+ }
+
+ @Test
+ @DisplayName("Should reject null dag config values")
+ fun shouldRejectNullDagConfigValue() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("description", null)
+ }
+
+ Assertions.assertEquals("Value for Dag config key 'description' must not
be null", error.message)
+ }
+
+ @Test
+ @DisplayName("Should reject non-integral values for integer dag config keys")
+ fun shouldRejectFractionalIntegerValue() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("max_active_runs", 1.5)
+ }
+
+ Assertions.assertEquals(
+ "Value for Dag config key 'max_active_runs' must be an integral Number,
got: java.lang.Double",
+ error.message,
+ )
+ }
+
+ @Test
+ @DisplayName("Should reject an ISO-8601 string where a duration dag config
key expects a Duration")
+ fun shouldRejectStringForDurationDagConfigKey() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("dagrun_timeout", "PT5M")
+ }
+
+ Assertions.assertEquals(
+ "Value for Dag config key 'dagrun_timeout' must be a java.time.Duration,
got: java.lang.String",
+ error.message,
+ )
+ }
+
+ @Test
+ @DisplayName("Should reject an ISO-8601 string where a date-time dag config
key expects a date-time")
+ fun shouldRejectStringForDateTimeDagConfigKey() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("start_date", "2026-01-01")
+ }
+
+ Assertions.assertEquals(
+ "Value for Dag config key 'start_date' must be a
java.time.OffsetDateTime or java.time.Instant, " +
+ "got: java.lang.String",
+ error.message,
+ )
+ }
+
+ @Test
+ @DisplayName("Should reject a string-array dag config key whose elements are
not strings")
+ fun shouldRejectNonStringElementsForArrayDagConfigKey() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ DagDef("dag").config("tags", arrayListOf(1, 2))
+ }
+
+ Assertions.assertEquals(
+ "Value for Dag config key 'tags' must be an Iterable of String, got:
java.util.ArrayList",
+ error.message,
+ )
+ }
+
+ @Test
+ @DisplayName("Should store validated task config values on the task
definition")
+ fun shouldStoreTaskConfigValues() {
+ val def =
+ TaskDef("extract", NoOp::class.java)
+ .config("retries", 2)
+ .config("queue", "q")
+ .config("retry_delay", Duration.ofMinutes(5))
+ .config("retry_exponential_backoff", 1.5)
+
+ Assertions.assertEquals(
+ mapOf(
+ "retries" to 2,
+ "queue" to "q",
+ "retry_delay" to Duration.ofMinutes(5),
+ "retry_exponential_backoff" to 1.5,
+ ),
+ def.configValues,
+ )
+ }
+
+ @Test
+ @DisplayName("Should reject unknown task config keys")
+ fun shouldRejectUnknownTaskConfigKey() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ TaskDef("extract", NoOp::class.java).config("retrys", 1)
+ }
+
+ Assertions.assertEquals("Unknown task config key: 'retrys'", error.message)
+ }
+
+ @Test
+ @DisplayName("Should reject an ISO-8601 string where a duration task config
key expects a Duration")
+ fun shouldRejectStringForDurationTaskConfigKey() {
+ val error =
+ Assertions.assertThrows(IllegalArgumentException::class.java) {
+ TaskDef("extract", NoOp::class.java).config("retry_delay", "PT5M")
+ }
+
+ Assertions.assertEquals(
+ "Value for task config key 'retry_delay' must be a java.time.Duration,
got: java.lang.String",
+ error.message,
+ )
+ }
+
+ @Test
+ @DisplayName("Should configure a task through the handle dag.task returns")
+ fun shouldConfigureTaskThroughHandle() {
+ val dag = DagDef("dag")
+
+ val extract = dag.task<Long>("extract",
NoOp::class.java).config("retries", 2)
+
+ Assertions.assertEquals(mapOf("retries" to 2),
dag.tasks.getValue("extract").configValues)
+ Assertions.assertEquals(listOf(dag.tasks.getValue("extract")),
extract.nodes())
+ }
+
private fun upstreamsOf(
dag: DagDef,
taskId: String,