Skip to content

Async state store accessors block the event loop on sync ref (de)serialization #72129

Description

@jason810496

Summary

BaseStoreBackend exposes async aget/aset/adelete/aclear, but its ref (de)serialization hooks are sync-only:

  • shared/state/src/airflow_shared/state/__init__.py:251 serialize_task_state_store_to_ref
  • shared/state/src/airflow_shared/state/__init__.py:271 deserialize_task_state_store_from_ref
  • shared/state/src/airflow_shared/state/__init__.py:281 serialize_asset_state_store_to_ref
  • shared/state/src/airflow_shared/state/__init__.py:301 deserialize_asset_state_store_from_ref

The async accessors in task-sdk/src/airflow/sdk/execution_time/context.py await the supervisor round-trip but then call those sync hooks directly on the event-loop thread:

accessor path sync hook called on the loop thread
TaskStateStoreAccessor.aget context.py:593 deserialize_task_state_store_from_ref
TaskStateStoreAccessor.aset context.py:652 serialize_task_state_store_to_ref
AssetStateStoreAccessor.aget context.py:780 deserialize_asset_state_store_from_ref
AssetStateStoreAccessor.aset context.py:814 serialize_asset_state_store_to_ref

With a worker backend that does real I/O, aget/aset therefore still stall the loop. The async API promises non-blocking behaviour it does not deliver.

adelete/aclear are not affected — they go through backend.adelete/aclear, which are genuinely async on the base class.

Blast radius

The gap only manifests when [workers] state_store_backend is configured. With no worker backend, _get_worker_state_store_backend() returns None and the accessors never touch a serialization hook, so the async methods are fully non-blocking.

When a backend is configured, StateStoreObjectStorageBackend (providers/common/io) is the shipped case, and the blocking path is the default one, not an edge case: serialize_*_to_ref returns the value inline only when len(serialized) < [common.io] state_store_objectstorage_threshold, and that threshold defaults to 0, documented as "always offload to object storage". So every aset performs an object-storage write, and every aget of an offloaded key performs a read, on the event-loop thread.

Two callers are exposed:

  1. Watcher triggersBaseEventTrigger.run() is a coroutine, and every trigger on a triggerer shares one event loop, so one trigger's stall delays all of them.
  2. async def tasks — the task runner's loop stalls for the duration of the object-storage call.

Why this needs measurement before a fix

The severity is unquantified. The obvious fix (below) is cheap, but we should not land a perf-motivated change without numbers, and we do not currently know how visible the stall is in practice:

  • Object-storage latency varies by orders of magnitude — local filesystem, MinIO on the same host, and real S3 across a region are not the same problem. A fix justified only by a local-filesystem benchmark proves nothing.
  • The triggerer's blocking detector may or may not fire. TriggerRunner.block_watchdog (airflow-core/src/airflow/jobs/triggerer_job_runner.py:1576) samples every 100ms and reports when the gap exceeds [triggerer] blocked_main_thread_warning_threshold (default 0.2 s). A single S3 round-trip may land under that threshold while still hurting a triggerer running hundreds of watchers. We need to know which cases actually trip it.
  • Note on observability: the watchdog emits via log.ainfo, i.e. info level, despite the message calling 0.2s a "warning threshold". Anyone reproducing this must not filter to WARNING or they will see nothing. It also increments the triggers.blocked_main_thread statsd counter, which is the more reliable signal to assert on.
  • There is no equivalent detector on the task-runner side. An async def task that stalls its own loop is silent — no watchdog, no metric. DeadlockImminentError in comms.py is unrelated: it catches sync send() racing an in-flight asend(), not slow I/O.

Investigation

1. Benchmark

Measure wall-clock loop-stall per call for aget and aset, across:

  • backends: no worker backend (baseline), StateStoreObjectStorageBackend on local filesystem, on MinIO, on real S3 (cross-region)
  • payload sizes: below and above state_store_objectstorage_threshold, plus one large value (~1 MB) to exercise compression
  • compression: off and gzip
  • concurrency: 1, 10, 100 concurrent coroutines on one loop

Report the loop-stall distribution (p50/p95/max), not just the mean — the tail is what starves co-tenant coroutines. Compare against adelete/aclear, which already thread-offload once the companion provider change lands, to isolate the serialization cost from the supervisor round-trip.

2. Trigger blocking-warning matrix

For each cell, record whether block_watchdog logs and whether triggers.blocked_main_thread increments:

accessor no backend objectstorage / local fs objectstorage / MinIO objectstorage / S3
aget (key offloaded)
aget (key inline)
aset (below threshold)
aset (above threshold)
adelete
aclear

Run each at the default blocked_main_thread_warning_threshold (0.2s) and at a tightened value, so we can distinguish "does not block" from "blocks below the detection threshold". Also record the numbers with PYTHONASYNCIODEBUG=1, which the watchdog message itself points at.

3. Task-runner side

Confirm the stall is silent for async def tasks, and decide whether that asymmetry is worth its own follow-up (a task-runner block watchdog, or reusing the triggerer's).

Proposed fix

Add non-abstract async variants to BaseStoreBackend that default to a thread offload, so no existing backend breaks:

# shared/state/src/airflow_shared/state/__init__.py
    async def aserialize_asset_state_store_to_ref(
        self, *, value: JsonValue, key: str, scope: AssetScope
    ) -> str:
        """Async variant of ``serialize_asset_state_store_to_ref``.

        Defaults to offloading the sync implementation to a thread, so backends that
        have not overridden it stay correct and stop blocking the caller's loop.
        """
        return await asyncio.to_thread(
            self.serialize_asset_state_store_to_ref, value=value, key=key, scope=scope
        )

...and the matching adeserialize_asset_state_store_from_ref / aserialize_task_state_store_to_ref / adeserialize_task_state_store_from_ref. Then context.py's aget/aset await those instead of the sync hooks.

This means _build_set_message and _extract_get_response can no longer be shared verbatim between the sync and async paths (they would need to become async, or split at the backend call). Whichever way that lands, it touches the merged TaskStateStoreAccessor methods from #68232 as well as the asset ones, which is why it is a separate PR rather than a fixup.

A backend whose transport is natively async can override the a* variants directly instead of inheriting the thread offload.

Acceptance criteria

  • Benchmark results published in this issue, including the p95/max loop stall per backend and the concurrency sweep.
  • The trigger blocking-warning matrix filled in, with the pre-fix baseline recorded before any code change.
  • BaseStoreBackend gains async ref (de)serialization variants with a thread-offload default; no existing backend is forced to implement them.
  • TaskStateStoreAccessor.aget/aset and AssetStateStoreAccessor.aget/aset no longer call a sync backend hook on the event-loop thread.
  • A regression test asserts the blocking work runs off the loop thread (the threading.get_ident() pattern, not a timing assertion).
  • Post-fix run of the same matrix shows triggers.blocked_main_thread no longer increments for any state store accessor cell.

Metadata

Metadata

Assignees

Type

No type

Projects

No projects

Milestone

No milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions