Node.js WebSocket 实时通信全栈实践

深入 Node.js WebSocket 实时通信:协议握手与帧格式、ws vs Socket.IO 选型对比、Room 与 Namespace 广播、Redis Adapter 多服务器扩展、JWT 认证、心跳重连、限流防刷、消息队列离线投递,以及聊天室、协同编辑等实战案例。

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是否是最后一帧
opcode0x1=文本, 0x2=二进制, 0x8=关闭, 0x9=Ping, 0xA=Pong
MASK客户端发送必须置 1,服务端发送置 0
Payload len0-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

特性wsSocket.IO
协议原生 WebSocketWebSocket + HTTP 长轮询降级
自动重连手动实现内置,支持指数退避
广播/Room手动实现 Set 管理原生支持 Room 和 Namespace
多服务器扩展手动集成 Redis Pub/Subsocket.io-redis-adapter 一键扩展
二进制传输Buffer/ArrayBufferBuffer + 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').emitNamespace 内所有 Socket
socket.to(room).emitRoom 内除当前外的 Socket
io.to(room).emitRoom 内所有 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:

参数默认值说明
pingInterval25000ms服务端发送 Ping 间隔
pingTimeout20000msPing 超时时间
connectTimeout20000ms连接超时
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 方案
AWSAPI Gateway WebSocket API + Lambda
Vercel不支持原生 WebSocket;使用 Server-Sent Events 或第三方服务(Pusher、Ably)
Cloudflare WorkersDurable Objects 支持有状态连接
FirebaseCloud 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 光标
ServerlessAPI Gateway WebSocket / SSE / 第三方服务

WebSocket 实时通信看似简单,但要在生产环境中稳定运行,需要在连接管理、状态同步、故障恢复、安全防护等领域做大量工程投入。从单服务器的 Socket.IO 快速原型,到多节点 Redis Adapter 集群扩展,再到离线消息队列和限流防刷,每一步都是保障用户体验的必经之路。


延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「nodejs」更多文章

  1. Node.js ORM 深度对比:Prisma、TypeORM、Sequelize 与 Drizzle
  2. Node.js 设计模式与最佳实践:从 SOLID 到六边形架构
  3. Node.js 高级测试策略:从单元测试到混沌工程的完整实践