Skip to content

Java SDK: Serialize native Dags to DagSerialization v3 - #71190

Merged
jason810496 merged 13 commits into
apache:mainfrom
jason810496:feature/java-sdk-dag-serialization
Oct 8, 2026
Merged

jason810496 merged 13 commits into
apache:mainfrom
jason810496:feature/java-sdk-dag-serialization

Conversation

@jason810496

@jason810496 jason810496 commented Aug 5, 2026 •

Copy link
Copy Markdown
Member

Merge order:

  1. Java SDK: Group a native Dag's tasks with task groups #74230 — Group a native Dag's tasks with task groups
  2. Java SDK: Serialize native Dags to DagSerialization v3 #71190 — Serialize native Dags to DagSerialization v3 (current one)
  3. Java SDK: Pack each native Dag's source file into the bundle JAR #74096 — Pack each native Dag's source file into the bundle JAR

Why

Make the runtime can answer the coordinator's DagFileParseRequest itself.

How

  • parseDags serializes every Dag registered on the Bundle to DagSerialization v3.
  • Server.dispatchTask gains a DagFileParseRequest branch alongside StartupDetails, so a bundle process serves either a task run or a parse request.
  • A field the Dag leaves unset stays out, including the ones Python reads from [core] config, such as max_active_tasks and catchup. Airflow fills those in from its own config when it receives the Dag.
  • Task groups serialize into the nested task_group object Python writes, with each group carrying its own edges.
  • Every task carries is_stub, and one the Dag called with arguments also carries the _arg_bindings spec ADR-0007 defines, so the serialized Dag records what each task is called with.
  • A Java task still reads its arguments from the Dag in its own bundle rather than from that spec: the bundle that runs the task also holds the call that wired it, so nothing has to travel at run time.
  • A cron preset expands before it is written, and the timezone comes from the Dag's start date, both as Python does. A schedule that is neither a preset nor a cron expression is rejected where the Dag is written.
  • A Dag that cannot be serialized becomes an import error keyed by the bundle-relative path, so the other Dags in the same JAR still parse.

Example

load(transform(extracted, lit(1.5))) in a wiring class serializes as:

"_arg_bindings": [
  {"name": "extracted", "kind": "xcom", "task_id": "extract"},
  {"name": "factor", "kind": "literal", "value": 1.5}
]

The spec travels as JSON, so lit(...) takes a string, number, boolean, list or map; anything else is rejected when the Dag is parsed.

Cross-language validation

check-java-sdk-serialization-conformance builds the shared test Dags of scripts/ci/lang_sdk_serialization/test_dags.yaml with both this SDK and Airflow's own serializer, loads the Java output through DagSerialization.validate_schema and from_dict, and compares the two field by field.

Known limitation

Cron schedules map to CronTriggerTimetable only, as the TypeScript SDK does, until the supervisor forwards the [scheduler] timetable flags over the coordinator protocol. There is a TODO at the timetable serializer.


Was generative AI tooling used to co-author this PR?

@jason810496
jason810496 force-pushed the feature/java-sdk-dag-serialization branch 2 times, most recently from efba9d2 to c95e521 Compare October 6, 2026 06:03
@jason810496
jason810496 force-pushed the feature/java-sdk-dag-serialization branch 2 times, most recently from 2fbc57b to 50f810b Compare October 7, 2026 07:44
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt Outdated
Comment thread java-sdk/scripts/ci/prek/check_serialization_conformance.py Outdated
Comment thread .pre-commit-config.yaml

@henry3260 henry3260 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Just want to help unblock the technical side here, not a formal review.

@jason810496
jason810496 force-pushed the feature/java-sdk-dag-serialization branch from 50f810b to b8938f8 Compare October 7, 2026 17:22

@jason810496 jason810496 left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the review.

Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt Outdated
Comment thread java-sdk/scripts/ci/prek/check_serialization_conformance.py Outdated
Comment thread .pre-commit-config.yaml
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).
…steers it

`Bundle.register` expands the group edges, `Refs` decides which group each
task lands in, and `sdk/build.gradle.kts` generates the schema fields the
serializer writes from. All three change the serialized output, so add
them to the hook's trigger set.

Report a failed classpath build with the Gradle output instead of an
`IndexError` on an empty line list or a bare `CalledProcessError`.
Expand a cron preset before writing it, so `@daily` serializes as
`0 0 * * *` and the Dag hashes the same as the Python one. Reject a
schedule that is neither a preset nor a cron expression where the Dag is
written, rather than letting it reach a scheduler that cannot build a
timetable from it.

Take the timezone from the Dag's start date, as Python's `DAG.timezone`
does, instead of always writing UTC: a Dag started at a non-UTC offset
now fires at that offset rather than at the same wall clock in UTC.
A native Dag's serialized form recorded its task edges but nothing about
what each task was called with, so changing a literal argument left the
Dag byte-identical: no new version, and nothing in the version diff.

Write the same `is_stub` flag and `_arg_bindings` spec the TypeScript SDK
writes, as ADR-0007 decision G says a natively authored Dag should. The
wiring view now passes its parameter names along with the arguments, so
each binding can name the parameter it feeds.

A Java task still resolves its arguments from the Dag in its own bundle
rather than reading the spec back: the bundle that runs the task also
holds the call that wired it, so the values never have to travel, and a
`TaskInput` keeps binding as one whole input. The spec is what Airflow
records and shows.

Since the spec travels as JSON, a literal with no JSON form is now
rejected where the Dag is parsed instead of reaching a task that cannot
receive it.
…arse

A Dag that could not be serialized threw out of the parse response, so
every other Dag in the same JAR vanished with it and `import_errors` was
never filled. Report the failure under the bundle-relative path the way
the TypeScript SDK does, naming the Dag, and serialize the rest.
`native-dag-authoring`, `task-args`, `taskflow-dependencies` and
`task-group` are new in this release, not in 3.3, which is what the
matrix's "supported since" column says.
…uts first

The cron shape check now accepts comma lists of names, month and weekday
names, the `#`, `L` and `W` qualifiers and the `@midnight` and `@annually`
aliases, which Python stores unexpanded. Each of these previously turned
the Dag into an import error. The schema-default comment records why
explicit default values are dropped, and the timezone TODO and the Java
docs note that a cron Dag with no start date is scheduled in UTC.

A wired input wins over a runtime binding in ArgValues. The Java docs and
ADR 0007 now say so, and a test pins it.

The conformance run now passes `--supports literal_inputs` and the Java
serializer wires each task's upstream handles and literals as call
arguments, so the run exercises `_arg_bindings`. The hook's file filter
covers `internal/Fields.kt`. Selective checks skip the hook unless a
java-sdk file, Airflow's serializer or schema, or the shared harness
changed, through a new file group because serializer changes do not force
full tests.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants