Skip to content

Commit 196ed3a

Browse files
committed
Java SDK: Serialize native Dags to DagSerialization v3 on parse requests
A Java-authored Dag could not reach the scheduler on its own: the runtime answered task-execution requests only, so a Python stub file still had to exist purely to describe the Dag's structure. With dependency edges and schema-keyed configuration now recorded on the Java Dag model, the runtime has everything it needs to answer the coordinator's DagFileParseRequest the same way the Go SDK does, and the nativedag examples become real schedulable Dags with no Python counterpart. Native Java tasks deliberately emit no _arg_bindings: the execution API delivers bindings only for Python _StubOperator tasks, and a Java task always runs inside the JVM bundle that already holds its wired inputs, so the runtime resolves them locally. Cron schedules map to CronTriggerTimetable only, mirroring the Go SDK until the supervisor forwards the [scheduler] timetable flags over the coordinator protocol (see the TODO at the timetable serializer). (cherry picked from commit 651791a)
1 parent 65b09de commit 196ed3a

7 files changed

Lines changed: 645 additions & 28 deletions

File tree

‎java-sdk/README.md‎

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -619,17 +619,17 @@ prek hook regenerate it.
619619
| capability: `asset-event-emit` | MAY | ✗ | – | runtime does not emit asset events yet |
620620
| capability: `asset-event-read` | MAY | ✗ | – | no task-facing asset-event API yet |
621621
| **Native-Dag authoring** | | | | |
622-
| capability: `native-dag-authoring` | SHOULD | ✗ | – | native Dag authoring not implemented yet |
623-
| capability: `task-args` | MUST † | n/a | – | |
624-
| capability: `dag-params` | MUST † | n/a | – | |
625-
| capability: `taskflow-dependencies` | MUST † | n/a | – | |
626-
| capability: `branching` | SHOULD † | n/a | – | |
627-
| capability: `dag-test` | SHOULD † | n/a | – | |
628-
| capability: `task-group` | MAY † | n/a | – | |
629-
| capability: `dynamic-task-mapping` | MAY † | n/a | – | |
630-
| capability: `asset-inlets-outlets` | MAY † | n/a | – | |
631-
| capability: `asset-scheduling` | MAY † | n/a | – | |
632-
| capability: `object-store` | MAY † | n/a | – | |
622+
| capability: `native-dag-authoring` | SHOULD | ✓ | 3.3 | |
623+
| capability: `task-args` | MUST † | ✓ | 3.3 | |
624+
| capability: `dag-params` | MUST † | ✗ | – | |
625+
| capability: `taskflow-dependencies` | MUST † | ✓ | 3.3 | |
626+
| capability: `branching` | SHOULD † | ✗ | – | |
627+
| capability: `dag-test` | SHOULD † | ✗ | – | |
628+
| capability: `task-group` | MAY † | ✗ | – | |
629+
| capability: `dynamic-task-mapping` | MAY † | ✗ | – | |
630+
| capability: `asset-inlets-outlets` | MAY † | ✗ | – | |
631+
| capability: `asset-scheduling` | MAY † | ✗ | – | |
632+
| capability: `object-store` | MAY † | ✗ | – | |
633633

634634
*Marks: ✓ supported · ✗ not supported · n/a not applicable. A tier marked † applies only when `native-dag-authoring` is supported.*
635635

‎java-sdk/capabilities.yaml‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -52,8 +52,8 @@ states:
5252
supported: true
5353
since: "3.3"
5454

55-
# Runtime capabilities reflect the task-facing Client surface; native-Dag authoring is not
56-
# implemented yet, so every native capability is unsupported.
55+
# Runtime capabilities reflect the task-facing Client surface. Native-Dag capabilities describe
56+
# what a Java-authored Dag can declare and deliver through Dag serialization.
5757
capabilities:
5858
mixed-lang-stub-target:
5959
supported: true
@@ -95,14 +95,16 @@ capabilities:
9595
supported: false
9696
note: "no task-facing asset-event API yet"
9797
native-dag-authoring:
98-
supported: false
99-
note: "native Dag authoring not implemented yet"
98+
supported: true
99+
since: "3.3"
100100
task-args:
101-
supported: false
101+
supported: true
102+
since: "3.3"
102103
dag-params:
103104
supported: false
104105
taskflow-dependencies:
105-
supported: false
106+
supported: true
107+
since: "3.3"
106108
branching:
107109
supported: false
108110
dag-test:

