A WebSocket server scales differently from an HTTP server. An HTTP request lasts milliseconds, so the load balancer re-decides placement constantly and new capacity is used at once. A WebSocket connection lasts minutes or days: placement is decided once, at the upgrade, and then frozen. Everything easy about scaling stateless HTTP becomes a deliberate engineering problem.

This page is about scaling out: what caps a single node, why adding nodes does not relieve hot ones, how to move connections without a stampede, how messages reach the right node, and what to autoscale on. It assumes you already know the protocol and the basic single-tier architecture; for the upgrade path, per-connection loop and draining on deploy, read WebSocket architecture, in depth first.

Advertisement

The shape of the problem

Clientsbrowsers, appsL4/L7 balancerleast-connectionsWS node 142k connsWS node 241k connsWS node 3 (new)0 connsAutoscalerconns per nodeBackplanesharded by topicConnection registryuser -> nodeProducersapp servicesRebalance controllerclose 1012 to a jittered fraction of over-full nodes; clients reconnect with full-jitter backoffmetricspublishsubscribed topicsScale-out adds capacity only for NEW connections; existing sockets stay where they landed until something moves them.
A horizontally scaled WebSocket tier. The balancer places each connection once; a backplane carries messages between nodes; a rebalance controller and an autoscaler watch connection counts, not CPU.

Three properties drive every decision here. Connections are sticky: a socket on node 2 stays there until one side closes it. Load is idle state plus bursts: one inbound message to a popular channel becomes thousands of outbound frames. Clients reconnect, all at once if you let them: a crash, deploy or network blip turns them into one synchronised herd.

What caps a single node

Every open WebSocket is a TCP socket, which is a file descriptor, kernel buffers, and whatever per-connection state your application keeps. Work through the limits in order, because the first one you hit is usually not the one you expected.

File descriptors. Many distributions still start services with a soft limit of 1,024 open files. A server that silently stops accepting at about a thousand connections is almost always this. Raise the per-process limit (LimitNOFILE in systemd, or limits.conf) and the system ceilings, then verify from inside the running process, because a limit set in the wrong place is the most common way this fix fails.

# /etc/security/limits.d/ws.conf  (or LimitNOFILE= in the systemd unit)
ws  soft  nofile  1048576
ws  hard  nofile  1048576

# sysctl: system-wide ceilings, then the per-process ceiling
fs.file-max = 4194304
fs.nr_open  = 1048576

# On the proxy tier only: more source ports per (src IP, dst IP, dst port)
net.ipv4.ip_local_port_range = 10240 65000

# Check what you actually got, from inside the running process
cat /proc/$(pgrep -f ws-server)/limits | grep "open files"
ss -s                     # socket totals by state

Memory per connection. The kernel side is socket buffers, which autotune within net.ipv4.tcp_rmem and tcp_wmem: small when idle, large behind a slow reader. The application side is your buffers, write queue, session and subscriptions, plus, if you enable permessage-deflate (RFC 7692), a compression context per connection that can dwarf the rest. Do not take a number from a blog post. Open 10,000 idle test connections against one node, measure resident memory, repeat with your real subscriptions and compression, and size nodes from bytes per connection with headroom for burst queues.

Ephemeral ports at the proxy. Clients bring their own source IPs, so they do not exhaust ports. But a proxy tier (NGINX, Envoy, HAProxy, a sidecar) that opens a second connection to the backend uses a source port per backend connection, and the four-tuple must be unique. With the common Linux default range of 32768 to 60999 that is about 28,000 concurrent connections from one proxy IP to one backend address. Widen the range, give backends several addresses or ports, or add proxy source IPs.

Idle timeouts in the path. NGINX's default proxy_read_timeout is 60 seconds and cloud load balancers have their own idle timeouts. Ping comfortably inside the shortest timeout on the path, and confirm the values for your products.

Advertisement

Why scale-out does not relieve hot nodes

Suppose six nodes each hold 50,000 connections and you add two. Least-connections balancing sends every new connection to the empty nodes, but the original six keep their 300,000 sockets. If they were near their memory or fan-out limit, you have bought nothing until clients reconnect for other reasons, which with long sessions can take hours. Round-robin is worse: it keeps feeding the hot nodes.

To move existing load you must close connections on purpose. Restarting a hot node creates exactly the herd you want to avoid. The safe pattern is controlled shedding: a controller compares per-node counts with the mean and asks over-full nodes to close a small random fraction of sockets, spread over time, with a close code that marks a planned move. The IANA WebSocket close code registry lists 1012 (Service Restart) and 1013 (Try Again Later) for this kind of purpose; they are registry entries rather than part of RFC 6455 itself, so check that your client libraries surface them.

# Rebalance controller: runs every 30 s, moves at most a small fraction per tick.
TARGET_SLACK = 1.10        # tolerate 10% above the mean before acting
MAX_MOVE_FRACTION = 0.02   # never close more than 2% of a node's sockets per tick

def plan_moves(conns_by_node: dict[str, int]) -> dict[str, int]:
    mean = sum(conns_by_node.values()) / len(conns_by_node)
    moves = {}
    for node, n in conns_by_node.items():
        if n > mean * TARGET_SLACK:
            excess = n - mean
            moves[node] = int(min(excess, n * MAX_MOVE_FRACTION))
    return moves

def apply(node: str, count: int):
    # The node picks `count` random sockets (not the oldest: those are often the
    # most valuable sessions) and closes them with 1012 spread over the next 30 s.
    send_admin(node, {"op": "shed", "count": count, "spread_s": 30, "code": 1012})

Choose victims at random, not oldest-first (long sessions are often your most engaged users), cap moves per tick so a controller bug cannot empty a node, and skip nodes already draining.

Reconnect storms, worked through

A node holding 50,000 connections crashes. If every client reconnects immediately, the survivors receive 50,000 TLS handshakes, upgrades, authentication calls and resubscriptions in about a second. The authentication service, sized for login traffic, usually breaks first; handshakes time out, clients retry, and the storm outlives the original fault.

Now give clients capped exponential backoff with full jitter: each attempt waits a uniformly random time between zero and the current window. A 500 ms window would still mean 100,000 arrivals per second, so a client that loses an open connection starts at a 4 s window, and at 8 s on a planned close (1012 or 1013). With a starting window of 4 s the 50,000 reconnects average 12,500 per second; if the cluster can absorb about 2,000 handshakes per second per node and seven nodes remain, it can take 14,000 per second, so the storm clears in about four seconds with no retry cascade. Without jitter, the same backoff schedule merely moves the whole herd to the same later instant.

// Client reconnect with capped exponential backoff and FULL jitter.
// The random draw over the whole window is what spreads a herd out.
function backoffMs(attempt, baseMs = 500, capMs = 30000) {
  const window = Math.min(capMs, baseMs * 2 ** attempt);
  return Math.random() * window;
}

function connect(url, attempt = 0) {
  const ws = new WebSocket(url);
  let opened = false;
  ws.onopen = () => { opened = true; attempt = 0; resubscribe(ws); };
  ws.onclose = (ev) => {
    // 1012 (Service Restart) / 1013 (Try Again Later) are IANA-registered codes the
    // server uses for planned moves. Losing an open connection starts from a wide window
    // so a whole node's clients do not stampede the survivors.
    const planned = ev.code === 1012 || ev.code === 1013;
    const next = opened ? (planned ? 4 : 3) : attempt + 1;   // 4 s window, 8 s if planned
    setTimeout(() => connect(url, next), backoffMs(next));
  };
}

On the server, put a token-bucket limit on new upgrades per node and reject excess with a fast HTTP 503 before expensive work; a quick rejection is far cheaper than a connection that times out halfway through authentication. Validate short-lived signed tokens locally instead of calling an identity service per connection, and let clients resubscribe in one message.

Fan-out: the backplane and its amplification

With twenty nodes, a room's members are spread everywhere, so a message published on node 4 must reach every node holding a member. That is the job of the backplane: a pub/sub system (Redis, NATS, Kafka, a cloud bus) that nodes subscribe to.

Let a topic have M members over N nodes and receive R messages per second. Socket writes are M × R whatever you do. Backplane deliveries are R × (nodes subscribed to the topic). Under naive broadcast every node receives every message, so traffic is R × N per topic and grows with the cluster even for tiny topics; nodes burn CPU decoding messages for rooms they have no members in.

Two techniques fix it. Interest-based subscription: a node subscribes to a topic only while it holds at least one local member, and unsubscribes when the last leaves. Sharding the backplane: hash topics onto several broker shards so no single broker carries all traffic. The sketch below does both.

# Topic-sharded backplane: each topic hashes to one of N broker shards.
# A node subscribes only to shards that carry topics its local sockets care about.
import hashlib

N_SHARDS = 16

def shard_for(topic: str) -> int:
    return int.from_bytes(hashlib.blake2b(topic.encode(), digest_size=4).digest(), "big") % N_SHARDS

