Skip to content

Java SDK: answer TaskHandlerParseRequest with registered handlers - #74317

Draft
jason810496 wants to merge 7 commits into
jason/core-taskhandler-refactor/06-go-sdk-task-handler-parsefrom
jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse
Draft

jason810496 wants to merge 7 commits into
jason/core-taskhandler-refactor/06-go-sdk-task-handler-parsefrom
jason/core-taskhandler-refactor/08-java-sdk-task-handler-parse

Conversation

@jason810496

@jason810496 jason810496 commented Oct 6, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top), on native #74037 (stack #74170): #73973, #73974, #73975, #74317, #73976, #74067, #74135, #74140, #74141, #73971, #74030, #74031, #74032, #74136, #74137, #74138, #74139

A Java bundle now answers #73974's probe with every handler registered through @Builder.TaskHandler or Bundle.register(dagId, taskId, cls) and how its params bind, and JavaCoordinator can find and start the JAR a stub task would run. Until now a Java bundle failed on a TaskHandlerParseRequest as an unexpected first frame. Nothing calls the probe yet: #74135 does.

public static class ReportInput implements TaskInput {
  @ArgName("run_label") public String label;
  public long transformed;
}

@Builder.TaskHandler(dag = "etl", task = "transform")
public long transform(Client client, long extracted) { ... }

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

answers:

TaskHandlerParsingResult(
    fileloc="/bundles/java/etl.jar",
    task_handlers={"etl": [
        TaskHandlerDeclaration(task_id="transform", binding="positional", params=[
            TaskHandlerParam(name="extracted", value_schema={"type": "integer", "format": "int64"}),
        ]),
        TaskHandlerDeclaration(task_id="report", binding="named", params=[
            TaskHandlerParam(name="run_label", exact_name=True,
                             value_schema={"anyOf": [{"type": "string"}, {"type": "null"}]}),
            TaskHandlerParam(name="transformed", value_schema={"type": "integer", "format": "int64"}),
        ]),
    ]},
)
  • Only task handlers are declared. Tasks of a Dag declared in Java are left out, and Bundle.register(dagId, taskId, cls) now rejects a class generated for a @Builder.Dag task and points to @Builder.TaskHandler.
  • Flat data params are positional with their Java names. A TaskInput is named, with @ArgName fields matching exactly. A handler that reads no argument is named with empty params.
  • Java drops parameter names at compile time unless built with -parameters, so the annotation processor gives each generated @Builder.TaskHandler class an AIRFLOW_TASK_PARAMS field with its params' names and declared types.
  • value_schema describes the declared Java type in the vocabulary build_arg_bindings emits; only primitives exclude null. It does not widen for the decoder's lenient coercions, such as "5" into a long.
  • The answer is a plain map, because the generated models cannot hold task_handlers, and a test checks it against the vendored schema. The runtime waits for the parent to acknowledge it before exiting.
  • The vendored supervisor schema was refreshed only by downloading a published version, and versions are published only at release, so messages added during a release cycle never reached the Java models. :sdk:syncSupervisorSchema now copies the Task SDK snapshot when it declares the configured version. It runs only when asked and the build only checks the version, so a later Task SDK schema change neither rewrites the file nor fails ktlint. This refresh also brings in main's other schema changes since the last vendoring.
  • The probe checks the JAR a worker would run: _find_task_handler_artifact is the worker's own scan (main_class included) and returns the JAR that sets the Main-Class, and _build_parse_task_handler_command is a task's own command.
  • The example, e2e, Scala Spark and k8s bundles registered the Java tasks of Python Dags as a Dag declared in Java, which the answer never reports. They now register them as task handlers, and runs do not change.

Compatibility: a JAR built from main's Java SDK before this change reports 2026-10-30 but cannot answer, so its probe fails and #74135 logs a warning. Rebuilding fixes it.


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code (Opus 5.5) following the guidelines

The vendored schema.json was refreshed only when airflowSupervisorSchemaVersion changed, by downloading the published file. A version is published only at release (2026-10-30 returns 404), so messages added during a release cycle, such as TaskHandlerParseRequest, never reach the Java models. The Gradle task :sdk:syncSupervisorSchema now copies the Task SDK snapshot when it declares the configured version, and downloads the published file as before otherwise.

