Skip to content

Fix streaming on multi-worker ASGI with Redis shared storage - #4054

Open
T4rk1n wants to merge 2 commits into
devfrom
fix/asgi-redis-streaming
Open

T4rk1n wants to merge 2 commits into
devfrom
fix/asgi-redis-streaming

Conversation

@T4rk1n

@T4rk1n T4rk1n commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

On uvicorn --workers 4 with RedisSharedStorage, FastAPI and Quart failed 24-57% of streaming callbacks at only 25 browsers, with p95 frame latency of 10-35 s. The same apps were fine on one worker with local storage, and Flask/gunicorn was fine on the same Redis.

Cause

Redis subscriptions used PollingSubscription, whose async iterator ran a blocking XREAD (5 s block) in the loop's default executor: one thread per open downlink, out of 12 on an 8-core machine. aget/aset/apublish used the same executor, so at a dozen downlinks per worker every publish queued behind the parked reads for 5-15 s and streams ran into the client's timeout. Instrumenting the workers showed nothing else: no stream was cancelled by the cross-worker connection-record check.

Changes

  • Redis: a native redis.asyncio path, one client per event loop. aget/aset/adelete/apublish are loop-native. One reader task per loop serves every async subscription with a single multi-stream XREAD, so it is one Redis connection per worker, not one per browser. A new subscription wakes the read through a private wake stream. Gaps are detected by sequence contiguity; a stream reset under a subscriber is caught on quiet cycles. The sync path is unchanged.
  • Diskcache: its async subscriptions held an executor thread for a 1 s sleep-poll each. They now poll without blocking and wait on the loop.
  • Downlink: on ASGI it wrote its connection record with sync calls on the event loop. It now uses the async calls, and writes its closed record from a task of its own, since a client disconnect cancels the response.

Shared storage and streaming are unreleased (#3930, #3931), so no CHANGELOG entry.

Numbers

benchmarks/streaming load sweep (#4055), Redis, 4 workers, 8-core laptop:

browsers FastAPI p50 / p95 Quart p50 / p95 errors
25 0.6 / 3.3 ms 0.7 / 5.5 ms 0%
1000 3.5 / 40.9 ms 3.2 / 41.2 ms 0%
4000 227 / 518 ms 249 / 528 ms 0%

Before: 24-57% errors at 25 browsers.

Tests

  • tests/streaming/test_stream_asgi_redis.py: uvicorn --workers 2 on Redis, 80 idle downlinks open, then a real browser stream must finish in under 4 s (FastAPI and Quart). Before the change it was stuck at the first frame after 15 s.
  • tests/shared_storage/test_redis_backend.py: async subscriptions hold no executor threads (fails before), replay, gap, close from another thread, a passed client, a subscriber joining behind a read in flight (no false gap), per-loop clients dropped once their loop closes.
  • tests/shared_storage/test_diskcache_backend.py: async subscriptions hold no executor threads (fails before).
  • tests/streaming/test_stream_transport.py: a cancelled async downlink still records itself closed.

Not in this PR

  • With --workers N, uvicorn binds its socket with protocol 0, so asyncio never sets TCP_NODELAY and every first frame waits ~40 ms on Nagle plus delayed ACK (3 ms on one worker). That is uvicorn, not Dash.
  • A downlink closing reads its connection record and then writes "closed" in two steps; a new downlink opening in between can be marked closed. This predates the change and needs a compare-and-set in the shared-storage API.

Redis subscriptions ran a blocking XREAD in the loop's default executor, one
thread per open downlink, and aget/aset/apublish used the same executor. At a
dozen downlinks per worker every publish queued behind the parked reads for
5-15 s, so on uvicorn --workers 4 a quarter to half of streams failed at 25
browsers.

The Redis backend now has a native redis.asyncio path: one client per event
loop, loop-native key/value and publish, and one reader task per loop that
serves every async subscription with a single multi-stream XREAD (woken for
new subscriptions through a private wake stream). Diskcache's async
subscriptions poll without blocking and sleep on the loop. The ASGI downlink
writes its connection record with the async calls, and its closed record
from a task of its own so a client disconnect can't skip it.
@github-actions

github-actions Bot commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

Dash performance benchmarks

✅ all within thresholds

scenario metric p90 (ms) median growth baseline p90 note
✅ callback_chain chain_ms 444.7 437.9 0.93x 499.1
✅ callback_chain graph_ms 2.8 2.8 1.0x 2.4
✅ callback_fanout fanout_ms 86.8 81.7 0.91x 92.5
✅ deep_nesting render_ms 58.7 57.2 0.97x 56.8
✅ full_children_replace replace_ms 5465.2 1955.9 18.3x 4697.1
✅ initial_render_large render_ms 616.1 599.6 1.01x 694.4
✅ initial_render_small render_ms 96.5 90.7 0.95x 104.0
✅ patch_append_nested append_ms 132.8 96.6 2.68x 192.3
✅ patch_append_toplevel append_ms 126.9 83.5 2.49x 140.2
✅ patch_scalar_update_large update_ms 156.8 141.3 0.94x 202.9
✅ wildcard_all_resolve wildcard_ms 299.7 284.0 0.95x 313.6
✅ wildcard_all_resolve graph_ms 1.3 1.3 1.0x 1.3

growth = late-third / early-third per-op time; ~1 is flat, a large value means the per-op cost scales with accumulated state.

machine scale vs baseline: 0.94x - divided out of the baseline ratios so they compare like for like (the absolute warn/fail ceilings are left un-scaled); calibrated on initial_render_small.

Split the reader's dispatch and failure paths out of its run loop, drop
copies of subscriber sets that are never mutated while iterated, keep one
possible raise per pytest.raises block, and silence pylint's import-error
on the redis.asyncio import, which CI's lint environment can't resolve.
@sonarqubecloud

sonarqubecloud Bot commented Oct 7, 2026

Copy link
Copy Markdown

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants