One WebSocket server is easy: every connection lives in one process, and sending a message to user 42 is a lookup in a local map. Add a second server and the map splits in two. The service that produces an event for user 42 now has to find out which process holds that user's sockets, possibly several of them on different machines, and get the message there even while nodes restart.

This article is about that addressing problem. It does not repeat per-node connection limits, reconnect storms, backplane amplification or autoscaling, which scaling WebSocket servers: limits, storms and autoscaling covers. Here you will build the routing layer: the three delivery patterns and what each costs, a connection registry with leases, per-node inbox channels, a replay log for messages that would otherwise be lost, and the cleanup of routes that point at dead nodes.

Advertisement

The problem after the second server

A WebSocket is a long-lived TCP connection terminated by one process. The load balancer picks that process once, at the handshake, and every frame afterwards goes to the same place. That is what pins state: the socket object, its send buffer and its subscriptions exist in exactly one node's memory. Meanwhile the events you want to push are produced elsewhere, by an order service, a chat backend or a job worker, which have no idea where any socket lives.

So a horizontally scaled real-time system has two layers. Gateway nodes hold sockets and know only their own connections. A delivery layer takes an event addressed to a user, a room or everyone, and makes it arrive at the right gateways. Load balancing at the handshake, including draining gateways during deploys, is a separate concern covered in load balancing WebSockets. Everything below is about the delivery layer.

Three delivery patterns and what each costs

Broadcast to every node. Every event is published on one shared channel; every gateway receives it and checks its local map. It is the simplest thing that works and needs no registry, but every gateway processes every event. Topic subscription. Each gateway subscribes to a channel per room or per user it currently serves, and events are published to that channel. The backplane does the filtering, at the cost of subscription churn as users connect and disconnect. Directed delivery. A registry records which nodes hold each user's connections; a router looks the user up and publishes only to those nodes' inbox channels. It costs one lookup per event and a registry to keep correct.

PatternBackplane deliveries per eventState to maintainBest fit
BroadcastOne per gateway nodeNoneFew nodes, or events most users need
Topic per room or userOne per subscribed nodeSubscriptions on the backplaneRooms with members spread over a few nodes
Directed via registryOne per node holding the userRegistry with leasesPer-user notifications at scale

Worked numbers make the choice concrete. Suppose 40 gateway nodes and 50,000 per-user events per second, with users averaging 1.5 connected devices. Broadcast delivers 50,000 x 40 = 2,000,000 messages per second into gateways, nearly all of them discarded after a map miss. Directed delivery delivers at most 50,000 x 1.5 = 75,000, plus 50,000 registry reads. Broadcast stays reasonable for a handful of nodes or for genuine broadcasts such as a status banner; directed delivery is what lets per-user traffic grow with the fleet.

Advertisement

The architecture of directed delivery

Directed delivery: registry lookup, then one publish per owning nodeOrder serviceevent for user 42Routerlookup, then publishRegistryuser 42: A/c17, C/c911 event2 lookupinbox:Anode A channelinbox:Bnothing sentinbox:Cnode C channel3 publish3 publishGateway Aconn c17 to socketGateway Bother usersGateway Cconn c91 to socketPhoneuser 42Laptopuser 42Per-user stream keeps a replay log for resume
The router reads the registry for user 42, finds connections on nodes A and C, and publishes once to each node's inbox. Node B receives nothing.

Each gateway subscribes to exactly one inbox channel named after itself, for example inbox:node-a, when it starts. When a client connects and authenticates, the gateway writes a registry entry saying that user 42 has connection c17 on node A. When the order service emits an event for user 42, the router reads all of the user's entries, groups them by node and publishes one message per node, carrying the user ID and the payload. Each gateway looks up its local connections for that user and writes the frame to each socket. The registry and inboxes can live in Redis, which is the example below, or in any store with expiry plus any pub/sub system.

A connection registry with leases

The registry must survive the thing that breaks it most often: a gateway dying without cleaning up. So every entry is a lease with an expiry, renewed by the gateway while the connection is alive. A sorted set per user works well, with the member identifying node and connection and the score holding the lease expiry time. Readers ignore expired members, and a periodic trim deletes them.

