HTTP 协议自诞生以来一直是 Web 通信的基石,其请求-响应模型让浏览器与服务器之间的数据交换变得简单而有序。然而,这种模型的本质决定了服务端无法主动向客户端推送数据——每一次数据传递都必须由客户端发起请求,服务端被动响应。在实时性要求越来越高的今天,这种模式显然力不从心。直播弹幕、股票即时报价、多人文档协同编辑、在线客服系统、物联网设备上报等场景,无一不需要服务端在数据产生的第一时间主动触达客户端。WebSocket 正是在这样的背景下被设计出来的全双工通信协议,它在 TCP 之上定义了一套轻量的帧格式和握手规范,使得一个持久连接即可承载双向实时数据流。
Node.js 凭借事件驱动的非阻塞 I/O 模型,天然适合处理海量并发长连接,成为构建 WebSocket 服务端的首选平台之一。本文从协议底层原理出发,逐步深入到工程实践:对比 ws 与 Socket.IO 的选型维度,剖析 Socket.IO 的 Room 与 Namespace 广播机制,演示 Redis Adapter 多服务器横向扩展方案,设计完整的 JWT 握手认证与 Token 刷新策略,实现可靠的心跳检测与断线重连,构建多层级的限流与消息防刷体系,搭建基于 Redis 的离线消息队列,最后通过聊天应用、实时通知和协同编辑三个案例串联所有知识点,并探讨 Serverless 环境下 WebSocket 的限制与替代方案。
1. WebSocket 协议原理
1.1 握手升级
WebSocket 连接始于一次 HTTP 升级请求。客户端发送标准 HTTP 请求,携带特殊的 Upgrade 头部:
GET /chat HTTP/1.1
Host: server.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Sec-WebSocket-Version: 13
服务端验证 Sec-WebSocket-Key,返回 101 Switching Protocols 响应:
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
握手完成后,TCP 连接保持不变,后续通信切换为 WebSocket 帧格式,双方可在同一连接上随时发送数据。值得注意的是,WebSocket 建立过程中复用了 HTTP 端口(通常是 80 或 443),这意味着不需要额外开放端口,防火墙穿透能力远优于裸 TCP 长连接。
1.2 帧格式解析
与 HTTP 使用文本头部加可选 Body 不同,WebSocket 采用二进制帧传输数据,每一帧都包含完整的数据类型、长度和传输状态信息。这种设计让协议本身非常轻量,在传输小消息时几乎没有协议开销。帧头部最小 2 字节,最大 14 字节:
0 1 2 3
0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-------+-+-------------+-------------------------------+
|F|R|R|R| opcode|M| Payload len | Extended payload length |
|I|S|S|S| (4) |A| (7) | (16/64) |
|N|V|V|V| |S| | (if payload len==126/127) |
| |1|2|3| |K| | |
+-+-+-+-+-------+-+-------------+ - - - - - - - - - - - - - - - +
| Extended payload length continued, if payload len == 127 |
+ - - - - - - - - - - - - - - - +-------------------------------+
| |Masking-key, if MASK set to 1 |
+-------------------------------+-------------------------------+
| Masking-key (continued) | Payload Data |
+-------------------------------- - - - - - - - - - - - - - - -+
: Payload Data continued ... :
+ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - +
| Payload Data continued ... |
+---------------------------------------------------------------+
| 字段 | 说明 |
|---|---|
FIN | 是否是最后一帧 |
opcode | 0x1=文本, 0x2=二进制, 0x8=关闭, 0x9=Ping, 0xA=Pong |
MASK | 客户端发送必须置 1,服务端发送置 0 |
Payload len | 0-125 直接表示长度;126 表示后续 2 字节;127 表示后续 8 字节 |
Masking-key | 客户端随机生成的 4 字节掩码,用于异或混淆 payload |
1.3 Ping/Pong 保活机制
WebSocket 内置心跳检测:一端发送 Ping 帧(opcode 0x9),另一端必须回复 Pong 帧(opcode 0xA)。Ping/Pong 可携带最多 125 字节的 payload,通常用于传递服务端时间戳或健康状态。心跳机制不仅用于检测连接是否存活,还能及时发现网络异常,避免因中间代理超时断开而导致的数据丢失。在生产环境中,合理设置心跳间隔至关重要——间隔过短会增加不必要的网络开销长连接资源消耗,间隔过长则可能导致僵尸连接长时间占用服务器资源。
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
wss.on('connection', (ws) => {
// 每 30 秒发送一次 Ping
const pingInterval = setInterval(() => {
if (ws.readyState === WebSocket.OPEN) {
ws.ping(Date.now().toString());
}
}, 30000);
ws.on('pong', (data) => {
const pingTime = parseInt(data.toString());
console.log(`RTT: ${Date.now() - pingTime}ms`);
});
ws.on('close', () => clearInterval(pingInterval));
});
2. ws vs Socket.IO 选型对比
Node.js 生态中有两个主流的 WebSocket 库:原生轻量的 ws 和功能丰富的 Socket.IO。
| 特性 | ws | Socket.IO |
|---|---|---|
| 协议 | 原生 WebSocket | WebSocket + HTTP 长轮询降级 |
| 自动重连 | 手动实现 | 内置,支持指数退避 |
| 广播/Room | 手动实现 Set 管理 | 原生支持 Room 和 Namespace |
| 多服务器扩展 | 手动集成 Redis Pub/Sub | socket.io-redis-adapter 一键扩展 |
| 二进制传输 | Buffer/ArrayBuffer | Buffer + Base64 兼容层 |
| 中间件 | 需自行封装 | 类似 Express 的中间件体系 |
| Payload 大小限制 | 默认无限 | 默认 1MB,可配置 |
| 学习成本 | 低,接近协议层 | 中,封装了更多抽象 |
| 包体积(客户端) | 极小 | 约 30KB gzip |
| 适用场景 | IoT、游戏、高并发推送 | 聊天应用、实时协作、通知 |
2.1 ws 库快速上手
const WebSocket = require('ws');
// 服务端
const wss = new WebSocket.Server({ port: 8080 });
const clients = new Set();
wss.on('connection', (ws, req) => {
clients.add(ws);
console.log(`Connected: ${req.socket.remoteAddress}`);
ws.on('message', (data) => {
// 广播给所有客户端
for (const client of clients) {
if (client.readyState === WebSocket.OPEN) {
client.send(data.toString());
}
}
});
ws.on('close', () => clients.delete(ws));
});
// 客户端
const client = new WebSocket('ws://localhost:8080');
client.on('open', () => client.send('Hello Server'));
client.on('message', (data) => console.log('Received:', data.toString()));
2.2 Socket.IO 快速上手
// 服务端
const { Server } = require('socket.io');
const io = new Server(3000, { cors: { origin: '*' } });
io.on('connection', (socket) => {
console.log('User connected:', socket.id);
socket.on('chat message', (msg) => {
io.emit('chat message', { user: socket.id, text: msg });
});
socket.on('disconnect', () => {
console.log('User disconnected:', socket.id);
});
});
// 客户端
import { io } from 'socket.io-client';
const socket = io('http://localhost:3000');
socket.on('connect', () => console.log('Connected:', socket.id));
socket.on('chat message', (data) => console.log(data.text));
socket.emit('chat message', 'Hello World');
选型建议:如果项目需要极致的实时性、海量并发连接(如高频量化交易、MMO 游戏后端、IoT 设备网关),且团队有能力从零实现完整的重连、广播和负载均衡逻辑,选择
ws最为合适;对于大多数常规业务场景(即时聊天、系统通知、在线协作、实时数据看板),直接选择Socket.IO能够显著缩短开发周期,并自动获得协议降级和多服务器扩展能力。
3. Socket.IO Rooms、Namespaces 与广播
Socket.IO 提供了强大的频道管理能力,是实现多房间聊天、权限隔离的核心工具。
3.1 Room(房间)
Room 是 Socket 的逻辑分组,一个 Socket 可以加入多个 Room。消息可按 Room 精确投递。
io.on('connection', (socket) => {
// 加入房间
socket.join('room-42');
socket.join('room-premium');
// 向房间广播(不包含发送者)
socket.to('room-42').emit('message', 'Hello room!');
// 向房间广播(包含发送者)
io.to('room-42').emit('announcement', 'Server update');
// 向多个房间广播
socket.to('room-42').to('room-43').emit('message', 'Multi-room');
// 离开房间
socket.leave('room-42');
// 获取房间列表
const rooms = Array.from(socket.rooms);
console.log('My rooms:', rooms); // 包含默认的 socket.id 房间
// 获取房间内所有 Socket ID
const roomSockets = io.sockets.adapter.rooms.get('room-42');
console.log('Room size:', roomSockets?.size || 0);
});
3.2 Namespace(命名空间)
Namespace 用于逻辑隔离不同的业务线或权限域,创建独立的 Socket 实例:
// /chat 命名空间
const chatNs = io.of('/chat');
chatNs.use(authMiddleware); // Namespace 级别中间件
chatNs.on('connection', (socket) => {
socket.on('send', (msg) => {
chatNs.emit('broadcast', msg);
});
});
// /admin 命名空间
const adminNs = io.of('/admin');
adminNs.use(adminAuthMiddleware);
adminNs.on('connection', (socket) => {
socket.on('notify-all', (msg) => {
io.emit('global-notify', msg); // 跨 Namespace 广播需使用根 io
});
});
// 客户端连接指定 Namespace
const chatSocket = io('http://localhost:3000/chat');
const adminSocket = io('http://localhost:3000/admin');
3.3 广播模式总结
| 方法 | 目标范围 |
|---|---|
socket.emit | 仅当前 Socket |
socket.broadcast.emit | 除当前外的所有 Socket |
io.emit / io.of('/ns').emit | Namespace 内所有 Socket |
socket.to(room).emit | Room 内除当前外的 Socket |
io.to(room).emit | Room 内所有 Socket |
socket.except(room).emit | 排除指定 Room |
socket.volatile.emit | 不保证发送(丢弃比阻塞更重要时) |
socket.compress(true).emit | 启用 permessage-deflate 压缩 |
4. Redis Adapter 多服务器扩展
单台服务器的 Socket.IO 进程受限于单机的 CPU 核心数和内存容量,通常最多支撑数万并发连接。当用户规模扩大、消息吞吐量激增时,必须横向扩展为服务器集群。然而,WebSocket 连接是有状态的——每个客户端都固定连接到某一台服务器的某个进程上,如果仅仅在多台机器上各跑一个 Socket.IO 实例,那么服务器 A 上的客户端发送的消息,服务器 B 和 C 上的客户端将无法收到,因为广播事件只发生在本地内存中。Redis Adapter 正是为解决这一分布式难题而设计的:它将 Socket.IO 的内存内事件总线替换为 Redis 的发布订阅通道,实现跨服务器实例的广播同步。
4.1 架构原理
┌─────────┐ ┌─────────┐ ┌─────────┐
│ Server 1│ │ Server 2│ │ Server 3│
│ :3001 │ │ :3002 │ │ :3003 │
└────┬────┘ └────┬────┘ └────┬────┘
│ │ │
└───────────────┼───────────────┘
│
┌──────▼──────┐
│ Redis │
│ Pub/Sub │
└─────────────┘
每台服务器将本地广播事件发布到 Redis Channel,其他服务器订阅该 Channel,在本地 Socket 中执行广播。
4.2 完整集成代码
const { createServer } = require('http');
const { Server } = require('socket.io');
const { createAdapter } = require('@socket.io/redis-adapter');
const { createClient } = require('redis');
const express = require('express');
const app = express();
const server = createServer(app);
const io = new Server(server, {
cors: { origin: process.env.CLIENT_ORIGIN || 'http://localhost:5173' },
});
// 创建 Redis Pub/Sub 客户端
async function setupAdapter() {
const pubClient = createClient({ url: process.env.REDIS_URL || 'redis://localhost:6379' });
const subClient = pubClient.duplicate();
await Promise.all([pubClient.connect(), subClient.connect()]);
io.adapter(createAdapter(pubClient, subClient));
console.log('Redis adapter connected');
}
setupAdapter().catch(console.error);
// 房间管理中间件
io.use(async (socket, next) => {
const token = socket.handshake.auth.token;
try {
const user = await validateToken(token); // JWT 验证
socket.user = user;
next();
} catch (err) {
next(new Error('Authentication error'));
}
});
io.on('connection', (socket) => {
console.log(`[${process.pid}] User connected: ${socket.user.id}`);
// 加入用户专属房间(用于一对一消息)
socket.join(`user:${socket.user.id}`);
// 加入公开聊天室
socket.on('join-room', (roomId) => {
socket.join(roomId);
socket.to(roomId).emit('user-joined', {
userId: socket.user.id,
username: socket.user.name,
});
});
// 发送消息
socket.on('send-message', ({ roomId, content }) => {
const message = {
id: generateId(),
roomId,
userId: socket.user.id,
username: socket.user.name,
content,
timestamp: new Date().toISOString(),
};
// 持久化到数据库(异步,不阻塞广播)
saveMessage(message).catch(console.error);
// 广播给房间内所有人(跨服务器)
io.to(roomId).emit('new-message', message);
});
// 一对一私聊
socket.on('private-message', ({ toUserId, content }) => {
const message = {
id: generateId(),
from: socket.user.id,
to: toUserId,
content,
timestamp: new Date().toISOString(),
};
// 目标用户可能在任意服务器上,Redis Adapter 确保送达
io.to(`user:${toUserId}`).emit('private-message', message);
savePrivateMessage(message).catch(console.error);
});
// 离开房间
socket.on('leave-room', (roomId) => {
socket.leave(roomId);
socket.to(roomId).emit('user-left', { userId: socket.user.id });
});
socket.on('disconnect', () => {
console.log(`[${process.pid}] User disconnected: ${socket.user.id}`);
});
});
// 获取房间成员(跨服务器接口)
app.get('/rooms/:roomId/members', async (req, res) => {
const sockets = await io.in(req.params.roomId).fetchSockets();
const members = sockets.map((s) => ({
id: s.user.id,
name: s.user.name,
}));
res.json({ count: members.length, members });
});
const PORT = process.env.PORT || 3000;
server.listen(PORT, () => console.log(`Server ${process.pid} on port ${PORT}`));
4.3 负载均衡配置
Nginx 需要配置 ip_hash 或基于 Cookie 的会话亲和性,确保同一客户端始终路由到同一服务器:
upstream websocket_backend {
ip_hash; # 或 hash $cookie_io consistent;
server 127.0.0.1:3001;
server 127.0.0.1:3002;
server 127.0.0.1:3003;
}
server {
listen 80;
location /socket.io/ {
proxy_pass http://websocket_backend;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_read_timeout 86400;
}
}
5. WebSocket 认证机制
与 HTTP 请求的每一次交互都可以通过 Cookie 或 Authorization 头部携带身份凭证不同,WebSocket 协议本身并未定义认证机制。这意味着身份验证需要在握手阶段或连接建立后另行实现。在握手阶段拒绝非法连接是最理想的方案,因为一旦长连接建立,维持连接的内存和 CPU 成本远高于一次 HTTP 请求,让未认证用户占用连接资源是一种明显的浪费。
5.1 握手阶段传递 Token
// 服务端中间件
io.use(async (socket, next) => {
const token = socket.handshake.auth.token || socket.handshake.query.token;
if (!token) {
return next(new Error('Missing token'));
}
try {
const decoded = jwt.verify(token, process.env.JWT_SECRET);
const user = await User.findById(decoded.userId);
if (!user || !user.isActive) {
return next(new Error('Invalid user'));
}
socket.user = user;
next();
} catch (err) {
next(new Error('Authentication failed'));
}
});
// 客户端
const socket = io('http://localhost:3000', {
auth: { token: localStorage.getItem('accessToken') }
});
5.2 Token 过期刷新策略
// 连接错误处理
socket.on('connect_error', async (err) => {
if (err.message === 'Authentication failed') {
// 尝试刷新 Token
const newToken = await refreshAccessToken();
localStorage.setItem('accessToken', newToken);
// 使用新 Token 重连
socket.auth.token = newToken;
socket.connect();
}
});
// 服务端主动断开过期连接后刷新
socket.on('disconnect', (reason) => {
if (reason === 'io server disconnect') {
// 服务端强制断开,尝试重连并刷新令牌
refreshAndReconnect();
}
});
5.3 Session 认证方案
对于基于 Session 的应用(如 Express + connect-redis),可通过共享 Cookie 验证:
const session = require('express-session');
const sessionMiddleware = session({
store: new RedisStore({ client: redisClient }),
secret: process.env.SESSION_SECRET,
resave: false,
saveUninitialized: false,
});
// Socket.IO 共享 Express Session
io.use((socket, next) => {
sessionMiddleware(socket.request, socket.request.res || {}, next);
});
io.use((socket, next) => {
if (socket.request.session?.userId) {
socket.userId = socket.request.session.userId;
next();
} else {
next(new Error('Unauthorized'));
}
});
6. 心跳与重连策略
6.1 Socket.IO 内置心跳
Socket.IO 引擎层自动处理 Ping/Pong:
| 参数 | 默认值 | 说明 |
|---|---|---|
pingInterval | 25000ms | 服务端发送 Ping 间隔 |
pingTimeout | 20000ms | Ping 超时时间 |
connectTimeout | 20000ms | 连接超时 |
const io = new Server(server, {
pingInterval: 10000, // 10 秒一次心跳
pingTimeout: 5000, // 5 秒未收到 Pong 则断开
});
6.2 自定义重连策略
const socket = io('http://localhost:3000', {
reconnection: true,
reconnectionAttempts: 10, // 最大重连次数
reconnectionDelay: 1000, // 初始延迟 1 秒
reconnectionDelayMax: 5000, // 最大延迟 5 秒
randomizationFactor: 0.5, // 抖动因子
});
// 监听重连事件
socket.on('reconnect', (attempt) => {
console.log(`Reconnected after ${attempt} attempts`);
// 重新加入房间、同步未读消息
socket.emit('join-room', currentRoomId);
socket.emit('sync-messages', { roomId: currentRoomId, since: lastMessageTime });
});
socket.on('reconnect_error', (err) => {
console.error('Reconnection failed:', err.message);
});
socket.on('reconnect_failed', () => {
// 所有重连尝试均失败
showOfflineMode();
});
6.3 连接状态完整生命周期
首次连接: disconnected → connecting → connect →connected
网络抖动: connected → disconnect → reconnecting → reconnect → connected
服务端拒绝: connecting → connect_error → reconnecting → ...
无限重连失败: ... → reconnect_failed → disconnected
7. 限流与防刷措施
WebSocket 连接一旦建立,客户端可以在没有任何预热的情况下持续发送消息。与 HTTP 请求相比,WebSocket 省略了 TCP 握手和 TLS 协商的开销,攻击者能够以更低的成本产生更高的请求密度。一个精心编写的恶意脚本可以在一秒内向服务端发送数千条消息,如果每条消息都触发数据库写入、Redis 广播和业务逻辑处理,很快就会耗尽服务器资源,导致正常用户的消息被延迟甚至丢包。因此,在连接层之上构建消息层级的限流和防刷机制,是保障服务稳定性的必要防线。
7.1 消息速率限制
基于滑动窗口实现每用户每秒消息数限制:
const rateLimit = require('ws-rate-limit')('40s', 20); // 40 秒内最多 20 条
// 或使用自定义实现
class RateLimiter {
constructor(maxMessages = 20, windowMs = 40000) {
this.max = maxMessages;
this.window = windowMs;
this.clients = new Map();
}
isAllowed(clientId) {
const now = Date.now();
const history = this.clients.get(clientId) || [];
// 清理过期记录
const valid = history.filter((t) => now - t < this.window);
this.clients.set(clientId, valid);
if (valid.length >= this.max) {
return false;
}
valid.push(now);
return true;
}
reset(clientId) {
this.clients.delete(clientId);
}
}
const limiter = new RateLimiter(30, 60000); // 每分钟 30 条
io.on('connection', (socket) => {
socket.on('send-message', (data) => {
if (!limiter.isAllowed(socket.user.id)) {
socket.emit('error', { code: 'RATE_LIMITED', message: '消息发送过于频繁' });
return;
}
// 处理消息...
});
});
7.2 连接数限制
// 每 IP 最大连接数
const connectionsPerIP = new Map();
io.on('connection', (socket) => {
const ip = socket.handshake.address || socket.request.socket.remoteAddress;
const count = (connectionsPerIP.get(ip) || 0) + 1;
if (count > 5) {
socket.disconnect(true);
return;
}
connectionsPerIP.set(ip, count);
socket.on('disconnect', () => {
const current = connectionsPerIP.get(ip) || 1;
if (current <= 1) connectionsPerIP.delete(ip);
else connectionsPerIP.set(ip, current - 1);
});
});
7.3 消息大小与格式验证
const Ajv = require('ajv');
const ajv = new Ajv();
const messageSchema = {
type: 'object',
properties: {
roomId: { type: 'string', maxLength: 50 },
content: { type: 'string', minLength: 1, maxLength: 2000 },
},
required: ['roomId', 'content'],
additionalProperties: false,
};
const validate = ajv.compile(messageSchema);
socket.on('send-message', (data) => {
if (!validate(data)) {
socket.emit('validation-error', validate.errors);
return;
}
// 业务处理...
});
8. 离线用户消息队列
在即时通信系统中,消息投递必须做到可靠到达,不能因为用户短暂离线就造成消息丢失。然而 WebSocket 连接天然是瞬时的——用户关闭浏览器、切换网络或者杀死 App 进程都会导致连接中断。此时如果服务端直接向该用户的专属房间广播消息,由于没有任何活跃连接订阅该房间,消息将直接丢失。解决这一问题的标准方案是引入离线消息队列:用户离线时将消息暂存,用户再次上线时按序拉取并投递。这套机制是现代聊天应用(如微信、Telegram、Slack)消息可达性的核心保障。
8.1 Redis List 实现离线队列
const redis = require('ioredis');
const redisClient = new Redis(process.env.REDIS_URL);
async function queueOfflineMessage(userId, message) {
const key = `offline:${userId}`;
await redisClient.rpush(key, JSON.stringify({
...message,
queuedAt: new Date().toISOString(),
}));
// 设置过期时间(30 天)
await redisClient.expire(key, 60 * 60 * 24 * 30);
}
async function getOfflineMessages(userId) {
const key = `offline:${userId}`;
const messages = await redisClient.lrange(key, 0, -1);
await redisClient.del(key); // 消费后删除
return messages.map((m) => JSON.parse(m));
}
// 发送消息时检查用户在线状态
async function sendOrQueueMessage(toUserId, message) {
const sockets = await io.in(`user:${toUserId}`).fetchSockets();
if (sockets.length > 0) {
// 用户在线,直接投递
io.to(`user:${toUserId}`).emit('private-message', message);
} else {
// 用户离线,入队
await queueOfflineMessage(toUserId, message);
}
}
// 用户上线时拉取离线消息
io.on('connection', async (socket) => {
const offlineMessages = await getOfflineMessages(socket.user.id);
if (offlineMessages.length > 0) {
socket.emit('offline-messages', offlineMessages);
}
});
8.2 消息确认与幂等性
// 客户端收到消息后发送 ack
socket.on('private-message', async (msg) => {
displayMessage(msg);
socket.emit('ack', { messageId: msg.id });
});
// 服务端记录已送达
socket.on('ack', async ({ messageId }) => {
await markAsDelivered(messageId);
});
// 投递时防止重复(幂等键)
async function deliverOnce(messageId, userId, payload) {
const lockKey = `deliver:${messageId}:${userId}`;
const locked = await redisClient.set(lockKey, '1', 'EX', 3600, 'NX');
if (locked === 'OK') {
io.to(`user:${userId}`).emit('private-message', payload);
}
}
9. 实战案例
9.1 聊天应用完整架构
浏览器客户端 Nginx 负载均衡 Node.js 服务器集群
┌─────────┐ ┌──────────────┐ ┌──────┐ ┌──────┐ ┌──────┐
│ Socket │─── WebSocket ────→│ ip_hash │────────→│ 3001 │ │ 3002 │ │ 3003 │
│ Client │ │ /socket.io │ └──┬───┘ └──┬───┘ └──┬───┘
└─────────┘ └──────────────┘ │ │ │
└────────┼────────┘
│
┌──────▼──────┐
│ Redis │
│ Pub/Sub │
└──────┬──────┘
│
┌──────▼──────┐
│ MongoDB │
│ 消息持久化 │
└─────────────┘
9.2 实时通知系统
// 通知服务模块
class NotificationService {
constructor(io) {
this.io = io;
}
async notify(userIds, notification) {
const message = {
id: generateId(),
type: notification.type,
title: notification.title,
body: notification.body,
data: notification.data,
createdAt: new Date().toISOString(),
read: false,
};
// 持久化
await Notification.insertMany(
userIds.map((userId) => ({ ...message, userId }))
);
// 实时推送
for (const userId of userIds) {
const sockets = await this.io.in(`user:${userId}`).fetchSockets();
if (sockets.length > 0) {
this.io.to(`user:${userId}`).emit('notification', message);
} else {
// 离线:推送至手机(FCM/APNs)或邮件队列
await pushToMobile(userId, message);
}
}
}
// 广播全局公告
async broadcastAnnouncement(announcement) {
await GlobalAnnouncement.create(announcement);
this.io.emit('announcement', announcement);
}
}
9.3 协同编辑(Operational Transformation)
// 简化版协同编辑:操作序列广播
const ot = require('ot.js');
io.of('/editor').on('connection', (socket) => {
socket.on('join-doc', (docId) => {
socket.join(`doc:${docId}`);
// 发送当前文档状态
const snapshot = getDocumentSnapshot(docId);
socket.emit('doc-state', snapshot);
});
socket.on('operation', ({ docId, operation, revision }) => {
// OT 变换:将用户操作与服务端历史合并
const transformed = ot.type.transform(operation, getHistorySince(docId, revision));
// 应用操作
applyOperation(docId, transformed);
// 广播变换后的操作(其他客户端据此更新)
socket.to(`doc:${docId}`).emit('operation', {
operation: transformed,
authorId: socket.user.id,
});
});
// 光标位置同步(非关键操作,使用 volatile)
socket.on('cursor-move', ({ docId, position }) => {
socket.volatile.to(`doc:${docId}`).emit('cursor-update', {
userId: socket.user.id,
username: socket.user.name,
position,
});
});
});
10. Serverless 下的 WebSocket
10.1 无状态架构的限制
传统 Serverless(AWS Lambda、Vercel Functions、Cloudflare Workers)无法保持长连接,因为函数实例在请求结束后销毁。WebSocket 需要专门的托管服务:
| 平台 | WebSocket 方案 |
|---|---|
| AWS | API Gateway WebSocket API + Lambda |
| Vercel | 不支持原生 WebSocket;使用 Server-Sent Events 或第三方服务(Pusher、Ably) |
| Cloudflare Workers | Durable Objects 支持有状态连接 |
| Firebase | Cloud Functions + Firebase Realtime Database 监听 |
10.2 AWS API Gateway WebSocket 示例
// Lambda 处理 WebSocket 事件
exports.handler = async (event) => {
const { routeKey, connectionId, body } = event.requestContext;
switch (routeKey) {
case '$connect':
// 连接建立,存储 connectionId 与 userId 映射
await DynamoDB.put({
TableName: 'Connections',
Item: {
connectionId,
userId: event.queryStringParameters.userId,
connectedAt: new Date().toISOString(),
},
});
return { statusCode: 200 };
case '$disconnect':
await DynamoDB.delete({
TableName: 'Connections',
Key: { connectionId },
});
return { statusCode: 200 };
case 'sendMessage':
const { roomId, content } = JSON.parse(body);
// 查询房间所有 connectionId
const connections = await getRoomConnections(roomId);
// 向所有连接推送消息
await Promise.all(connections.map((conn) =>
postToConnection(conn.connectionId, { roomId, content })
));
return { statusCode: 200 };
default:
return { statusCode: 400 };
}
};
// 使用 AWS SDK 向指定连接发送消息
const { ApiGatewayManagementApiClient, PostToConnectionCommand } = require('@aws-sdk/client-apigatewaymanagementapi');
const client = new ApiGatewayManagementApiClient({
endpoint: process.env.WEBSOCKET_ENDPOINT,
});
async function postToConnection(connectionId, data) {
await client.send(new PostToConnectionCommand({
ConnectionId: connectionId,
Data: Buffer.from(JSON.stringify(data)),
}));
}
10.3 替代方案:Server-Sent Events(SSE)
对于单向推送场景(服务端 → 客户端),如股票行情推送、新闻订阅、系统监控告警等不需要客户端频繁回传数据的业务,SSE 是比 WebSocket 更轻量的选择。SSE 建立在标准 HTTP 之上,自动利用浏览器原生重连机制(EventSource 会在连接断开三秒后自动尝试重连),且天然支持断点续传(通过 Last-Event-ID 头部恢复事件流),在实现成本和运维复杂度上都低于全双工的 WebSocket。
// 服务端 SSE
app.get('/events', (req, res) => {
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
const clientId = req.query.clientId;
const send = (data) => res.write(`data: ${JSON.stringify(data)}\n\n`);
eventEmitter.on(`update:${clientId}`, send);
req.on('close', () => {
eventEmitter.off(`update:${clientId}`, send);
});
});
// 客户端
const es = new EventSource('/events?clientId=123');
es.onmessage = (e) => console.log(JSON.parse(e.data));
总结
| 维度 | 推荐方案 |
|---|---|
| 基础通信 | Socket.IO(含自动降级) |
| 多服务器扩展 | Redis Adapter + ip_hash 负载均衡 |
| 认证 | 握手阶段 JWT Token,配合 Refresh Token |
| 心跳 | Socket.IO 内置 + 自定义业务心跳 |
| 限流 | 滑动窗口 + IP 连接数限制 + Schema 校验 |
| 离线消息 | Redis List 队列 + 上线拉取机制 |
| 大型聊天室 | Room 隔离 + 消息持久化 + 分页拉取历史 |
| 协同编辑 | Socket.IO Namespace + OT 算法 + Volatile 光标 |
| Serverless | API Gateway WebSocket / SSE / 第三方服务 |
WebSocket 实时通信看似简单,但要在生产环境中稳定运行,需要在连接管理、状态同步、故障恢复、安全防护等领域做大量工程投入。从单服务器的 Socket.IO 快速原型,到多节点 Redis Adapter 集群扩展,再到离线消息队列和限流防刷,每一步都是保障用户体验的必经之路。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。