大規模リアルタイムチャットシステムのWebSocket設計とスケール戦略
数百万人の同時接続を支える分散リアルタイムチャットシステムのアーキテクチャ設計ガイド:WebSocket Gateway、Redis Pub/Sub、Kafkaバッファの活用法。
WebSocketを使用して100人の同時接続ユーザー向けにチャットアプリを作成するのは簡単です。単一のNode.jsやGoサーバーがあれば十分に機能します。しかし、システムが10万人や数百万人の同時接続ユーザー(CCU)に拡大すると、TCP接続のステートフルな性質とI/Oのボトルネックにより、単一サーバーは即座にリソース枯渇でクラッシュします。
サーバー1に接続しているユーザーAからのメッセージを、サーバー2に接続しているユーザーBへ瞬時に届けるにはどうすればよいでしょうか?本記事では、大規模チャットシステムを支える実践的な分散アーキテクチャを解説します。
1. メンタルモデル(ELI5):郵便局と受付窓口の仕組み#
巨大な集合住宅における郵便配達をイメージしてください:
- 基本チャットアプリ(単一サーバー): 小さなロビーに受付係が1人だけいます。全員がその部屋に入って直接話します。即座に伝わりますが、1万人が一斉に入ってくると受付係はパンクします。
- 分散チャットアプリ(マルチサーバー): 複数の受付窓口(WebSocket Gateway)を設けます。窓口1で窓口2にいる友人宛てに手紙を出すと:
- 窓口1は手紙を高速コンベア(メッセージを絶対に失わないKafka)に流します。
- 館内放送(Redis Pub/Sub)で「B号室宛てに新着メッセージがあります!」と一斉通知します。
- 友人がいる窓口2が通知を受信し、開いたままの接続(WebSocket)を通じて友人に手紙を直接手渡します。
2. WebSocketをスケールさせる際の本質的課題#
HTTPプロトコル(リクエストごとに切断されるステートレスな通信)と異なり:
- ステートフルな接続: 各WebSocket接続は、クライアントと特定のサーバーインスタンス間でTCPソケットを維持し続けます。
- リソースの制約: 開かれた各ソケットはファイルディスクリプタ(FD)とメモリ(約10KB〜50KB/socket)を消費します。一般的なサーバーでは5万〜10万接続で上限に達します。
- ノード間ルーティング: 分散環境では、送信者Aと受信者Bが別々の物理サーバーに接続されていることが普通です。サーバー1はサーバー2へメッセージを中継する仕組みを必要とします。
3. 分散チャットシステムの全体アーキテクチャ#
以下はプロダクション環境で採用される分散メッセージングアーキテクチャです:
flowchart TD
subgraph Clients["クライアント層 (Clients)"]
UserA["Client A (送信者)"]
UserB["Client B (受信者)"]
end
subgraph Edge["ロードバランサ & Gateway層"]
LB["Layer 7 Load Balancer (Nginx / HAProxy / Envoy)"]
WS1["WebSocket Gateway #1"]
WS2["WebSocket Gateway #2"]
WSS["WebSocket Gateway #N"]
end
subgraph RealTime["リアルタイム調整層 (Backplane)"]
Redis[(Redis Cluster / PubSub / Dragonfly)]
Presence[(Presence & Routing Cache)]
end
subgraph EventStream["永続化 & メッセージキュー層"]
Kafka[(Apache Kafka Topic: chat-messages)]
Worker["Chat Worker / Storage Consumer"]
DB[(Primary DB: PostgreSQL / ScyllaDB)]
end
UserA -->|1. WSS 接続| LB
UserB -->|1. WSS 接続| LB
LB -->|Consistent Hash / Sticky| WS1
LB -->|Consistent Hash / Sticky| WS2
WS1 -->|2. メッセージ受信| Kafka
Kafka -->|3. 消費 & 永続化| Worker
Worker -->|履歴保存| DB
WS1 -->|4. イベント発行| Redis
Redis -->|5. 全ノード配信| WS2
WS2 -->|6. クライアントへ送信| UserB
WS1 -.->|セッション追跡| Presence
WS2 -.->|セッション追跡| Presence
4. 主要コンポーネントの設計ポイント#
4.1. WebSocket Gateway クラスタ#
- 単一責任の原則: TCP/TLS接続の維持、データの暗号化・復号(JSON / Protobuf)、内部バスへの転送に特化させます。
- 重いビジネスロジックをGatewayから完全に分離し、メモリ消費を最小限に抑えます。
4.2. メッセージバックプレーン: Redis Pub/Sub#
- 超低レイテンシ(サブミリ秒): 受信したチャットイベントを全Gatewayノードに一斉ブロードキャストし、受信者が接続しているノードを探し出します。
- 注意点: Fire-and-forget(投げっぱなし)のため、ノード再起動時のメッセージ消失を防ぐためにKafkaと併用します。
4.3. メッセージキューと永続化: Apache Kafka#
- 分散コミットログによりメッセージ損失ゼロ(Zero Message Loss)を保証。
room_idやconversation_idをパーティションキーとして設定し、チャットルーム内の順序(Message Ordering)を厳格に保持。- バックグラウンドワーカーが非同期でDB(PostgreSQL / ScyllaDB)に書き込むため、リアルタイムの配信性能に影響を与えません。
4.4. プレゼンスと接続管理 (Presence & Session State)#
Redis Hashを使用してユーザーがどのノードに接続しているかを管理します:
Key: user:session:<user_id>
Value: { "server_id": "ws-node-02", "status": "online", "last_heartbeat": 1772183200 }text5. 実装例: WebSocket Gateway + Redis Pub/Sub (Node.js & TypeScript)#
分散クラスタ内の1つのWebSocketノードの実装例です:
ws-gateway-node.ts
import { WebSocketServer, WebSocket } from 'ws';
import Redis from 'ioredis';
import { v4 as uuidv4 } from 'uuid';
const SERVER_ID = `node-${process.env.NODE_ID || uuidv4().slice(0, 8)}`;
const PORT = Number(process.env.PORT) || 8080;
const redisPub = new Redis(process.env.REDIS_URL || 'redis://localhost:6379');
const redisSub = new Redis(process.env.REDIS_URL || 'redis://localhost:6379');
// このサーバーインスタンスに接続されているローカルソケットのMap
const localConnections = new Map<string, WebSocket>();
const wss = new WebSocketServer({ port: PORT });
wss.on('connection', (ws: WebSocket, req) => {
const userId = new URL(req.url || '', `http://${req.headers.host}`).searchParams.get('userId');
if (!userId) {
ws.close(4001, 'UserId required');
return;
}
localConnections.set(userId, ws);
console.log(`[${SERVER_ID}] User ${userId} connected.`);
// Redisにプレゼンス(接続先サーバー情報)を登録
redisPub.set(`presence:${userId}`, SERVER_ID, 'EX', 60);
ws.on('message', async (data: string) => {
try {
const payload = JSON.parse(data.toString());
const { recipientId, content } = payload;
const chatEvent = {
senderId: userId,
recipientId,
content,
timestamp: Date.now(),
};
// 分散メッセージバックプレーンに発行
await redisPub.publish('chat:messages', JSON.stringify(chatEvent));
} catch (err) {
console.error('Failed to process message:', err);
}
});
ws.on('close', () => {
localConnections.delete(userId);
redisPub.del(`presence:${userId}`);
console.log(`[${SERVER_ID}] User ${userId} disconnected.`);
});
});
// 他のノードから発行されたメッセージを購読
redisSub.subscribe('chat:messages');
redisSub.on('message', (_channel, messageStr) => {
const message = JSON.parse(messageStr);
const targetWs = localConnections.get(message.recipientId);
// 受信者がこのサーバーに接続されている場合、直接送信
if (targetWs && targetWs.readyState === WebSocket.OPEN) {
targetWs.send(JSON.stringify(message));
}
});
console.log(`WebSocket Gateway [${SERVER_ID}] running on port ${PORT}`);typescript6. 数百万接続を支えるOS & インフラ最適化#
1ノードあたり数十万の同時接続を処理するための必須設定:
- Linuxファイルディスクリプタ上限の拡張:
bash# /etc/security/limits.conf * soft nofile 1048576 * hard nofile 1048576 - TCPカーネルパラメータのチューニング (
/etc/sysctl.conf):
inifs.file-max = 2097152 net.ipv4.tcp_max_syn_backlog = 65536 net.core.somaxconn = 65536 net.ipv4.tcp_rmem = 4096 87380 4194304 net.ipv4.tcp_wmem = 4096 65536 4194304 - Heartbeatプロトコル(Ping/Pong): クライアントとサーバー間で定期的なping/pong(30秒間隔など)を実施し、切断されたデッドソケットを速やかに検知・解放します。
[!TIP] 過去のチャット履歴をWebSocket経由で取得しないようにしましょう!過去ログの読み込みにはカーソルベースのページネーションを備えた通常のREST APIやGraphQLを使用し、WebSocketは純粋なリアルタイムイベント配信専用に保つのが定石です。