Skip to content

Java SDK: Honor TaskFlow arg bindings sent by the supervisor - #71188

Merged
jason810496 merged 6 commits into
apache:mainfrom
jason810496:feature/java-sdk-arg-bindings-runtime
Sep 25, 2026
Merged

jason810496 merged 6 commits into
apache:mainfrom
jason810496:feature/java-sdk-arg-bindings-runtime

Conversation

@jason810496

@jason810496 jason810496 commented Aug 5, 2026 •

Copy link
Copy Markdown
Member

Why

The Python @task.stub call site already defines task data flow. Java tasks should consume those bindings directly instead of repeating upstream IDs with @Builder.XCom.

The Java call site is settled by ADR-0001 (merged in #72019); this PR implements its binding half. The wire contract is ADR-0007.

Example

Three authoring syntaxes, exactly the three ADR-0001 specifies. All resolve to the same bindings.

1. Annotation based with positional injection

@Builder.TaskHandler(dag = "etl", task = "score")
public long score(Client client, long rows, double threshold, List<String> regions) { ... }

2. Annotation based with an explicit struct

report(run_label=..., transformed=...) binds with nothing declared — names match ignoring case and underscores:

public static class ReportInput implements TaskInput {
  public String runLabel; // binds run_label through the fold
  public long transformed;
}

@Builder.TaskHandler(dag = "etl", task = "report")
public void report(ReportInput input) { ... }

3. Interface based with an explicit TaskInput

public static class SummarizeInput implements TaskInput {
  // Pinned so the field can be called region. Or drop it: public String regionCode;
  @ArgName("region_code")
  public String region;

  public long transformed;
}

public static class Summarize implements InputTask<SummarizeInput> {
  public void execute(@NotNull Context context, Client client, SummarizeInput input) { ... }
}

Syntax 1 compiles to this — copied out of the example project's generated sources:

public static final class ConsumeDoubleList implements Task {
  @Override
  public void execute(Context context, Client client) throws Exception {
    TaskArgs args = TaskArgs.of(context, client);
    List<Double> values = args.get(0, new TypeRef<List<Double>>() {});
    new XComCastingExample().consumeDoubleList(values);
  }
}

How

  • @Builder.TaskHandler(dag, task) names the pair a handler binds to, rather than @Builder.Task(id) inside a @Builder.Dag class. Python declares the task and Java supplies only its body, so the two are different things and should not share a name. @Builder.Dag/@Builder.Task are left for a Dag authored in Java (Java SDK: Declare a Dag's task graph with @Builder.Deps #71189).
  • Bundle is created and registered into: register takes the class a handler lives on.
  • Data parameters bind by position; injected Client/Context take no position.
  • Keyword arguments bind by name through one TaskInput. A task declares flat parameters or one TaskInput, never both.
  • Names match ignoring case and underscores — the same fold Go and TS use, so one Python signature binds identically in all three SDKs with nothing declared. @ArgName is for what the fold cannot reach (a Python keyword, an unusable identifier) and is matched literally.
  • Declared type arguments survive the decode, so a List<String> parameter or field does not arrive as List<LinkedHashMap>.
  • TaskArgs is internal — org.apache.airflow.sdk.internal, not a TaskInput, emitted by the processor only. InputTask<TaskArgs> is a compile error; interface tasks declare a TaskInput.
  • Missing reference or boxed values become null; primitive inputs fail with MissingXComException, naming the stub argument.
  • @Builder.XCom is removed, keeping the Python call site as the single source of data-flow wiring.

Follow-up

ADR-0001 writes the generic read as Jackson's TypeReference<T>. This PR ships an SDK-owned TypeRef<T> instead, because Jackson is an implementation dependency and exposing it would put Jackson on every consumer's compile classpath. Promoting jackson-core to api and using TypeReference directly is worth doing on its own — it is a deliberate next step, not a limitation this PR works around.

ADR-0001's examples are also over-annotated under the fold: @ArgName("run_label") on a runLabel field is now redundant. Worth a pass over the ADR.


Was generative AI tooling used to co-author this PR?
  • Yes, with help of Claude Code Opus 5 and Codex (GPT-5) following the guidelines

@FrankYang0529 FrankYang0529 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Overall LGTM. Thanks for the PR.

Comment thread java-sdk/README.md
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt Outdated

@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, I just addressed the comments.

Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt Outdated
Comment thread java-sdk/README.md
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/ArgBinding.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt Outdated

@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.

Thanks, just leave few concerns

A declared parameter must find an argument; a passed argument need not find
a parameter. The two directions are not symmetric, and the SDK now treats
them that way.

- `TaskInput`: a field no argument supplies fails the task, whatever its
  type. Previously only a primitive field failed and a boxed or reference
  one silently took `null`, which hid a Java signature that had drifted from
  the Python stub. Once a field has claimed its argument, a null literal or
  an absent XCom still gives `null` to a boxed or reference field and still
  fails for a primitive, so declaring a boxed type stays the way to say the
  value is optional. A field two arguments fold onto now names the ambiguity
  instead of failing as if nothing matched.
- `TaskInput`: an argument no field claims is logged rather than failed. The
  Go SDK errors here, but it has to, because its unmatched fields are
  silently zero-valued; Java's primitive and boxed types already say per
  field whether a value is required.
- Positional binding: read `from_default` and drop those entries before the
  arity check when the counts disagree. ADR-0007 captures an unpassed
  parameter with its default and expects consumers to leave it unclaimed, so
  a method that omits a defaulted trailing parameter is not reading shifted
  arguments.

ADR-0001 records the rule and why it differs from go-sdk ADR-0006.
@jason810496
jason810496 force-pushed the feature/java-sdk-arg-bindings-runtime branch from e8b8a5e to fd97648 Compare September 23, 2026 12:25
@jason810496
jason810496 marked this pull request as ready for review September 23, 2026 14:55

@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/ArgBinding.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt Outdated
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt
@henry3260

henry3260 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Thanks for solving my comments :)

…istration at serve

A TaskInput field binds by name, so neither direction of a name mismatch
can hand a task the wrong value: a field nothing supplies keeps its Java
default, and an argument no field claims changes nothing the task reads.
Failing the unfilled direction, as this PR did until now, made Java stricter
than the Go SDK for a disagreement neither language can act on. Both
directions are logged instead, with the wording the Go SDK uses. Two fields
whose names fold alike still fail when the bundle is built, since the fold
cannot tell them apart and neither value is safe to pick.

Bundle is mutable now that registration is incremental, which left two ways
to build one that cannot work:

- A Dag declared in Java owns its own tasks, so its ID cannot also hold task
  handlers. The two kinds live in their own maps and each rejects an ID the
  other holds, whichever order they arrive in. Registering a handler for a
  Dag the Java side declared used to succeed and quietly add a task no
  Python file could supply arguments for.
- Registration ends when Server.serve starts, so a register left below it is
  reported rather than racing the running task.
@jason810496
jason810496 force-pushed the feature/java-sdk-arg-bindings-runtime branch from 9afb19f to ef4bb3c Compare September 24, 2026 08:38
@jason810496
jason810496 marked this pull request as ready for review September 24, 2026 08:41

@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.

Nice catch, both concerns are valid. TS and Go side did validate your findings.
I addressed them all in the latest commit, thanks.

Additionally, the latest commit filled the gap of mixed language task TaskFlow binding semantic like #73648.

Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt
Comment thread java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Bundle.kt

@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.

Thanks for the update!

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.

3 participants