Zum Inhalt springen

WebSocket-Verbindungen über mehrere Server hinweg skalieren

Wie du WebSocket-Verbindungen über einen Server hinaus skalierst: Sticky Sessions, Pub/Sub-Fan-out mit Redis, Verbindungsstatus und Reconnect.

5 Min. Lesezeit
Architekturdiagramm, das mehrere WebSocket-Server zeigt, die über eine Redis-Pub/Sub-Schicht für die Nachrichtenverteilung miteinander verbunden sind

Ein einzelner Node.js-Server bewältigt problemlos 50.000 bis 100.000 gleichzeitige WebSocket-Verbindungen. Doch sobald du mehr Kapazität oder Redundanz brauchst, bricht mit einem zweiten Server alles zusammen. Client A verbindet sich mit Server 1, Client B mit Server 2. Wenn Client A eine Nachricht an Client B schickt, weiß Server 1 nichts von dessen Verbindung.

Das Kernproblem ist, dass WebSocket-Verbindungen zustandsbehaftet und an einen einzelnen Server gebunden sind. Deine Nachrichten brauchen einen Weg, diese Servergrenzen zu überwinden. Dieser Leitfaden stellt die bewährten Muster vor, mit denen sich WebSocket-Systeme von einem Server auf viele skalieren lassen.

Die Ausgangslage mit einem einzelnen Server

Bevor du skalierst, solltest du verstehen, wie eine Implementierung mit nur einem Server aussieht und an welcher Stelle sie versagt.

tstypescript
import { WebSocketServer, WebSocket } from "ws";
 
// Single-server implementation — works fine until you add a second server
const connections = new Map<string, WebSocket>();
const rooms = new Map<string, Set<string>>();
 
const wss = new WebSocketServer({ port: 8080 });
 
wss.on("connection", (ws, req) => {
  const userId = authenticateConnection(req);
  connections.set(userId, ws);
 
  ws.on("message", (data) => {
    const message = JSON.parse(data.toString());
 
    switch (message.type) {
      case "join_room":
        joinRoom(userId, message.roomId);
        break;
      case "room_message":
        broadcastToRoom(message.roomId, message.payload, userId);
        break;
      case "direct_message":
        sendToUser(message.targetId, message.payload);
        break;
    }
  });
 
  ws.on("close", () => {
    connections.delete(userId);
    removeFromAllRooms(userId);
  });
});
 
function sendToUser(targetId: string, payload: unknown): void {
  const target = connections.get(targetId);
  if (target?.readyState === WebSocket.OPEN) {
    target.send(JSON.stringify(payload));
  }
  // ❌ If targetId is on another server, this silently fails
}
 
function broadcastToRoom(
  roomId: string,
  payload: unknown,
  senderId: string
): void {
  const members = rooms.get(roomId);
  if (!members) return;
 
  for (const memberId of members) {
    if (memberId === senderId) continue;
    sendToUser(memberId, payload);
  }
  // ❌ Only reaches members connected to THIS server
}

Pub/Sub-Fan-out mit Redis

Die Standardlösung ist eine Pub/Sub-Schicht. Empfängt ein Server eine Nachricht, veröffentlicht er sie in einem Redis-Kanal. Alle Server abonnieren diesen Kanal und stellen die Nachricht ihren lokalen Verbindungen zu.

tstypescript
import Redis from "ioredis";
import { WebSocketServer, WebSocket } from "ws";
 
class ScalableWebSocketServer {
  private connections = new Map<string, WebSocket>();
  private pub: Redis;
  private sub: Redis;
  private serverId: string;
 
  constructor(port: number) {
    this.serverId = `server-${port}-${Date.now()}`;
    this.pub = new Redis(process.env.REDIS_URL);
    this.sub = new Redis(process.env.REDIS_URL);
    this.setupSubscriptions();
    this.setupWebSocket(port);
  }
 
  private setupSubscriptions(): void {
    this.sub.subscribe("ws:broadcast", "ws:direct", "ws:room");
 
    this.sub.on("message", (channel, data) => {
      const message = JSON.parse(data);
 
      // Skip messages we published ourselves
      if (message.sourceServer === this.serverId) return;
 
      switch (channel) {
        case "ws:direct":
          this.deliverLocal(message.targetId, message.payload);
          break;
        case "ws:room":
          this.deliverToLocalRoomMembers(
            message.roomId,
            message.payload,
            message.senderId
          );
          break;
        case "ws:broadcast":
          this.deliverToAll(message.payload);
          break;
      }
    });
  }
 
