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 793eda354121a..c4edaef79bd56 100644 --- a/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst +++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst @@ -352,7 +352,8 @@ Annotate a plain Java class and let the SDK generate the boilerplate at compile * - ``@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`. + explicitly are applied. ``queue`` is the queue each task runs on unless the task sets its + own. 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 @@ -595,6 +596,11 @@ 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. +Every task of the Dag runs on the Java coordinator, so it needs a queue that +:ref:`queue_to_coordinator ` sends there. Set ``queue`` once on the Dag, with +``@Builder.Dag(queue = "java")`` or ``dag.config("queue", "java")``, and each task inherits it, including a +``TriggerDagRun`` task. A task's own ``queue`` wins over the Dag's. + Building the Dag in Java ~~~~~~~~~~~~~~~~~~~~~~~~ diff --git a/java-sdk/example/src/java/org/apache/airflow/example/nativedag/AnnotationExample.java b/java-sdk/example/src/java/org/apache/airflow/example/nativedag/AnnotationExample.java index 0bb55ceacd0f5..973cf09948a17 100644 --- a/java-sdk/example/src/java/org/apache/airflow/example/nativedag/AnnotationExample.java +++ b/java-sdk/example/src/java/org/apache/airflow/example/nativedag/AnnotationExample.java @@ -31,6 +31,7 @@ id = "java_native_annotation_example", description = "Pure-Java Dag authored with annotations", schedule = "@daily", + queue = "java", startDate = "2026-01-01T00:00:00Z", catchup = false, tags = {"example", "java-sdk"}) 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 fcf095ee72ec7..06365e155f5b4 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 @@ -101,6 +101,7 @@ public static DagDef build() { new DagDef("java_native_interface_example") .config("description", "Pure-Java Dag authored with the interface API") .config("schedule", "@daily") + .config("queue", "java") .config("catchup", false) .config("tags", List.of("example", "java-sdk")); diff --git a/java-sdk/example/src/java/org/apache/airflow/example/nativedag/TargetExample.java b/java-sdk/example/src/java/org/apache/airflow/example/nativedag/TargetExample.java index cf8d451429545..f428650dc51b8 100644 --- a/java-sdk/example/src/java/org/apache/airflow/example/nativedag/TargetExample.java +++ b/java-sdk/example/src/java/org/apache/airflow/example/nativedag/TargetExample.java @@ -41,6 +41,7 @@ public static DagDef build() { var dag = new DagDef("java_native_target_example") .config("description", "Pure-Java Dag that the other native examples trigger") + .config("queue", "java") .config("catchup", false) .config("tags", List.of("example", "java-sdk")); dag.task("receive", Receive.class); 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 5d8c89a058b06..0c4d920e06ea3 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 @@ -403,7 +403,7 @@ class BuilderTest { """ package org.apache.airflow.example; import org.apache.airflow.sdk.Builder; - @Builder.Dag(id = "cfg", schedule = "@daily", tags = {"a", "b"}, catchup = true, + @Builder.Dag(id = "cfg", schedule = "@daily", queue = "java", tags = {"a", "b"}, catchup = true, startDate = "2026-01-01T00:00:00Z") public class TestExample { @Builder.Task(retries = 2, queue = "q", retryDelay = "PT5M", retryExponentialBackoff = 1.5) @@ -442,6 +442,7 @@ class BuilderTest { public static DagDef build() { var dag = DagSource.declaredBy(new DagDef("cfg"), TestExample.class); dag.config("schedule", "@daily"); + dag.config("queue", "java"); dag.config("tags", List.of("a", "b")); dag.config("catchup", true); dag.config("start_date", OffsetDateTime.parse("2026-01-01T00:00:00Z")); diff --git a/java-sdk/sdk/build.gradle.kts b/java-sdk/sdk/build.gradle.kts index e30facde61e40..cf7d584650fae 100644 --- a/java-sdk/sdk/build.gradle.kts +++ b/java-sdk/sdk/build.gradle.kts @@ -461,6 +461,23 @@ abstract class GenerateDagDslTask : DefaultTask() { "`\"@once\"`, `\"@continuous\"`, a cron expression, or empty for no schedule.", ), ) + // The schema has no Dag-level queue; a Python Dag sets one through + // default_args instead. Every task of a Java Dag runs on the Java + // coordinator, so one queue on the Dag routes all of them there. + if (!dagProps.path("queue").isMissingNode) { + throw GradleException("The schema now has a Dag-level 'queue'; resolve it from the schema instead") + } + add( + DslField( + "queue", + "queue", + "STRING", + "String", + quote(""), + null, + "Queue each task of the Dag runs on, unless the task sets its own `queue`.", + ), + ) dagAllowlist.forEach { key -> val prop = dagProps.path(key) if (prop.isMissingNode) { 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 37015ca91978f..aadb2d706da4f 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 @@ -77,6 +77,9 @@ class DagDef( * mismatched value types are rejected on the call, so mistakes surface where * the Dag is defined. * + * `"queue"` is the one key Airflow's Dag has no setting for: it is the queue + * every task of the Dag runs on, unless the task sets its own `"queue"`. + * * @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 diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/TriggerDagRun.kt b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/TriggerDagRun.kt index c78d1a9fd5fc2..9165fd440fe5e 100644 --- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/TriggerDagRun.kt +++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/TriggerDagRun.kt @@ -66,6 +66,7 @@ private val TRIGGER_FIELDS: Map = * The task runs no Java code and takes no arguments. It pushes the triggered * run's ID, and the link the "Triggered DAG" extra link reads. It renders no * templates, so a value such as `"{{ ds }}"` reaches the new run unchanged. + * Like a Java task, it runs on the Dag's `queue` unless the task sets its own. * * @param dagId `trigger_dag_id`: the Dag to trigger. * @throws IllegalArgumentException if [dagId] is empty. diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt index e5ad4d6eacf6a..268afb866f2fd 100644 --- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt +++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt @@ -107,7 +107,10 @@ internal fun serializeDag( "relative_fileloc" to relativeFileloc, "timezone" to dagTimezone(dag.dagConfig), "timetable" to serializeTimetable(dag.id, dag.dagConfig), - "tasks" to dag.tasks.map { (taskId, def) -> serializeTask(taskId, def, downstream[taskId]) }, + "tasks" to + dag.tasks.map { (taskId, def) -> + serializeTask(taskId, def, downstream[taskId], dag.dagConfig["queue"] as String?) + }, "dag_dependencies" to serializeDagDependencies(dag), "task_group" to serializeTaskGroups(dag, expansion), "edge_info" to emptyMap(), @@ -122,11 +125,13 @@ internal fun serializeDag( /** * Converts one task to the Airflow serialization format. `downstream` is the * inverted view of the Dag's upstream edges, sorted for stable JSON. + * `dagQueue` is the Dag's queue, which the task takes unless it sets its own. */ private fun serializeTask( taskId: String, def: TaskDef, downstream: List?, + dagQueue: String?, ): Map { val data = linkedMapOf("task_id" to taskId) val trigger = def.trigger @@ -154,7 +159,11 @@ private fun serializeTask( // __type encoding is stripped. If core grows a task-level fill_config_defaults, // every SDK has to keep explicitly set values instead, or an explicit retries=0 // reads as unset and picks up the configured default. - def.configValues.forEach { (key, value) -> + // The Dag's queue is merged in first, so one equal to the schema default is + // left out too. A trigger task takes it as well, because the Java runtime runs it. + val config = + if (dagQueue == null || "queue" in def.configValues) def.configValues else def.configValues + ("queue" to dagQueue) + config.forEach { (key, value) -> if (key !in OMITTED_TASK_KEYS && !matchesSchemaDefault(SchemaFields.TASK[key], value)) { data[key] = unwrapTypeEncoding(serializeValue(value)) } 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 70bb24552e09d..5ed165580c94f 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 @@ -175,6 +175,7 @@ internal class DagDefTest { .config("dagrun_timeout", Duration.ofMinutes(5)) .config("start_date", OffsetDateTime.parse("2026-01-01T00:00:00Z")) .config("tags", listOf("a", "b")) + .config("queue", "java") Assertions.assertEquals( mapOf( @@ -185,6 +186,7 @@ internal class DagDefTest { "dagrun_timeout" to Duration.ofMinutes(5), "start_date" to OffsetDateTime.parse("2026-01-01T00:00:00Z"), "tags" to listOf("a", "b"), + "queue" to "java", ), dag.dagConfig, ) diff --git a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt index e025d9a5b154f..b813afd1079f5 100644 --- a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt +++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt @@ -22,10 +22,13 @@ package org.apache.airflow.sdk.execution import org.apache.airflow.sdk.Arg import org.apache.airflow.sdk.Bundle import org.apache.airflow.sdk.Client +import org.apache.airflow.sdk.ConditionTask import org.apache.airflow.sdk.Context import org.apache.airflow.sdk.DagDef +import org.apache.airflow.sdk.SwitchTask import org.apache.airflow.sdk.Task import org.apache.airflow.sdk.TaskDef +import org.apache.airflow.sdk.TriggerDagRun import org.apache.airflow.sdk.execution.comm.DagFileParseRequest import org.apache.airflow.sdk.internal.Refs import org.junit.jupiter.api.Assertions.assertEquals @@ -44,6 +47,20 @@ private class SerdeNoopTask : Task { ) = Unit } +private class SerdeCondition : ConditionTask { + override fun decide( + context: Context, + client: Client, + ) = true +} + +private class SerdeSwitch : SwitchTask { + override fun choose( + context: Context, + client: Client, + ) = SerdeNoopTask::class.java +} + @Suppress("UNCHECKED_CAST") private fun taskData( serialized: Map, @@ -54,6 +71,12 @@ private fun taskData( return tasks[index]["__var"] as Map } +@Suppress("UNCHECKED_CAST") +private fun queuesByTaskId(serialized: Map): Map = + (serialized["tasks"] as List>) + .map { it["__var"] as Map } + .associate { it["task_id"] as String to it["queue"] } + internal class SerdeTest { @Test @DisplayName("Should emit required dag fields and leave unset config-backed fields out") @@ -186,6 +209,48 @@ internal class SerdeTest { ) } + @Test + @DisplayName("Should give every task the Dag's queue unless the task sets its own") + fun shouldGiveEachTaskTheDagQueue() { + val dag = DagDef("d").config("queue", "java") + val extract = dag.task("extract", SerdeNoopTask::class.java) + dag.task("heavy", SerdeNoopTask::class.java).config("queue", "java_large") + dag.task("on_default", SerdeNoopTask::class.java).config("queue", "default") + dag.If("has_rows", SerdeCondition::class.java).after(extract) + dag.Switch("pick", SerdeSwitch::class.java).after(extract) + dag.task("trigger", TriggerDagRun("reports")) + dag.task("trigger_on_python", TriggerDagRun("reports")).config("queue", "python") + + val serialized = serializeDag(dag, "", ".") + + assertEquals( + mapOf( + "extract" to "java", + "heavy" to "java_large", + // Set on the task, so it wins; at the schema default, so it is left out. + "on_default" to null, + "has_rows" to "java", + "pick" to "java", + "trigger" to "java", + "trigger_on_python" to "python", + ), + queuesByTaskId(serialized), + ) + assertFalse("queue" in serialized) + } + + @Test + @DisplayName("Should leave the queue out when the Dag's queue is the schema default or unset") + fun shouldLeaveOutDefaultOrUnsetDagQueue() { + val onDefault = DagDef("d").config("queue", "default") + onDefault.task("t", SerdeNoopTask::class.java) + val unset = DagDef("d") + unset.task("t", SerdeNoopTask::class.java) + + assertFalse("queue" in taskData(serializeDag(onDefault, "", "."), 0)) + assertFalse("queue" in taskData(serializeDag(unset, "", "."), 0)) + } + @Test @DisplayName("Should serialize nested task groups with their own edges") fun shouldSerializeTaskGroups() {