Zum Inhalt springen

WebSocket-Architekturmuster für skalierbare Echtzeit-Apps

Skalierbare WebSocket-Architekturen entwerfen: Verbindungsverwaltung, raumbasiertes Pub/Sub, Heartbeats, Reconnect und horizontale Skalierung.

4 Min. Lesezeit
Ein Diagramm einer WebSocket-Serverarchitektur, das einen Load Balancer zeigt, der Verbindungen auf mehrere Serverinstanzen mit Redis-Pub/Sub verteilt.

Die Grenze eines einzelnen Servers

WebSocket-Verbindungen sind dauerhaft und zustandsbehaftet. Ein einzelner Server, der alle Verbindungen hält, funktioniert – bis ein zweiter Server nötig wird. Von da an erreicht eine Nachricht, die an Server A gesendet wird, die Clients auf Server B nicht mehr. Die Skalierung von WebSocket-Anwendungen erfordert architektonische Entscheidungen, mit denen klassische HTTP-Anwendungen nie konfrontiert sind.

Verbindungsverwaltung

Jede WebSocket-Verbindung muss nachverfolgt werden. Ohne eine Registry kannst du Nachrichten nicht gezielt an bestimmte Nutzer, Räume oder Kanäle senden.

tstypescript
// ❌ Storing connections in a plain array — no way to target specific users
const connections: WebSocket[] = [];
 
// ✅ Connection registry with metadata for targeted messaging
interface ConnectionMeta {
  userId: string;
  rooms: Set<string>;
  connectedAt: number;
  lastPing: number;
}
 
class ConnectionRegistry {
  private connections = new Map<string, WebSocket>();
  private metadata = new Map<string, ConnectionMeta>();
 
  add(connectionId: string, ws: WebSocket, userId: string): void {
    this.connections.set(connectionId, ws);
    this.metadata.set(connectionId, {
      userId,
      rooms: new Set(),
      connectedAt: Date.now(),
      lastPing: Date.now(),
    });
  }
 
  remove(connectionId: string): void {
    const meta = this.metadata.get(connectionId);
    if (meta) {
      for (const room of meta.rooms) {
        this.leaveRoom(connectionId, room);
      }
    }
    this.connections.delete(connectionId);
    this.metadata.delete(connectionId);
  }
 
  getByUser(userId: string): WebSocket[] {
    const sockets: WebSocket[] = [];
    for (const [connId, meta] of this.metadata) {
      if (meta.userId === userId) {
        const ws = this.connections.get(connId);
        if (ws) sockets.push(ws);
      }
    }
    return sockets;
  }
 
  getByRoom(room: string): WebSocket[] {
    const sockets: WebSocket[] = [];
    for (const [connId, meta] of this.metadata) {
      if (meta.rooms.has(room)) {
        const ws = this.connections.get(connId);
        if (ws) sockets.push(ws);
      }
    }
    return sockets;
  }
 
  joinRoom(connectionId: string, room: string): void {
    this.metadata.get(connectionId)?.rooms.add(room);
  }
 
  leaveRoom(connectionId: string, room: string): void {
    this.metadata.get(connectionId)?.rooms.delete(room);
  }
}

Heartbeat und Erkennung toter Verbindungen

TCP meldet nicht immer, wenn eine Verbindung abbricht. Mobilfunknetze, zugeklappte Laptop-Deckel und Netzwerkwechsel beenden Verbindungen häufig unbemerkt. Ohne Heartbeats sammeln sich in deiner Registry tote Verbindungen an, die Speicher verschwenden und Sendefehler verursachen.

tstypescript
class HeartbeatManager {
  private intervals = new Map<string, NodeJS.Timeout>();
  private readonly PING_INTERVAL = 30_000;
  private readonly PONG_TIMEOUT = 10_000;
 
  start(
    connectionId: string,
    ws: WebSocket,
    onDead: (connectionId: string) => void
  ): void {
    const interval = setInterval(() => {
      if (ws.readyState !== WebSocket.OPEN) {
        this.stop(connectionId);
        onDead(connectionId);
        return;
      }
 
      let pongReceived = false;
 
      const onPong = () => {
        pongReceived = true;
      };
      ws.once("pong", onPong);
      ws.ping();
 
      setTimeout(() => {
        ws.removeListener("pong", onPong);
        if (!pongReceived) {
          ws.terminate();
          this.stop(connectionId);
          onDead(connectionId);
        }
      }, this.PONG_TIMEOUT);
    }, this.PING_INTERVAL);
 
    this.intervals.set(connectionId, interval);
  }
 
  stop(connectionId: string): void {
    const interval = this.intervals.get(connectionId);
    if (interval) {
      clearInterval(interval);
      this.intervals.delete(connectionId);
    }
  }
}

Wiederverbindung auf Client-Seite

Clients müssen Verbindungsabbrüche sauber abfangen. Exponentielles Backoff mit Jitter verhindert, dass sich nach einem Server-Neustart Tausende Clients gleichzeitig neu verbinden.

tstypescript
// ❌ Reconnect immediately in a tight loop — hammers the server
// ws.onclose = () => { connect(); };
 
// ✅ Exponential backoff with jitter
class ReconnectingWebSocket {
  private ws: WebSocket | null = null;
  private attempt = 0;
  private readonly maxDelay = 30_000;
  private readonly baseDelay = 1_000;
 