  sendToUser(targetId: string, payload: unknown): void {
    // Try local delivery first
    if (this.deliverLocal(targetId, payload)) return;
 
    // Publish for other servers to deliver
    this.pub.publish(
      "ws:direct",
      JSON.stringify({
        sourceServer: this.serverId,
        targetId,
        payload,
      })
    );
  }
 
  broadcastToRoom(
    roomId: string,
    payload: unknown,
    senderId: string
  ): void {
    // Deliver to local members
    this.deliverToLocalRoomMembers(roomId, payload, senderId);
 
    // Publish for other servers
    this.pub.publish(
      "ws:room",
      JSON.stringify({
        sourceServer: this.serverId,
        roomId,
        payload,
        senderId,
      })
    );
  }
 
  private deliverLocal(targetId: string, payload: unknown): boolean {
    const ws = this.connections.get(targetId);
    if (ws?.readyState === WebSocket.OPEN) {
      ws.send(JSON.stringify(payload));
      return true;
    }
    return false;
  }
 
  private deliverToLocalRoomMembers(
    roomId: string,
    payload: unknown,
    senderId: string
  ): void {
    // Room membership stored in Redis (shared across servers)
    // Local delivery only for connections on this server
  }
 
  private deliverToAll(payload: unknown): void {
    for (const ws of this.connections.values()) {
      if (ws.readyState === WebSocket.OPEN) {
        ws.send(JSON.stringify(payload));
      }
    }
  }
}

Verbindungszustand in Redis speichern

Raumzugehörigkeit und Anwesenheitsstatus müssen serverübergreifend geteilt werden. Nutze Redis-Sets und -Hashes, um festzuhalten, welche Nutzer sich in welchen Räumen befinden und welcher Server welche Verbindung hält.

tstypescript
class ConnectionRegistry {
  constructor(private redis: Redis, private serverId: string) {}
 
  async registerConnection(userId: string): Promise<void> {
    const pipeline = this.redis.pipeline();
 
    // Track which server holds this connection
    pipeline.hset("ws:connections", userId, this.serverId);
 
    // Track all connections on this server (for cleanup on crash)
    pipeline.sadd(`ws:server:${this.serverId}`, userId);
 
    // Set presence with TTL (heartbeat will refresh)
    pipeline.set(`ws:presence:${userId}`, "online", "EX", 60);
 
    await pipeline.exec();
  }
 
  async removeConnection(userId: string): Promise<void> {
    const pipeline = this.redis.pipeline();
    pipeline.hdel("ws:connections", userId);
    pipeline.srem(`ws:server:${this.serverId}`, userId);
    pipeline.del(`ws:presence:${userId}`);
    await pipeline.exec();
  }
 
  async joinRoom(userId: string, roomId: string): Promise<void> {
    await this.redis.sadd(`ws:room:${roomId}`, userId);
  }
 
  async leaveRoom(userId: string, roomId: string): Promise<void> {
    await this.redis.srem(`ws:room:${roomId}`, userId);
  }
 
  async getRoomMembers(roomId: string): Promise<string[]> {
    return await this.redis.smembers(`ws:room:${roomId}`);
  }
 
  // Clean up stale connections when a server crashes
  async cleanupServer(deadServerId: string): Promise<void> {
    const users = await this.redis.smembers(
      `ws:server:${deadServerId}`
    );
 
    const pipeline = this.redis.pipeline();
    for (const userId of users) {
      pipeline.hdel("ws:connections", userId);
      pipeline.del(`ws:presence:${userId}`);
    }
    pipeline.del(`ws:server:${deadServerId}`);
    await pipeline.exec();
  }
}

Reconnection auf Client-Seite

Clients werden die Verbindung verlieren — Netzwechsel, Server-Deployments, Timeouts am Load Balancer. Ein robuster Client kümmert sich transparent um die Reconnection und spielt dabei auch verpasste Nachrichten nach.

tstypescript
// ❌ Naive WebSocket client — no reconnection
const ws = new WebSocket("wss://api.example.com/ws");
ws.onmessage = (e) => handleMessage(JSON.parse(e.data));
// Connection drops → app is broken → user refreshes page
 
// ✅ Resilient WebSocket client with exponential backoff
class ResilientWebSocket {
  private ws: WebSocket | null = null;
  private reconnectAttempts = 0;
  private maxReconnectDelay = 30_000;
  private lastMessageId: string | null = null;
  private messageHandlers: Array<(data: unknown) => void> = [];
 
  constructor(private url: string) {
    this.connect();
  }
 
  private connect(): void {
    // Include last message ID for server to replay missed messages
    const connectUrl = this.lastMessageId
      ? `${this.url}?after=${this.lastMessageId}`
      : this.url;
 
    this.ws = new WebSocket(connectUrl);
 
    this.ws.onopen = () => {
      this.reconnectAttempts = 0;
    };
 
    this.ws.onmessage = (event) => {
      const data = JSON.parse(event.data);
      if (data.id) this.lastMessageId = data.id;
      this.messageHandlers.forEach((h) => h(data));
    };
 
    this.ws.onclose = (event) => {
      if (event.code === 1000) return; // Normal close
      this.scheduleReconnect();
    };
 
    this.ws.onerror = () => {
      this.ws?.close();
    };
  }
 
  private scheduleReconnect(): void {
    const delay = Math.min(
      1000 * Math.pow(2, this.reconnectAttempts) +
        Math.random() * 1000,
      this.maxReconnectDelay
    );
 
    this.reconnectAttempts++;
    setTimeout(() => this.connect(), delay);
  }
 
  onMessage(handler: (data: unknown) => void): void {
    this.messageHandlers.push(handler);
  }
 
  send(data: unknown): void {
    if (this.ws?.readyState === WebSocket.OPEN) {
      this.ws.send(JSON.stringify(data));
    }
  }
}

Load Balancing für WebSocket-Verbindungen

HTTP-Load-Balancer müssen für WebSocket-Unterstützung konfiguriert werden. Der anfängliche HTTP-Upgrade-Request baut die Verbindung auf, und alle folgenden Frames müssen zum selben Server geroutet werden.

nginxnginx
# Nginx configuration for WebSocket load balancing
upstream websocket_servers {
    # ip_hash ensures same client always hits same server
    # Alternative: use a shared session store and round-robin
    ip_hash;
 
    server ws-server-1:8080;
    server ws-server-2:8080;
    server ws-server-3:8080;
}
 
server {
    listen 443 ssl;
    server_name ws.example.com;
 
    location /ws {
        proxy_pass http://websocket_servers;
 
        # Required for WebSocket upgrade
        proxy_http_version 1.1;
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection "upgrade";
        proxy_set_header Host $host;
 
        # Increase timeouts for long-lived connections
        proxy_read_timeout 86400s;
        proxy_send_timeout 86400s;
 
        # Forward real client IP
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
    }
}
shbash
# ❌ Default nginx config drops WebSocket after 60 seconds
# proxy_read_timeout default is 60s
# Long-lived WebSocket connections timeout silently
 
# ✅ Increase timeouts and configure health checks
# proxy_read_timeout 86400s (24 hours)
# Combine with application-level heartbeat every 30 seconds
# If no heartbeat response in 90 seconds, client reconnects

Die wichtigsten Erkenntnisse

  1. WebSocket-Verbindungen sind an einen Server gebunden — ohne Pub/Sub-Schicht können Nachrichten keine Servergrenzen überqueren; Redis Pub/Sub ist die Standardlösung
  2. Raumzugehörigkeit und Anwesenheit in Redis speichern — jeder Server muss wissen, wer in welchem Raum ist; Redis-Sets liefern O(1)-Zugehörigkeitsprüfungen, die über alle Server hinweg geteilt werden
  3. Verwaiste Verbindungen nach einem Server-Absturz aufräumen — verfolge, welche Verbindung zu welchem Server gehört, damit ein Health-Checker aufräumen kann, wenn ein Server unerwartet ausfällt
  4. Die Reconnection des Clients muss automatisch ablaufen — exponentielles Backoff mit Jitter verhindert einen Thundering-Herd-Effekt beim Neustart des Servers; sende die ID der letzten Nachricht mit, um verpasste Nachrichten nachzuliefern
  5. Timeouts des Load Balancers für langlebige Verbindungen konfigurieren — Standard-Timeouts (60 s) beenden WebSocket-Verbindungen stillschweigend; kombiniere sie mit Heartbeats auf Anwendungsebene und großzügigeren Proxy-Timeouts
Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX