Socket.IO Cluster Engine 实战指南:无需粘性会话的多进程横向扩展方案
【免费下载链接】socket.ioBidirectional and low-latency communication for every platform项目地址: https://gitcode.com/gh_mirrors/so/socket.io
导读
本文介绍 Socket.IO 官方仓库中的@socket.io/cluster-engine包——一个集群友好的 Engine.IO 引擎,它让 Socket.IO 服务可以在多个 Node.js 进程(甚至多台服务器)之间共享负载,完全不需要配置粘性会话(sticky sessions)。读完本文,你将掌握三种部署形态(单机 cluster、多机 Redis、两者组合)的完整接入代码,理解其底层的分布式锁与延迟连接机制,并能通过配置项精确调优升级与响应时序。
背景:为什么需要 cluster engine?
传统的 Socket.IO 水平扩展面临一个经典困境:Engine.IO 会话(以sid标识)一旦建立,后续请求必须始终路由到同一进程,否则会因 "unknown sid" 而失败。为此,常见方案是启用反向代理(如 Nginx)的 sticky sessions / IP hash。但这会带来负载不均、故障转移复杂等问题。
@socket.io/cluster-engine的定位(见 package.json 描述)是 "A cluster-friendly engine to share load between multiple Node.js processes (without sticky sessions)":它扩展了engine.io包自带的Server,让各进程通过IPC 通道(Node.js cluster 模式)或Redis Pub/Sub(多机模式)互相协作——当某个请求的会话不在本进程时,先向集群广播查询会话归属,再把数据包转发给持有该会话的 worker,从而彻底摆脱对粘性会话的依赖。
安装
npm i @socket.io/cluster-engine该包发布在 npm 上,名为@socket.io/cluster-engine(当前仓库版本 0.1.1)。它的运行时依赖包括:
engine.io ~6.6.0:底层传输层,ClusterEngine即继承自它的Server;engine.io-parser ~5.2.3:Engine.IO 数据包编解码;@msgpack/msgpack ~2.8.0:Redis 模式下集群消息的序列化格式;debug ~4.4.1:调试日志(可通过DEBUG=engine:*环境变量开启)。
包对外导出的 API 集中在 lib/index.ts:NodeClusterEngine、setupPrimary、RedisEngine、setupPrimaryWithRedis以及类型ClusterEngineOptions。
用法一:Node.js cluster(单机多进程)
核心思路:主进程(primary)fork多个 worker,所有 worker 共享同一个监听端口(由内核负载均衡分配连接);setupPrimary()在主进程中建立消息路由,worker 内实例化NodeClusterEngine并通过 IPC 与主进程交换集群消息。
import cluster from "node:cluster"; import process from "node:process"; import { availableParallelism } from "node:os"; import { setupPrimary, NodeClusterEngine } from "@socket.io/cluster-engine"; import { createServer } from "node:http"; import { Server } from "socket.io"; if (cluster.isPrimary) { console.log(`Primary ${process.pid} is running`); const numCPUs = availableParallelism(); // fork workers for (let i = 0; i < numCPUs; i++) { cluster.fork(); } // setup connection within the cluster setupPrimary(); // needed for packets containing Buffer objects (you can ignore it if you only send plaintext objects) cluster.setupPrimary({ serialization: "advanced", }); cluster.on("exit", (worker, code, signal) => { console.log(`worker ${worker.process.pid} died`); }); } else { const httpServer = createServer((req, res) => { res.writeHead(404).end(); }); const engine = new NodeClusterEngine(); engine.attach(httpServer, { path: "/socket.io/" }); const io = new Server(); io.bind(engine); // workers will share the same port httpServer.listen(3000); console.log(`Worker ${process.pid} started`); }要点说明:
- 端口共享:
httpServer.listen(3000)写在每个 worker 里,多个 worker 共享同一端口,新连接由操作系统分发给不同 worker——这就是负载均衡的来源; setupPrimary()必须且只需在主进程调用一次:它监听cluster.on("message"),根据消息中的recipientId定向转发给目标 worker,没有recipientId的广播类消息则转发给除发送方外的所有 worker(见 lib/cluster.ts);serialization: "advanced":当业务涉及发送Buffer(二进制)数据包时,主进程转发的消息可能携带 Buffer,需要开启 advanced 序列化;只收发纯文本对象时可省略;io.bind(engine):把自定义引擎绑定到 Socket.IOServer上,替代默认的 engine.io Server。
仓库 test/worker.js 展示了一个更简化的 worker 形态:直接对engine监听connection事件回显消息,其测试 test/cluster.test.ts 用 3 个 worker 验证了跨进程 ping/pong 与二进制收发。
用法二:Redis(多机横向扩展)
当进程分布在多台服务器上、无法共用 IPC 时,改用RedisEngine:集群消息通过 Redis Pub/Sub 广播,@msgpack/msgpack负责序列化。该模式不依赖node:cluster,每个 Node.js 进程都是独立实例,前置一个普通轮询/随机负载均衡器即可(无需粘性会话)。
import { createServer } from "node:http"; import { createClient } from "redis"; import { RedisEngine } from "@socket.io/cluster-engine"; import { Server } from "socket.io"; const httpServer = createServer((req, res) => { res.writeHead(404).end(); }); const pubClient = createClient(); const subClient = pubClient.duplicate(); await Promise.all([ pubClient.connect(), subClient.connect(), ]); const engine = new RedisEngine(pubClient, subClient); engine.attach(httpServer, { path: "/socket.io/" }); const io = new Server(); io.bind(engine); httpServer.listen(3000);要点说明:
- pub/sub 客户端分离:
pubClient用于发布消息,subClient(由pubClient.duplicate()派生)用于订阅,二者都必须成功connect(); - 通道命名:每个节点订阅
engine.io#(公共广播通道)与engine.io#<nodeId>#(定向通道)两个通道;发布时若消息带recipientId则发到对应定向通道,否则发到公共通道(见 lib/redis.ts); - 客户端兼容性:源码中的
SUBSCRIBE辅助函数同时兼容redis包(存在sSubscribe方法)与ioredis包两种客户端(见 lib/redis.ts),test/redis.test.ts 对两种客户端分别做了测试; - 通道前缀可配置:
RedisEngine与setupPrimaryWithRedis均接受channelPrefix选项,默认值为"engine.io",多套环境共用同一 Redis 时建议改名隔离。
仓库配套的 compose.yaml 提供了本地调试用的 Redis 7 服务定义。生产场景请自行部署 Redis(建议启用持久化与 ACL)。
用法三:Node.js cluster + Redis(组合模式)
这是最完整的形态:单机多 worker 之间走 IPC,跨机集群通过 Redis 桥接。主进程调用setupPrimaryWithRedis(pubClient, subClient),同时扮演"IPC 转发器"和"Redis 网关"——worker 的 IPC 消息由主进程统一 publish 到 Redis,远端主进程收到后转发给本机对应 worker。
import cluster from "node:cluster"; import process from "node:process"; import { availableParallelism } from "node:os"; import { createClient } from "redis"; import { setupPrimaryWithRedis, NodeClusterEngine } from "@socket.io/cluster-engine"; import { createServer } from "node:http"; import { Server } from "socket.io"; if (cluster.isPrimary) { console.log(`Primary ${process.pid} is running`); const numCPUs = availableParallelism(); // fork workers for (let i = 0; i < numCPUs; i++) { cluster.fork(); } const pubClient = createClient(); const subClient = pubClient.duplicate(); await Promise.all([ pubClient.connect(), subClient.connect(), ]); // setup connection between and within the clusters setupPrimaryWithRedis(pubClient, subClient); // needed for packets containing Buffer objects (you can ignore it if you only send plaintext objects) cluster.setupPrimary({ serialization: "advanced", }); cluster.on("exit", (worker, code, signal) => { console.log(`worker ${worker.process.pid} died`); }); } else { const httpServer = createServer((req, res) => { res.writeHead(404).end(); }); const engine = new NodeClusterEngine(); engine.attach(httpServer, { path: "/socket.io/" }); const io = new Server(); io.bind(engine); // workers will share the same port httpServer.listen(3000); console.log(`Worker ${process.pid} started`); }注意组合模式下worker 端仍然使用NodeClusterEngine(IPC),Redis 连接只存在于主进程(setupPrimaryWithRedis内部会订阅engine.io#与本机专属通道,并把收到的远端消息转发给本机 worker,见 lib/redis.ts)。这样本机内部通信零网络开销,跨机才走 Redis。
仓库 examples/cluster-engine-node-cluster/server.js 与 examples/cluster-engine-redis/server.js 还展示了与@socket.io/cluster-adapter/@socket.io/redis-adapter组合使用的完整示例:engine 层负责"连接"的负载均衡,adapter 层负责"消息广播"的跨进程分发,两者搭配才是完整的 Socket.IO 集群方案。
Options 配置项
| 名称 | 说明 | 默认值 |
|---|---|---|
responseTimeout | 等待其他节点响应(如锁应答、升级应答)的最大毫秒数 | 1000 ms |
noopUpgradeInterval | 客户端升级 WebSocket 期间,向轮询通道发送 "noop" 心跳包的间隔毫秒数 | 200 ms |
delayedConnectionTimeout | 等待升级成功的最大毫秒数,超时则在当前节点完成连接建立 | 300 ms |
这些选项定义在 lib/engine.ts 的ClusterEngineOptions接口中,构造时通过Object.assign与默认值合并(lib/engine.ts)。除这三个专属选项外,ClusterEngine构造器还透传 engine.ioServerOptions(如pingInterval、upgradeTimeout、path等,测试中就使用了upgradeTimeout与pingInterval)。
使用示例:
const engine = new NodeClusterEngine({ responseTimeout: 2000, // 跨机响应放宽到 2s,容忍更慢的网络 noopUpgradeInterval: 100, // 更频繁地推送 noop 包,加速升级 delayedConnectionTimeout: 500, upgradeTimeout: 1000, // 透传给 engine.io });在 test/in-memory.test.ts 中可以观察到它们各自生效的验证场景:delayedConnectionTimeout: 50配合sleep(100)验证"延迟连接后完成握手"、"升级失败后恢复轮询"等路径(如"should upgrade (delayed)"、"should resume after upgrade failure"两个用例)。
工作原理
README 中这样概括其设计(详见 README.md):
This engine extends the one provided by the
engine.iopackage, so that sticky sessions are not required when scaling horizontally. The Node.js workers communicate via the IPC channel (or via Redis pub/sub) to check whether the Engine.IO session exists on another worker. In that case, the packets are forwarded to the worker which owns the session. Additionally, when a client starts with HTTP long-polling, the connection is delayed to allow the client to upgrade, so that the WebSocket connection ends up on the worker which owns the session.
结合 lib/engine.ts 源码,可以拆解为三层机制:
1. 集群消息协议(7 种消息类型)
MessageType枚举定义了节点间传递的全部消息(lib/engine.ts),每种消息都携带senderId(节点随机 ID,由randomBytes(3)生成),部分携带recipientId用于定向投递:
| 消息类型 | 方向 | 用途 |
|---|---|---|
ACQUIRE_LOCK | 请求方 → 集群 | 询问哪个 worker 持有该sid的会话,以及该传输操作是否被允许 |
ACQUIRE_LOCK_RESPONSE | 归属方 → 请求方 | 授予或拒绝锁请求(success) |
DRAIN | 归属方 → 远端传输 worker | 把需要写回客户端的包转发过去 |
PACKET | 远端传输 worker → 归属方 | 把从客户端收到的包转发给会话归属方 |
UPGRADE | 升级 worker → 归属方 | 通知归属方 WebSocket 升级探测成功/失败 |
UPGRADE_RESPONSE | 归属方 → 升级 worker | 告知是否接管会话,可选携带缓冲的数据包 |
CLOSE | 远端传输 worker → 归属方 | 通知归属方远端传输关闭或出错 |
仓库 docs/sequence_diagrams.md 给出了四种关键场景(轮询读、轮询写、WebSocket 升级成功/失败)的完整 Mermaid 时序图,建议结合阅读。
2. 分布式"锁":会话归属仲裁
这是整个方案的核心。当请求(GET/POST 轮询或 WebSocket 升级)到达某个 worker,而该sid不在本地时,verify()方法(覆盖自 engine.io 的Server.verify,lib/engine.ts)会触发_acquireLock:
- 请求方发布
ACQUIRE_LOCK(含sid、transportName、锁类型read/write,其中 GET 为read、POST 为write); - 持有该会话的 worker 收到后,用
isClientLockable(lib/engine.ts)判断当前是否允许此操作:polling读锁:要求会话当前确实是 polling 传输且轮询传输不可写(即没有正在挂起的轮询响应);polling写锁:要求会话当前是 polling 传输;websocket/webtransport升级锁:要求会话当前是 polling 且未在升级中、未升级完成;
- 归属方回复
ACQUIRE_LOCK_RESPONSE;请求方在responseTimeout内收不到应答则按失败处理(返回 400)。
test/in-memory.test.ts 的"should acquire read lock (different process)"等用例直接验证了"持有读锁时另一进程的轮询请求返回 400"这一行为。
3. 延迟连接 + noop 加速:让 WebSocket 落到归属 worker
Engine.IO 客户端通常先以 HTTP long-polling 发起握手,随后升级到 WebSocket。若轮询连接落在 Worker A、升级请求却被负载均衡分到 Worker B,就必须让 WebSocket 连接最终落在持有会话的 worker 上。方案是延迟connection事件:
ClusterEngine覆盖emit("connection", socket)(lib/engine.ts):当新 socket 走的是非 WebSocket 传输时,不立即派发connection,而是把 socket 标记为kDelayed,将其收到的包缓冲进kBuffer,并启动delayedConnectionTimeout定时器;- 若在超时前收到来自归属方的升级授权,
UPGRADE消息处理逻辑会决定是否"接管"(takeOver):若连接仍处于延迟状态,归属方把缓冲包随UPGRADE_RESPONSE一并交给升级 worker,由升级 worker 创建本地 Socket 并派发connection;若connection已派发,则保持会话归属不变,WebSocket 传输只作为"远端传输"挂在升级 worker 上(见 lib/engine.ts); - 为了加速升级,归属方会以
noopUpgradeInterval为间隔向轮询通道写入noop包(client.sendPacket("noop")),促使客户端立即发起升级探测;若超过upgradeTimeout(engine.io 原生选项)仍未完成,则重置升级状态并调用_doConnect完成连接(lib/engine.ts); _doConnect(lib/engine.ts)是兜底逻辑:延迟期满仍未升级成功时,就地完成connection派发并回放缓冲包,保证客户端无论如何都能连上。
升级探测阶段由_tryUpgrade(lib/engine.ts)完成标准的ping "probe"/pong "probe"/upgrade三步握手,任一环节超时或出错即宣告升级失败并通知归属方。测试"should upgrade"、"should upgrade (delayed)"、"should resume after upgrade failure"分别覆盖了这三种走向。
4. 数据包转发路径
- 客户端 → 归属方:远端传输 worker 通过
_onPacket发布PACKET消息;归属方收到后,若连接仍延迟则压入缓冲,否则调用client.onPacket()处理(lib/engine.ts); - 归属方 → 客户端:归属方通过
_forwardFlushWhenPolling/_forwardFlushWhenWebSocket替换 transport 的send方法,把需要写回的包以DRAIN消息转发给持有远端传输的 worker,由其写回客户端(lib/engine.ts)。轮询传输每次只能 drain 一次,处理完即从_remoteTransports中移除。
验证与运行
仓库自带三套测试,可作为行为规范与调试参考:
- test/in-memory.test.ts:用
EventEmitter模拟集群总线,无需真实进程即可单测锁、转发、升级全流程; - test/cluster.test.ts:真实 fork 3 个 worker,验证跨进程 ping/pong 与二进制(Buffer)收发;
- test/redis.test.ts:分别使用
redis与ioredis客户端,验证 3 个独立进程经 Redis 协同工作。
运行测试(需先启动 Redis 以跑 redis 用例):
cd packages/socket.io-cluster-engine npm install npm test本地联调 Redis 可用仓库自带的 compose.yaml:
docker compose up -d排查问题时可开启 debug 日志观察集群消息流转:
DEBUG=engine:* node server.js结语
@socket.io/cluster-engine通过"分布式锁仲裁会话归属 + 延迟连接等待升级 + 跨节点数据包转发"三件套,把 Engine.IO 从单进程模型改造成了多进程友好模型。对于已经在用 Socket.IO 且受限于粘性会话的团队,它提供了一条平滑的横向扩展路径:单机多核用NodeClusterEngine+setupPrimary,多机部署用RedisEngine或组合模式,再搭配对应的@socket.io/cluster-adapter/@socket.io/redis-adapter完成消息广播层,即可构建完整的 Socket.IO 集群。
License
MIT
【免费下载链接】socket.ioBidirectional and low-latency communication for every platform项目地址: https://gitcode.com/gh_mirrors/so/socket.io
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考