Repository navigation
[DO NOT MERGE] Mirror of Java SDK #73596, #71189, #71190 for the native Dag stack - #73846
Closed
jason810496 wants to merge 15 commits into
Closed
jason810496 wants to merge 15 commits into
jason810496 wants to merge 15 commits into
Conversation
jason810496
added this pull request to stack #73848
September 28, 2026 13:20
This was referenced Sep 28, 2026
1 task done
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
September 28, 2026 15:57
a943bc2 to
daaccdf
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
September 28, 2026 15:57
888dc77 to
4a7d46e
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
September 29, 2026 01:25
daaccdf to
2547205
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
September 29, 2026 01:26
4a7d46e to
20dc5fd
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
September 29, 2026 01:29
2547205 to
39c93a6
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
2 times, most recently
from
September 29, 2026 13:18
5c237c0 to
32456fb
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
September 29, 2026 13:18
39c93a6 to
196ed3a
Compare
jason810496
removed this pull request from stack #73848
October 1, 2026 03:02
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
October 1, 2026 03:03
32456fb to
1d9eee4
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
October 1, 2026 03:03
196ed3a to
3d6e6a5
Compare
1 task done
jason810496
added this pull request to stack #74005
October 1, 2026 03:04
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
October 1, 2026 06:20
1d9eee4 to
e5d7124
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
2 times, most recently
from
October 1, 2026 12:41
e2459c7 to
41a2b06
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
October 1, 2026 12:41
e5d7124 to
56cc68c
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
October 5, 2026 09:12
37c2458 to
981651c
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
October 5, 2026 10:42
9ca9a23 to
20cece7
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
2 times, most recently
from
October 5, 2026 11:39
587f720 to
e876210
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
October 6, 2026 14:10
698fae9 to
c657459
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
October 6, 2026 14:10
e876210 to
5d91e38
Compare
jason810496
force-pushed
the
jason/lang-sdk-e2e/06-ts-sdk
branch
from
October 6, 2026 15:13
c657459 to
0480246
Compare
jason810496
removed this pull request from stack #74170
October 7, 2026 07:03
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
October 7, 2026 07:10
5d91e38 to
cba80f3
Compare
jason810496
added this pull request to stack #74395
October 7, 2026 07:27
…request DagDef remembers the outermost class that built it, and a generated @Builder.Dag builder names the annotated class instead of itself. Server accepts --describe-sources <file> to write each Java-declared Dag's declaring class and return without connecting, for the Gradle plugin. (cherry picked from commit 9eeda41)
The packDagSources task runs mainClass in describe mode, finds each Dag's source file through the class file's SourceFile attribute, and packs each unique file with a sources.json mapping under META-INF/airflow/. The JAR manifest names it with Airflow-Java-SDK-Sources. (cherry picked from commit 030a070)
Every `Jar` task gained a dependency on `packDagSources`, so a sources or javadoc JAR carried the Dag sources and claimed the `Airflow-Java-SDK-Sources` attribute that the Python side picks a bundle JAR out of a directory by. Only `jar` and `shadowJar` carry them now. The describe run had no time bound, so a `main` that does not route through `Server.create(args)` and starts serving instead held the build open with no way out. It now runs under a timeout and falls back to packing the entrypoint's source alone, which is what every other describe failure already does. Say in `serve` and `serveAsync` that a `--describe-sources` server returns without connecting, and record in `sources.json` that its paths are relative to the source directory, which is also where they sit under `META-INF/airflow/sources/`. (cherry picked from commit d36ca2c)
…layout The Java importer now passes the Dag ID to the source lookup, so the Code view reads dag_source_paths instead of always showing the entrypoint. The source locator falls back to a unique match on the SourceFile name when the package does not match the directory, using an index built only when the package path misses. An unresolvable entrypoint or Dag class is now logged as a warning. The describe task is no longer up to date after a failed run, so a transient failure does not stick until a source changes. Its hint now points at Server.create(args), which is what answers --describe-sources. Add a shadowJar test, correct the comment on which JARs carry the sources payload, and document that a Dag maps to the class that constructed it. (cherry picked from commit 10a6954)
A native Java Dag could only run every task it declared: it had no way to let a task's result decide what runs next, so a Dag that loads only when there are rows had to be split in two or fall back to a Python stub. `DagDef.If` and `@Builder.If` declare a task whose boolean names the side that runs; the other is skipped. The deciding task carries `_can_skip_downstream`, so Airflow skips a cleared side again rather than running it, and the SDK writes the `skipmixin_key` XCom before the skip, as Python's SkipMixin does. It skips only the side not taken, so a task after both sides needs a trigger rule such as `none_failed_min_one_success`. `If` is capitalized because `if` is a Java keyword and `@Builder.If` reads the same way. With no task id given, one is derived from the condition class's simple name, so `HasRows.class` becomes `hasRows`. The runtime reaches `SkipDownstreamTasks` through the transport client only, never through the task-facing `Client`, so a task body cannot skip tasks the Dag did not declare it may. (cherry picked from commit 30e5218)
A condition picks between two tasks, but a Dag that routes work by size or by region needs more than two, and splitting that into nested conditions makes the graph unreadable. `DagDef.Branch` and `@Builder.Branch` declare a task that chooses one of the tasks listed with `option(...)`; every other one is skipped. The chosen task's ID becomes the branch's return value, as it does for Python's branch operators. A branch chooses exactly one case, so a task that runs after several of them needs a trigger rule such as `none_failed_min_one_success`. Each authoring style names the case in the way that style can type-check. A `BranchTask` returns the case's own class, so javac checks it; two cases that share a class are rejected where the second is named, since the branch could not tell them apart. A `@Builder.Branch` method has no such class to name -- a task method's implementation is generated -- so it returns a `TaskId` from the generated `<Dag>Builder.TaskIds`, which holds one constant per task of the Dag. Either way nothing hard-codes a task id string, and the generated holder is emitted only for a Dag that declares a branch. (cherry picked from commit bc81537)
A native Java Dag could not start another Dag's run. Airflow's
TriggerDagRunOperator is Python, and a Java bundle has no Python Dag file to
declare one in, so a Java Dag that fans out to another Dag needed a Python
counterpart for that one task.
A `TriggerDagRun` passed to `dag.task` declares a task that has no Java body
to run; the Java runtime runs it as `TriggerDagRunOperator` does. Its
settings are named as Python names the operator's parameters:
```java
dag.task("trigger_downstream",
new TriggerDagRun("downstream_etl")
.config("wait_for_completion", true)
.config("conf", Map.of("rows", 2)));
```
An annotated Dag declares one from an ordinary `@Builder.Task` method that
returns a `TriggerDagRun`. The serialized Dag has to carry what is
triggered, so that method runs when the Dag is built rather than when the
task runs, and therefore takes no parameters.
The task serializes as Python's own `TriggerDagRunOperator`: its template
fields, its UI colour and extra link, and a `dag_dependencies` entry, so the
Airflow UI draws the dependency between the two Dags. Python writes only an
operator's template fields and keeps the rest in the Dag file; a Java Dag
has no such file, so every setting the author made is written too.
The runtime records the triggered run's link and ID as XComs, skips or fails
on a run that already exists, and waits for completion either by polling or,
with `deferrable`, by deferring to the Python triggerer's DagStateTrigger and
resuming from the event it fires. That makes this the first Java runtime to
emit DeferTask and a skipped TaskState.
`scripts/ci/lang_sdk_serialization/test_dags.yaml` now gives the shared
trigger Dag a real logical date instead of a Jinja template, because the SDK
takes one as a date-time. The Go SDK is making the same change in #74417;
whichever lands second drops that one line, and until then the Go
serialization conformance check reports that single difference.
(cherry picked from commit af569b2)
A TriggerDagRun task triggers a Dag run, reads its state and checks whether the Dag is paused. The runtime has no client for these three requests yet. An existing run_id comes back as a result instead of an error, because the task decides whether that skips or fails it. (cherry picked from commit 41be38c)
The Go runtime runs a TriggerDagRun task itself and renders no templates, so a string cannot hold a date it must compute. A time.Time lets the runtime send the date as it is, and the zero value keeps the Python default of now. The conformance case and the Go fixture move to a literal datetime. (cherry picked from commit a5ae692)
A task from airflow.TriggerDagRun had no Task in the bundle, so a Go Dag needed a Python worker to run it. The bundle now returns a TriggerTask, and RunTask runs it the way TriggerDagRunOperator does: the same run_id and date rules, the link and trigger_run_id XComs, polling, and a deferral to DagStateTrigger that resumes in the same runtime. It reads its config from the environment and renders no templates. (cherry picked from commit c347c12)
The Go runtime runs a TriggerDagRun task now, so it belongs on the coordinator that the Dag queue maps to. Before, the task kept no queue and went to a Python worker. A queue in the TaskSpec still wins, as for any other task, and a Dag with no queue leaves the task without one. (cherry picked from commit 4d025be)
Say how the task names its run, what Deferrable needs, and that the arguments are not templated and OpenLineage parent info is not added. (cherry picked from commit d73904c)
Shows Dags built with airflow.Dag, with task groups, If, Switch and TriggerDagRun, in a binary the Dag processor parses. The e2e test packs and runs it. (cherry picked from commit d5c4ce3)
Pack the native example into the Dags folder so the Dag processor runs it to parse, and give that service the worker's USER and HOME. Keep the handler-only bundle out of parsing with an .airflowignore. (cherry picked from commit 0de3491)
Parse, run and inspect the example native Dags through the stock stack. The TriggerDagRun case is a strict xfail until the Go runtime has a trigger runner. (cherry picked from commit 727e256)
jason810496
force-pushed
the
jason/lang-sdk-e2e/07-java-sdk
branch
from
October 8, 2026 06:34
cba80f3 to
01dee95
Compare
jason810496
removed this pull request from stack #74395
October 9, 2026 08:45
Member
Author
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Stack (bottom to top): #74042, #74035, #74043, #74036, #74037, #73845, #73846, #73847
Part of the native Dag e2e stack. Builds on the TS SDK mirror layer.
Replays the Java SDK PRs #71189, #71190 and #74096 at their current heads, in that order (one commit each for #71189 and #71190, two for #74096), so the native e2e layer can test native Java Dags end to end. No new behavior. Please review on the original PRs, not here.
main.ArgTestSupport.ktandArgValuesTest.kt, the test helpers wire the task withRefs.recordand settaskDefon the context, as Java SDK: Resolve a task's arguments from the Dag's own wiring #73596 does now that it is merged.Was generative AI tooling used to co-author this PR?