构建支持数百万用户的高并发实时WebSocket聊天系统架构
深入解析面向数百万并发用户的分布式实时聊天系统架构设计:WebSocket网关集群、Redis Pub/Sub通信与Kafka消息缓冲。
使用WebSocket为100个并发用户构建聊天应用非常简单:单台Node.js或Go服务器即可轻松搞定。然而,当系统规模扩展到10万甚至数百万同时在线用户(CCU)时,由于TCP长连接的状态性(Stateful)和I/O瓶颈,单机架构会迅速耗尽系统资源而崩溃。
如何将连接在服务器1上的用户A发出的消息,瞬间精准传递给连接在服务器2上的用户B?本文将逐步拆解支持大规模并发的生产级分布式实时聊天系统架构。
1. 心智模型(ELI5):大型社区的邮局与服务柜台#
想象一个超大型社区的内部信件传递系统:
- 基础聊天应用(单服务器): 门厅里只有一位接待员。所有人都在同一个房间里说话。虽然即时,但当一万人涌入时,接待员根本应付不过来。
- 分布式聊天应用(多服务器集群): 我们开设多个接待窗口(WebSocket Gateway)。当你在1号窗口寄信给2号窗口的朋友时:
- 1号窗口把信件放入高速履带(保证信件永不丢失的Kafka)。
- 中央广播系统(Redis Pub/Sub)喊话:“B房间有新信件到达!”。
- 你的朋友所在的2号窗口收到广播,通过保持连通的长连接(WebSocket)将信件直接递到朋友手中。
2. WebSocket横向扩展的核心技术挑战#
与传统的HTTP请求(Stateless,请求完毕即断开)不同:
- 有状态长连接(Stateful): 每个WebSocket连接都在客户端与特定的网关服务器之间保持持久的TCP套接字。
- 操作系统资源限制: 每个连接都会占用文件描述符(FD)和内存(每个Socket约10KB至50KB)。普通服务器通常在5万到10万个连接时就会达到瓶颈。
- 跨节点消息路由: 在集群中,用户A和用户B很可能连接在不同的物理服务器上。服务器1必须具备将数据路由给服务器2的能力。
3. 分布式实时聊天系统整体架构#
以下是生产环境中标准的高并发分布式消息流架构图:
flowchart TD
subgraph Clients["客户端层 (Clients)"]
UserA["Client A (发送方)"]
UserB["Client B (接收方)"]
end
subgraph Edge["负载均衡与网关层"]
LB["Layer 7 负载均衡器 (Nginx / HAProxy / Envoy)"]
WS1["WebSocket 网关节点 #1"]
WS2["WebSocket 网关节点 #2"]
WSS["WebSocket 网关节点 #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 -->|一致性哈希 / 会话保持| WS1
LB -->|一致性哈希 / 会话保持| 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 网关集群 (Connection Layer)#
- 职责单一: 仅负责维持TCP/TLS握手连接、消息编解码(JSON或Protobuf)以及转发内部消息总线。
- 剥离所有繁重的业务逻辑,保持网关节点尽可能轻量与高效。
4.2. 消息广播底座: Redis Pub/Sub#
- 超低延迟(亚毫秒级): 将收到的聊天事件秒级广播给所有的网关节点,寻找目标用户的连接。
- 局限性: 即发即弃(Fire-and-forget)。若网关节点恰好重启,广播数据将丢失。因此必须配合Kafka保证可靠性。
4.3. 消息队列与异步持久化: Apache Kafka#
- 基于分布式提交日志,保障消息零丢失(Zero Message Loss)。
- 将分区键(Partition Key)设置为
room_id或conversation_id,严格保证单聊/群聊内的消息顺序(Message Ordering)。 - 后台Worker异步从Kafka消费数据批量写入数据库(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网关 + Redis Pub/Sub (Node.js & TypeScript)#
以下是分布式网关节点的最小化实现示例:
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');
// 存储当前服务器上保持的本地客户端连接
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.`);
// 登记在线状态及所在服务器
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);
// 如果目标用户正在本节点保持连接,直接通过WebSocket推送
if (targetWs && targetWs.readyState === WebSocket.OPEN) {
targetWs.send(JSON.stringify(message));
}
});
console.log(`WebSocket Gateway [${SERVER_ID}] running on port ${PORT}`);typescript6. 支持数百万连接的操作系统与基础设施调优#
单机承载数十万长连接的关键系统级优化:
- 调高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 - 心跳机制 (Ping/Pong): 客户端与服务端每隔30秒进行心跳探活,迅速识别并释放移动网络断网导致的幽灵连接。
[!TIP] 绝不要通过WebSocket传输海量历史消息!对于历史消息加载,使用标准的REST API / GraphQL并结合基于游标(Cursor-based)的分页机制,让WebSocket通道纯粹保持在轻量级的实时事件传输上。