  constructor(
    private url: string,
    private onMessage: (data: string) => void
  ) {
    this.connect();
  }
 
  private connect(): void {
    this.ws = new WebSocket(this.url);
 
    this.ws.onopen = () => {
      this.attempt = 0; // Reset on successful connection
    };
 
    this.ws.onmessage = (event) => {
      this.onMessage(event.data as string);
    };
 
    this.ws.onclose = () => {
      this.scheduleReconnect();
    };
  }
 
  private scheduleReconnect(): void {
    const delay = Math.min(
      this.baseDelay * Math.pow(2, this.attempt),
      this.maxDelay
    );
    // Add jitter: random value between 0 and delay
    const jitter = Math.random() * delay;
    const finalDelay = delay + jitter;
 
    this.attempt++;
 
    setTimeout(() => this.connect(), finalDelay);
  }
 
  send(data: string): void {
    if (this.ws?.readyState === WebSocket.OPEN) {
      this.ws.send(data);
    }
  }
 
  close(): void {
    this.attempt = Infinity; // Prevent reconnection
    this.ws?.close();
  }
}

Horizontale Skalierung mit Pub/Sub

Sobald ein zweiter WebSocket-Server hinzukommt, müssen Nachrichten von einem Server auch die Clients auf dem anderen erreichen. Redis Pub/Sub oder ein vergleichbarer Message Broker schließt diese Lücke.

tstypescript
import { createClient } from "redis";
 
class ScalableMessageBroker {
  private publisher;
  private subscriber;
  private registry: ConnectionRegistry;
 
  constructor(registry: ConnectionRegistry, redisUrl: string) {
    this.registry = registry;
    this.publisher = createClient({ url: redisUrl });
    this.subscriber = createClient({ url: redisUrl });
  }
 
  async initialize(): Promise<void> {
    await this.publisher.connect();
    await this.subscriber.connect();
  }
 
  async subscribeToRoom(room: string): Promise<void> {
    await this.subscriber.subscribe(`room:${room}`, (message) => {
      // Deliver to local connections only
      const sockets = this.registry.getByRoom(room);
      for (const ws of sockets) {
        if (ws.readyState === WebSocket.OPEN) {
          ws.send(message);
        }
      }
    });
  }
 
  async publishToRoom(room: string, message: string): Promise<void> {
    // Publishes to ALL servers subscribed to this room
    await this.publisher.publish(`room:${room}`, message);
  }
 
  async publishToUser(userId: string, message: string): Promise<void> {
    await this.publisher.publish(`user:${userId}`, message);
  }
}

Design des Nachrichtenprotokolls

Rohe Strings über WebSocket werden auf Dauer unwartbar. Definiere ein typisiertes Nachrichtenprotokoll, das Versionierung und Routing unterstützt.

tstypescript
interface WsMessage<T = unknown> {
  type: string;
  payload: T;
  timestamp: number;
  correlationId?: string;
}
 
type MessageHandler<T = unknown> = (
  connectionId: string,
  payload: T
) => void | Promise<void>;
 
class MessageRouter {
  private handlers = new Map<string, MessageHandler>();
 
  on<T>(type: string, handler: MessageHandler<T>): void {
    this.handlers.set(type, handler as MessageHandler);
  }
 
  async route(connectionId: string, raw: string): Promise<void> {
    const message: WsMessage = JSON.parse(raw);
    const handler = this.handlers.get(message.type);
 
    if (!handler) {
      console.warn(`No handler for message type: ${message.type}`);
      return;
    }
 
    await handler(connectionId, message.payload);
  }
}
 
// Usage
const router = new MessageRouter();
 
router.on<{ room: string }>("join_room", (connId, payload) => {
  registry.joinRoom(connId, payload.room);
});
 
router.on<{ room: string; text: string }>(
  "chat_message",
  async (connId, payload) => {
    const message = JSON.stringify({
      type: "chat_message",
      payload: { text: payload.text, from: connId },
      timestamp: Date.now(),
    });
    await broker.publishToRoom(payload.room, message);
  }
);

Das Wichtigste in Kürze

Die Skalierung von WebSocket unterscheidet sich grundlegend von der Skalierung von HTTP, weil Verbindungen zustandsbehaftet und dauerhaft sind. Baue von Anfang an eine Verbindungsregistry, die Nutzer und Räume nachverfolgt. Implementiere Heartbeats, um tote Verbindungen zu erkennen, die TCP unbemerkt verliert.

Die Wiederverbindung auf Client-Seite muss exponentielles Backoff mit Jitter nutzen – ohne Jitter lösen Server-Neustarts eine Lawine gleichzeitiger Reconnects aus. Nutze für die horizontale Skalierung Redis Pub/Sub oder einen Message Broker, um Nachrichten zwischen den Serverinstanzen zu verbinden. Definiere frühzeitig ein typisiertes Nachrichtenprotokoll; unstrukturierte String-Nachrichten werden mit wachsender Anwendung unmöglich zu warten. Jede aufgeschobene Entscheidung – Verbindungsverfolgung, Aufräumen toter Verbindungen, serverübergreifende Kommunikation – lässt sich später nur schwerer nachrüsten, sobald Nutzer sich bereits auf das Verhalten eines einzelnen Servers verlassen.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX