Repository navigation
Add task state store API to the Java SDK - #73464
Conversation
|
Hello @Andrushika - thank you for your contributions to Apache Airflow! The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 7 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:
These pull requests stay open because maintainers are already engaged in them - they count towards your limit:
This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are. What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack. Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst. Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting |
249ca7e to
1b0ca65
Compare
FrankYang0529
left a comment
There was a problem hiding this comment.
Thanks for the PR. Leave one comment.
1b0ca65 to
4b9c632
Compare
jason810496
left a comment
There was a problem hiding this comment.
Nice, thanks! LGTM overall.
Let's also add a TODO before the serialization to respect the state_store/max_value_storage_bytes that will be passed from the coordinator side.
Airflow 3.3 added a task state store (AIP-103) and the supervisor already handles the Get/Set/Delete/ClearTaskStateStore messages, but the Java SDK had no task-facing API for it and capabilities.yaml listed it as unsupported. Java tasks could not keep state such as an external job ID across retries.
Review found that a @JvmField on Client cannot be stubbed by Mockito, so a Java unit test of a task that reads the store hit a NullPointerException. The docs also said entries survive later runs, but the store is scoped to one task instance, so they only survive retries within the same Dag run.
A key stored without a retention never expired, while Python applies [state_store] default_retention_days. The coordinator passes that setting to the JVM as AIRFLOW__STATE_STORE__DEFAULT_RETENTION_DAYS, so read it the same way the Go SDK does: absent means 30 days, 0 means never expire, and a malformed value fails instead of retaining for a different period. TaskStateStore.NEVER_EXPIRE replaces the old "no retention" meaning, and a zero or negative retention is rejected.
4b9c632 to
f775faa
Compare
Review on the Go and TS task state store PRs asked to drop the hardcoded 30-day fallback: the coordinator always passes AIRFLOW__STATE_STORE__DEFAULT_RETENTION_DAYS, so a missing variable means the JVM was not launched by the coordinator, and guessing a period there would drift from the deployment's config. Also leave a TODO for the max_value_storage_bytes warning the Python accessor emits.
|
Also applied the two points from #73420 and #73618 that apply here:
One note on ordering: the coordinator only passes that variable once #73420 is merged. If this PR goes in first, |
Would it be better to add the explicit |
|
Yes, we already have
The Sorry for the misleading wording :( |
…-store # Conflicts: # airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
Why
Airflow 3.3 added the task state store (AIP-103). A task can save a value under a key, and the value survives retries.
The supervisor already handles the four
*TaskStateStoremessages and the Java SDKschema.jsonalready has them. But nothing on the Kotlin side sends them. So Java tasks cannot use the store andcapabilities.yamlliststask-state-storeas unsupported.What
Add
client.getTaskStateStore()withget,set,deleteandclear. It is wired the same way as variables.getreturnsnullwhen the key is not there.settakes an optionaljava.time.Durationretention, orTaskStateStore.NEVER_EXPIRE. Without a retention it uses[state_store] default_retention_days, read fromAIRFLOW__STATE_STORE__DEFAULT_RETENTION_DAYSthe same way as the Go SDK in #73420 (0 means never expire, a missing variable fails).SetTaskStateStoreis added toREQUIRED_NULLABLE_REQUESTS. Without it Jackson drops a nullexpires_atand the supervisor rejects the message. I checked this with the supervisor decoder.One difference from Python is noted in
java.rst:[workers] state_store_backendis not supported. The backend hooks run inside the Python task process today, so supporting it in the Java SDK needs a further design discussion.related: #73420
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Fable 5.1) following the guidelines