import time
import redis.asyncio as redis

LEASE = 45          # seconds; renew every 15
r = redis.Redis()

def member(node, conn_id):
    return f"{node}|{conn_id}"

async def register(user, node, conn_id):
    await r.zadd(f"conns:{user}", {member(node, conn_id): time.time() + LEASE})
    await r.expire(f"conns:{user}", LEASE * 2)

async def renew(user, node, conn_id):          # called from the heartbeat loop
    await register(user, node, conn_id)

async def unregister(user, node, conn_id):     # removes only our own member
    await r.zrem(f"conns:{user}", member(node, conn_id))

async def route(user, payload):
    now = time.time()
    await r.zremrangebyscore(f"conns:{user}", "-inf", now)
    live = await r.zrangebyscore(f"conns:{user}", now, "+inf")
    nodes = {m.decode().split("|", 1)[0] for m in live}
    for node in nodes:
        if not await r.exists(f"alive:{node}"):   # node stopped heartbeating
            await drop_node_routes(user, node, live)
            continue
        await r.publish(f"inbox:{node}", f"{user}|{payload}")
    return len(nodes)

async def drop_node_routes(user, node, members):
    stale = [m for m in members if m.decode().startswith(node + "|")]
    if stale:
        await r.zrem(f"conns:{user}", *stale)

Two details in that code prevent real bugs, and a third handles dead nodes. First, unregister removes only its own member. A simpler design stores one value per user, user 42 is on node A, and deletes the key on disconnect. Then a phone that switches networks reconnects to node C, C writes its entry, and moments later node A notices the old socket closed and deletes the key, erasing the new route. Keying entries by connection makes removal precise. Second, the lease is renewed from the same heartbeat loop that keeps the socket alive, so a node that stops heartbeating stops renewing, and its routes expire on their own. Each gateway also refreshes a short-lived alive:<node> key every few seconds, which the router checks before publishing.

Pub/sub is at-most-once: add a replay log

Redis PUBLISH delivers a message to the clients subscribed at that instant and keeps nothing. If node C is restarting when the event arrives, or the client is between a dropped connection and its reconnect, the message is gone. For a typing indicator that is fine. For an order notification it is not, and no amount of registry accuracy fixes it, because the gap is in time rather than in addressing.

The standard fix is a per-user or per-room log alongside the live path. The router appends every durable event to a Redis stream before publishing, and the client remembers the ID of the last event it processed. On reconnect the client sends that ID. The new gateway registers the connection first, so live delivery is already on, and then replays everything after that ID; an event emitted between the two steps arrives by at least one path.

async def emit(user, payload):
    event_id = await r.xadd(f"events:{user}", {"p": payload}, maxlen=1000, approximate=True)
    await route(user, f"{event_id.decode()}|{payload}")

async def resume(user, last_id, send):
    # exclusive start: everything strictly after last_id
    for event_id, fields in await r.xrange(f"events:{user}", min=f"({last_id}", max="+"):
        await send(event_id.decode(), fields[b"p"].decode())

Because both paths run at once, the client must deduplicate by event ID: an event published live while replay is running can arrive twice. The stream is capped, so a client that was offline longer than the retention window needs a full state refresh instead of a replay, and the server should say so explicitly. Per-user streams also give per-user ordering for free, since IDs increase monotonically; ordering across users or rooms is a different problem, discussed in message ordering architecture. If you need longer retention or consumer groups, a durable broker such as NATS JetStream or Kafka plays the same role.

Dead nodes and stale routes

When a gateway crashes, three things happen at once. Its clients see the connection drop and start reconnecting elsewhere, with backoff as described in reconnection strategies. Its registry entries remain until their leases expire. And its inbox channel has no subscriber, so messages routed there vanish.