class Node:
    def __init__(self, broker):
        self.broker = broker
        self.local = {}            # topic -> set(socket)

    def join(self, sock, topic):
        if topic not in self.local:
            self.local[topic] = set()
            self.broker.subscribe(shard_for(topic), topic, self.on_message)
        self.local[topic].add(sock)

    def leave(self, sock, topic):
        subs = self.local.get(topic)
        if subs is not None:
            subs.discard(sock)
            if not subs:
                del self.local[topic]
                self.broker.unsubscribe(shard_for(topic), topic)

    def on_message(self, topic, payload):
        for sock in self.local.get(topic, ()):
            sock.enqueue(payload)  # bounded per-socket queue; drop or close on overflow

Worked numbers: a 5,000-member live event channel across 10 nodes at 20 messages per second is 100,000 socket writes per second (unavoidable) and 200 backplane deliveries per second (20 messages × 10 nodes). A thousand private two-person chats at one message per second each would cost 10,000 backplane deliveries per second under broadcast and about 2,000 with interest-based subscription, because each chat lives on at most two nodes. Small topics are where broadcast wastes most.

Direct messages to a user need a different structure: a connection registry mapping user id to the node or nodes holding that user's sockets, written on connect and removed on close, with a TTL refreshed by heartbeat so a crashed node's entries expire. Publish to the node-specific channel rather than a global one.

Autoscaling on the right signal

CPU is the wrong primary signal for WebSocket nodes. A node full of idle connections can be at 15% CPU and one popular broadcast away from saturation, and its memory is committed regardless of CPU. Scale on connections per node against a target derived from your measured per-connection memory and burst fan-out, with CPU and outbound queue depth as secondary guards.

SignalUse it forWatch out for
Open connections per nodePrimary scale-out triggerLags during reconnect storms; smooth over minutes
Upgrade rate and rejectionsDetecting storms and auth bottlenecksSpikes are expected after deploys
Outbound queue depth and dropsSlow consumers, fan-out saturationOne huge topic can dominate
Resident memory per nodeValidating the per-connection budgetCompression contexts inflate it
Backplane deliveries per nodeBroadcast waste as the cluster growsShould not grow linearly with N

Scale-in is the dangerous direction. Treat it as a drain: remove the node from the balancer, shed its connections with 1012 over minutes, and terminate when nearly empty or at a hard deadline. Use a long scale-in cooldown so the autoscaler does not remove a node the rebalancer just filled.

Failure modes

  • Herd after deploy. A rolling restart dumps each node's clients onto the next node to be restarted. Drain with spread-out closes and rebalance after the rollout.
  • Broadcast backplane wall. Adding nodes raises CPU on every node. A high ratio of backplane messages received to messages delivered locally means broadcast waste.
  • Stale registry entries. A crashed node's users appear online and direct messages vanish. Use TTL-refreshed entries and treat publish-to-nobody as a metric.
  • Lost messages during moves. A client that reconnects to a new node misses messages sent in the gap. Give messages sequence numbers per topic and let clients resume from the last one they saw, served from a short retention buffer (streams rather than fire-and-forget pub/sub). The backpressure guide covers what to do when a resumed client cannot keep up.

Trade-offs to decide explicitly

Stateless versus sticky sessions. If session state lives in a shared store, any node can serve any client and moves are cheap; state in node memory is faster but turns every move into a migration. Interest-based subscriptions cut backplane traffic but add churn when small rooms flap. WebSocket versus SSE. For mostly server-to-client traffic, SSE over HTTP/2 reuses ordinary HTTP infrastructure with the same fan-out and reconnect problems. Per-message overhead for small updates is covered in the frame format deep dive.

What to do next

  1. Read the open-files limit from /proc for your running server and raise it until it is well above your target connections per node.
  2. Load-test one node with idle connections and then with your real subscriptions and compression setting; record bytes per connection and set a per-node connection target with burst headroom.
  3. Switch the WebSocket tier to least-connections balancing and confirm every idle timeout on the path is longer than your ping interval.
  4. Ship full-jitter backoff in every client, with a larger starting window for close codes 1012 and 1013, and add a per-node upgrade rate limit that fails fast.
  5. Measure backplane messages received versus delivered per node; if the ratio grows with cluster size, move to interest-based subscriptions and shard the backplane by topic.
  6. Add a rebalance controller that sheds a capped random fraction of over-full nodes, and autoscale on connections per node with a long scale-in cooldown.
  7. Run a game day: kill a node at peak load and confirm the reconnect wave clears within your target without auth-service errors.
Key takeaway: WebSocket placement is decided once per connection, so scaling out only helps new connections until you move old ones deliberately. Raise and verify per-node limits, measure memory per connection instead of guessing, rebalance by shedding small random fractions with planned-close codes, make every client back off with full jitter, keep the backplane from broadcasting everything everywhere, and autoscale on connections per node rather than CPU.