kaxil opened a new pull request, #74421:
URL: https://github.com/apache/airflow/pull/74421
When `[workers] state_store_backend` is set, the Task SDK wraps every value
returned by `serialize_*_state_store_to_ref()` in `{"__airflow_state_ref__":
...}`
([`context.py`](https://github.com/apache/airflow/blob/3e302fd257a0cfe8deba2de529bdcae6606b83c8/task-sdk/src/airflow/sdk/execution_time/context.py#L815-L819)).
#67530 added the wrapper so a backend reference can be told apart from a user
value that happens to look like a path. A backend can also choose to keep a
value in the database, though. The common.io object-storage backend does that
below `state_store_objectstorage_threshold` by returning `json.dumps(value)`
([`backend.py`](https://github.com/apache/airflow/blob/3e302fd257a0cfe8deba2de529bdcae6606b83c8/providers/common/io/src/airflow/providers/common/io/state_store/backend.py#L218-L221)).
The SDK wraps that too, so a small value is stored as an escaped JSON string
inside a wrapper that says it is a reference, and that is what the UI and REST
API show:
| Value | Stored before | Stored after |
|---|---|---|
| `{"count": 5}`, below the threshold | `{"__airflow_state_ref__":
"{\"count\": 5}"}` | `{"count": 5}` |
| ~2 KB, above the threshold | `{"__airflow_state_ref__":
"file:///.../large"}` | unchanged |
Tasks still read the right value back. A key written through the REST API or
the UI's edit button also logged "was not written through the configured state
backend" on every read.
## Design rationale
`serialize_*_state_store_to_ref()` can now return `None`, meaning "store
this value as plain JSON", and the SDK wraps only the strings it returns.
References keep the wrapper, so the distinction #67530 introduced still holds.
The base class default changes from `json.dumps(value)` to `None`.
**The provider does not check the Airflow version.** Below the threshold,
common.io returns `super().serialize_*_state_store_to_ref(...)`.
`airflow.sdk.state.BaseStoreBackend` is the copy shipped with the installed
Task SDK, so that call returns `None` where the SDK supports inline values and
`json.dumps(value)` on older SDKs, which wrap every return value. Returning
`None` unconditionally would be stored as `{"__airflow_state_ref__": null}` on
those versions and the value would be lost. A version gate would have the same
problem on 3.4.0, which was branched before this change.
**Custom backends that return a string are unaffected.** The return type
widens from `str` to `str | None`, and returning `None` is opt-in. The docs now
tell backends that also support older Airflow versions to return `super()` for
values they keep inline.
**Existing rows still read back.** A row stored in the old format still has
the wrapper, so `get()` still passes it to `deserialize_*`, and both the base
class and common.io still JSON-decode it. Such a row keeps showing the wrapper
in the UI until the key is written again.
The "not written through the configured state backend" warning is removed.
With inline storage, a value without the wrapper is the normal case.
## Before / after
Run end to end with `airflow standalone` (LocalExecutor, sqlite), the
common.io backend on a `file://` path, and a 1024-byte threshold. One task
writes small and large task and asset state. A second task reads the asset
state back, including a key set through `PUT
/api/v2/assets/{id}/state-store/{key}` and a row seeded in the old format.
Task state for the writer task:
| Before | After |
|---|---|
|  |  |
Asset state:
| Before | After |
|---|---|
|  |  |
`legacy_inline` only exists in the after run. It is the old-format row, and
the reader task got `{'written_by': '3.3'}` back from it. In the before run,
reading `from_api` logged the warning; in the after run it did not.
## Gotchas
The SDK and shared-library changes need a backport to `v3-4-test` to ship in
3.4.0. The provider change is safe on any version, because it follows the
installed SDK.
---
* Read the **[Pull Request
Guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#pull-request-guidelines)**
for more information. Note: commit author/co-author name and email in commits
become permanently public when merged.
* For fundamental code changes, an Airflow Improvement Proposal
([AIP](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals))
is needed.
* When adding dependency, check compliance with the [ASF 3rd Party License
Policy](https://www.apache.org/legal/resolved.html#category-x).
* For significant user-facing changes create newsfragment:
`{pr_number}.significant.rst`, in
[airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments).
You can add this file in a follow-up commit after the PR is created so you
know the PR number.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]