The lease bounds how long stale routes last, which is why the lease should be a small multiple of the renewal interval: 45 seconds with renewal every 15 tolerates two missed renewals. The router can do better than waiting by checking a node liveness key, as the code above does: a gateway that has stopped refreshing its alive:<node> key within a few seconds is treated as gone, and the router deletes that node's members for the user. It is tempting to use the integer that PUBLISH returns instead, but the Redis documentation notes that in Redis Cluster it counts only subscribers connected to the same node as the publisher, so a healthy gateway subscribed through another cluster node reads as zero; even on a single server, a gateway whose subscriber connection is briefly reconnecting reads as zero. Treat that count as a hint at most. The replay log covers the messages that fell into the gap: the client reconnects, presents its last event ID and receives what it missed.

Keep gateway nodes disposable

The design works only if losing a gateway loses nothing that cannot be rebuilt. Audit what lives in gateway memory. Socket objects and send buffers are inherently local, and that is acceptable. Subscriptions to rooms should be re-declared by the client after reconnect rather than reconstructed from server state. Authentication should be re-validated at each handshake from the token the client presents, not cached across nodes. Cursors such as the last event ID belong to the client. Anything else, such as unsent message queues or per-user rate-limit counters, either moves to shared storage or is explicitly accepted as lossy.

Disposability is also what makes deploys and scale-in safe. A draining node stops accepting handshakes, closes its sockets with close code 1001 (going away) so clients reconnect promptly, and lets its leases expire or deletes its own members. Because resume is built in, the drain does not need to be perfect.

Scaling the backplane itself

In Redis Cluster, classic PUBLISH is propagated to every node in the cluster, so adding cluster nodes adds capacity for keys but not for pub/sub throughput. Redis 7.0 introduced sharded pub/sub, SSUBSCRIBE and SPUBLISH, which assigns each channel to a hash slot so a message touches only the shard that owns the channel. Per-node inbox channels suit this well: there are as many channels as gateways, they spread across shards, and each gateway holds one subscription. Registry keys per user spread across slots naturally. Watch the hottest keys, such as a celebrity user with thousands of connections or a room with tens of thousands of members, which may justify broadcast or a dedicated topic channel for that one entity; real-time fanout discusses those large fan-outs.

Failure modes

  • Reconnect race erases routes. A late disconnect handler deletes a newer registration. Key entries by connection ID and remove only your own.
  • Lease longer than tolerance. A five-minute lease means five minutes of messages routed to a dead node. Shorten it, and check node liveness keys before publishing.
  • No replay. Pub/sub alone loses every message sent during restarts and reconnects, which users report as missing notifications.
  • Unbounded logs. Streams without a cap grow with every inactive user. Cap by length or age and send a refresh instruction to clients beyond it.
  • Router as a bottleneck. One registry read per event is cheap until it is not; batch lookups for multi-user events and cache negative results briefly.
  • Hidden gateway state. A per-connection queue or counter that is not rebuilt after reconnect becomes a correctness bug that only appears during deploys.

Trade-offs to decide explicitly

Broadcast trades backplane volume for simplicity and is right for small fleets. Directed delivery trades a registry and its correctness work for cost that scales with recipients rather than nodes. Lease length trades registry write load against the window of stale routes. Replay retention trades memory against how long a client can be offline and still resume. And at-most-once live delivery plus a replay log is usually a better trade than making the live path durable, because it keeps the hot path fast and confines durability to one append per event.

What to do next

  1. Write down which events are droppable and which must survive a reconnect.
  2. Pick a delivery pattern per event type using the cost table, with your real node count and devices per user.
  3. Implement the registry with per-connection members and leases renewed from the heartbeat loop.
  4. Give every gateway one inbox channel and a short-lived liveness key, and drop routes to nodes whose key has expired.
  5. Add a capped per-user stream for durable events, client-side last-event IDs and deduplication.
  6. Kill a gateway under load in a test environment and measure messages lost and time to recovery.
Key takeaway: Scaling WebSocket servers horizontally turns delivery into an addressing problem: the socket lives on one node and the event is produced somewhere else. Broadcast is simple but costs one delivery per node; directed delivery uses a lease-based registry and per-node inboxes so cost follows recipients. Key registry entries by connection to survive reconnect races, treat pub/sub as at-most-once and add a capped replay log with client-held event IDs, use node liveness keys to drop dead routes, and keep gateways disposable so crashes and deploys cost a reconnect, not data.