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 |
   |---|---|
   | ![Task state before](./before_task_state.png) | ![Task state 
after](./after_task_state.png) |
   
   Asset state:
   
   | Before | After |
   |---|---|
   | ![Asset state before](./before_asset_state.png) | ![Asset state 
after](./after_asset_state.png) |
   
   `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]

Reply via email to