‎java-sdk/sdk/module.md‎

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -56,17 +56,17 @@ meaning of each dimension is defined in the
5656
| capability: `asset-event-emit` | MAY | ✗ | – | runtime does not emit asset events yet |
5757
| capability: `asset-event-read` | MAY | ✗ | – | no task-facing asset-event API yet |
5858
| **Native-Dag authoring** | | | | |
59-
| capability: `native-dag-authoring` | SHOULD | ✗ | – | native Dag authoring not implemented yet |
60-
| capability: `task-args` | MUST † | n/a | – | |
61-
| capability: `dag-params` | MUST † | n/a | – | |
62-
| capability: `taskflow-dependencies` | MUST † | n/a | – | |
63-
| capability: `branching` | SHOULD † | n/a | – | |
64-
| capability: `dag-test` | SHOULD † | n/a | – | |
65-
| capability: `task-group` | MAY † | n/a | – | |
66-
| capability: `dynamic-task-mapping` | MAY † | n/a | – | |
67-
| capability: `asset-inlets-outlets` | MAY † | n/a | – | |
68-
| capability: `asset-scheduling` | MAY † | n/a | – | |
69-
| capability: `object-store` | MAY † | n/a | – | |
59+
| capability: `native-dag-authoring` | SHOULD | ✓ | 3.3 | |
60+
| capability: `task-args` | MUST † | ✓ | 3.3 | |
61+
| capability: `dag-params` | MUST † | ✗ | – | |
62+
| capability: `taskflow-dependencies` | MUST † | ✓ | 3.3 | |
63+
| capability: `branching` | SHOULD † | ✗ | – | |
64+
| capability: `dag-test` | SHOULD † | ✗ | – | |
65+
| capability: `task-group` | MAY † | ✗ | – | |
66+
| capability: `dynamic-task-mapping` | MAY † | ✗ | – | |
67+
| capability: `asset-inlets-outlets` | MAY † | ✗ | – | |
68+
| capability: `asset-scheduling` | MAY † | ✗ | – | |
69+
| capability: `object-store` | MAY † | ✗ | – | |
7070

7171
*Marks: ✓ supported · ✗ not supported · n/a not applicable. A tier marked † applies only when `native-dag-authoring` is supported.*
7272

‎java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,8 +33,10 @@ import kotlinx.coroutines.runBlocking
3333
import org.apache.airflow.sdk.execution.CoordinatorComm
3434
import org.apache.airflow.sdk.execution.LogSender
3535
import org.apache.airflow.sdk.execution.Logger
36+
import org.apache.airflow.sdk.execution.comm.DagFileParseRequest
3637
import org.apache.airflow.sdk.execution.comm.ErrorResponse
3738
import org.apache.airflow.sdk.execution.comm.StartupDetails
39+
import org.apache.airflow.sdk.execution.parseDags
3840
import org.apache.airflow.sdk.execution.runTask
3941
import kotlin.text.substringAfterLast
4042
import kotlin.text.substringBeforeLast
@@ -181,6 +183,7 @@ class Server(
181183
val frame = coordinator.readMessage()
182184
when (val body = frame.body) {
183185
is StartupDetails -> runTaskAndReport(bundle, body, coordinator)
186+
is DagFileParseRequest -> parseDagsAndReport(bundle, body, coordinator)
184187
is ErrorResponse -> throw ApiError("[${body.error}] ${body.detail}")
185188
else -> throw ApiError("Unexpected initial frame (id=${frame.id})")
186189
}
@@ -194,4 +197,12 @@ class Server(
194197
val result = runTask(bundle, startup, coordinator)
195198
coordinator.communicate<Unit>(result)
196199
}
200+
201+
private suspend fun parseDagsAndReport(
202+
bundle: Bundle,
203+
request: DagFileParseRequest,
204+
coordinator: CoordinatorComm,
205+
) {
206+
coordinator.communicate<Unit>(parseDags(bundle, request))
207+
}
197208
}

0 commit comments

Comments
 (0)