The copy runs only when asked: java-sdk/gradlew -p java-sdk :sdk:syncSupervisorSchema, then commit the file. Code generation reads only the committed schema.json, guarded by a read-only :sdk:checkSupervisorSchemaVersion that fails when the file is not at the configured version, so a later change to the snapshot never rewrites it from a build step or the ktlint hook. The prek hook stays a version-only check and runs the task only when the configured version changes.

The refreshed file adds the task handler messages and main's changes since the last vendoring (DagSourceCode, UpdateDagRunNote, retry_reason, multi_team, more required fields); the Java SDK builds and its tests pass on it unchanged.
The Dag processor will ask a Java bundle how each task handler binds its arguments, but Java parameter names do not survive compilation and the generated task class reads its arguments by position only. The annotation processor now gives every class it generates for a @Builder.TaskHandler method a static AIRFLOW_TASK_PARAMS field: each flat data parameter's name and declared type (a TypeRef when it has type arguments, primitives kept primitive), or the TaskInput type. @Builder.Dag task classes are unchanged, because the tasks of a Dag declared in Java are never task handlers.
The Dag processor checks each stub task argument against the JSON Schema of the handler param it binds to. buildValueSchema maps a declared Java type to the values it accepts, in the vocabulary @task.stub emits for Python annotations: only primitives exclude null, collections, arrays and maps describe their items, a POJO is an object, and types that decode from more than one JSON shape (byte[], Object, Optional, Date) say nothing. Enum constants come sorted, because Class.getFields has no defined order and the same JAR should answer the same on every JVM. It describes the declared type, not every value the decoder coerces, as pydantic does for its lax mode. Nothing calls it yet.
The Dag processor asks the runtime a stub task would run which task handlers it registers, and a Java bundle fails on that first frame as an unexpected message. Server now replies with one TaskHandlerParsingResult listing every handler registered with @Builder.TaskHandler or Bundle.register(dagId, taskId, cls), keyed by Dag id in registration order, so the answer depends only on the JAR. The tasks of a Dag declared in Java are never task handlers and are left out, and Bundle.register(dagId, taskId, cls) now rejects a class the annotation processor generated for a @Builder.Dag task, which carries a GeneratedDagTask marker, telling the author to declare a task handler with @Builder.TaskHandler. Flat data parameters bind positional with their names; a TaskInput, as the only data parameter or through InputTask, binds named, with @ArgName fields matched exactly; a task that reads no argument binds named with no params, so the Dag processor only warns about an argument passed to it. The reply is a plain map, because the generated models cannot hold task_handlers, and the SDK's tests check its keys against the vendored supervisor schema. The runtime waits for the parent's acknowledgement before it exits.
The example, the e2e test bundle, the Scala Spark example and the k8s example registered the Java tasks of Dags that a Python file owns as a Dag declared in Java. The task handler parse reports only task handler registrations, so the Dag processor would find no handler for those stub tasks, and once the Java SDK answers native Dag parsing each such Dag would be a second definition of the Python Dag. They now register those tasks as task handlers. Execution already falls back to task handlers, so runs do not change.
The probe from #73974 needs each coordinator to find the artifact a stub task runs and to start its runtime. JavaCoordinator finds the JAR with the scan that execute_task uses, main_class included, and returns the resolved path of the JAR that sets the Main-Class a task would run. The scan can stop on another JAR, such as the airflow-sdk JAR that gives a thin JAR its schema version, so _JarInfo now records that path. The probe runs a task's own command, so it answers for the classpath and Main-Class a worker would run.
The coordinator and SDK tests each fake one side, so nothing showed that a JAR built with the Gradle plugin answers the way the Dag processor parses it. The test builds java-sdk/example as a thin bundle from a copy of the SDK sources, finds it as a stub task of java_annotation_example would, and probes it. The airflow-sdk JAR sets the schema version and the example's JAR the Main-Class, so the found JAR is the latter. It needs a JDK, so it runs only with AIRFLOW_LANG_SDK_REAL_PROBE_TESTS=1.

This branch has not been deployed

No deployments
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.

1 participant