Scaling WebSockets Across Processes with Redis Pub/Sub¶
A WebSocket server that broadcasts by looping over its connected clients works perfectly — until there are two server processes. Each process knows only its own connections, so a message sent through one reaches only the clients attached to that one. Tested with two websockets server processes and 500 clients split evenly between them, sending 100 broadcast messages through a client on the first process: in-process broadcasting delivered 25,000 of 50,000 messages, and clients on the second process received none. Relaying every broadcast through a Redis pub/sub channel that both processes subscribe to delivered all 50,000, with end-to-end latency of 19.6 ms at the median and 36.9 ms at p99, against 8.6 / 21.5 ms for the single-process path. This guide wires up the relay, keeps it resilient, and covers what pub/sub does not guarantee.
Prerequisites¶
- Python 3.11+,
pip install websockets redis; measured with websockets 17.1, redis-py 5.3 and Redis 7. - Broadcasting within a process, from broadcasting to thousands of WebSocket clients.
- Slow-consumer handling, from handling WebSocket backpressure with slow consumers.
1. See why local broadcast breaks with more processes¶
Scaling a WebSocket service means more worker processes or more pods, each with its own set of connections:
clients: set = set()
async def handler(ws):
clients.add(ws)
try:
async for message in ws:
broadcast(clients, message) # only this process's clients
finally:
clients.discard(ws)
Measured with two processes: every client on the sender's process received all 100 messages; every client on the other process received 0. The bug is invisible in development, where one process serves everything, and appears as "some users miss updates" in production — usually discovered when someone scales from one replica to two.
Verify: run two server processes behind a load balancer, connect a client to each, and send a broadcast; both must receive it.
2. Publish to Redis, deliver from a relay task¶
Make every broadcast go through Redis. Handlers publish; one relay task per process subscribes and delivers to that process's local clients:
import redis.asyncio as aioredis
from websockets.asyncio.server import broadcast, serve
r = aioredis.Redis(host="redis", port=6379)
clients: set = set()
async def handler(ws):
clients.add(ws)
try:
async for message in ws:
await r.publish("room:lobby", message) # to every process, including this one
finally:
clients.discard(ws)
async def relay():
pubsub = r.pubsub()
await pubsub.subscribe("room:lobby")
async for item in pubsub.listen():
if item["type"] == "message":
broadcast(clients, item["data"].decode()) # local fan-out
async def main():
async with asyncio.TaskGroup() as tg:
tg.create_task(relay())
async with serve(handler, "0.0.0.0", 8765):
await asyncio.Future()
Measured: all 50,000 deliveries arrived, with p50 latency 19.6 ms and p99 36.9 ms end to end — the extra hop through Redis plus one client process receiving on 500 sockets. The publishing process receives its own message back through the subscription, so there is exactly one delivery path and no special case for local clients. broadcast from websockets writes to each connection without awaiting, skipping clients whose buffers are full, which keeps one slow client from delaying the rest.
Verify: latency from publish to receipt is measured in production and stays within your target at peak fan-out.
3. Use channels per room, not one for everything¶
With one channel, every process receives every message and discards the ones for rooms it has no clients in. Subscribe per room, and only while the process has local members:
class Rooms:
def __init__(self, r: aioredis.Redis) -> None:
self.pubsub = r.pubsub()
self.members: dict[str, set] = {}
async def join(self, room: str, ws) -> None:
if room not in self.members:
self.members[room] = set()
await self.pubsub.subscribe(f"room:{room}") # first local member
self.members[room].add(ws)
async def leave(self, room: str, ws) -> None:
members = self.members.get(room)
if members is None:
return
members.discard(ws)
if not members:
del self.members[room]
await self.pubsub.unsubscribe(f"room:{room}") # last local member left
async def relay(self) -> None:
async for item in self.pubsub.listen():
if item["type"] == "message":
room = item["channel"].decode().removeprefix("room:")
broadcast(self.members.get(room, ()), item["data"].decode())
Per-room subscriptions keep each process's inbound traffic proportional to the rooms it serves. For per-user delivery (notifications, revocations), a channel per user works the same way, and it is also how a logout on one instance can close that user's sockets on every instance, as needed in authenticating WebSocket connections.
Verify: a process with no members in a room receives no messages for it (PUBSUB NUMSUB room:<name> counts only processes with members).
4. Survive Redis disconnects¶
The relay depends on one long-lived subscription connection. If it drops, that process silently stops receiving broadcasts. Reconnect in a loop and resubscribe:
async def relay_forever(rooms: Rooms, r: aioredis.Redis) -> None:
delay = 0.5
while True:
try:
rooms.pubsub = r.pubsub()
if rooms.members:
await rooms.pubsub.subscribe(*(f"room:{name}" for name in rooms.members))
delay = 0.5
await rooms.relay()
except (ConnectionError, aioredis.ConnectionError):
log.warning("pub/sub connection lost; resubscribing in %.1fs", delay)
await asyncio.sleep(delay)
delay = min(delay * 2, 10)
Redis pub/sub is at-most-once: messages published while a subscriber is disconnected are not stored and are never delivered to it. For live presence or chat typing indicators that is acceptable. When clients must not miss messages, have them resynchronize on reconnect — fetch the latest state over HTTP — or use Redis Streams, which store messages and let each process resume from its last position, as in reading Redis Streams with consumer groups.
Verify: restart Redis during a test; the relay reconnects, resubscribes to every room with local members, and new broadcasts arrive.
5. Protect the fan-out from slow clients and big rooms¶
Every published message turns into one write per local member. Large rooms and slow clients turn that into memory and latency problems:
async def handler(ws):
...
async with serve(
handler, "0.0.0.0", 8765,
max_queue=16, # inbound frames buffered per connection
write_limit=64 * 1024, # outbound buffer high-water mark per connection
):
...
broadcast skips connections whose write buffer is above the limit instead of letting memory grow — those clients miss messages and should resynchronize, or be disconnected if they stay behind. Rate-limit publishes per client so one user cannot flood a room, and size Redis for the publish rate times the number of subscribing processes. Very large rooms (tens of thousands of members) are better served by a dedicated fan-out tier than by every application process.
Verify: a test client that stops reading does not raise the server's memory, and other clients' latency is unaffected.
Verification¶
Cross-process broadcasting works when:
- Every broadcast goes through the shared channel, never only to local clients.
- Processes subscribe per room while they have local members.
- The relay reconnects and resubscribes after Redis disconnects.
- Slow clients are skipped or dropped, and missed messages are recoverable where they matter.
Diagnostic Hook: publish a heartbeat message to a monitoring room every few seconds and measure, per process, the time since the last one was relayed. A process whose heartbeat age keeps growing has lost its subscription; compare PUBSUB NUMSUB with the number of processes that should be subscribed.
Pitfalls & edge cases¶
- Local broadcast with more than one process. Measured: half the deliveries missing.
- One global channel for every room. Every process receives everything.
- No relay reconnection. A dropped subscription silently stops broadcasts to that process.
- Treating pub/sub as reliable. Messages during a disconnect are gone.
Frequently Asked Questions¶
How do I broadcast WebSocket messages across multiple server processes?
Publish each broadcast to a Redis pub/sub channel and run a relay task in every process that subscribes and sends to that process's local clients. In testing this reached all 50,000 deliveries across two processes, against 25,000 with local broadcast.
How much latency does Redis pub/sub add to WebSocket broadcasts?
In testing with 500 clients, p50 latency was 19.6 ms through Redis against 8.6 ms for a single process, and p99 36.9 ms against 21.5 ms.
Does Redis pub/sub guarantee delivery?
No. It is at-most-once: a subscriber that is disconnected when a message is published never receives it. Use Redis Streams or have clients resynchronize when they must not miss messages.
Do sticky sessions fix WebSocket broadcasting?
No. Sticky sessions keep a client on one process, but a broadcast still has to reach clients on all the other processes.
Related¶
- WebSocket & Real-Time Streams — up to the topic overview.
- Handling WebSockets in FastAPI and Starlette — the same relay in an ASGI application.
- Network I/O & Protocol Handling — the section overview.