FEATURED · 精选文章

Socket.IO Cluster Engine 实战指南:无需粘性会话的多进程横向扩展方案

发布时间 / 2026/9/18 2:41:54
来源 / 创域科博编辑部
栏目 / 资讯中心
Socket.IO Cluster Engine 实战指南:无需粘性会话的多进程横向扩展方案 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即继承自它的Serverengine.io-parser ~5.2.3Engine.IO 数据包编解码msgpack/msgpack ~2.8.0Redis 模式下集群消息的序列化格式debug ~4.4.1调试日志可通过DEBUGengine:*环境变量开启。包对外导出的 API 集中在 lib/index.tsNodeClusterEngine、setupPrimary、RedisEngine、setupPrimaryWithRedis以及类型ClusterEngineOptions。用法一Node.js cluster单机多进程核心思路主进程primaryfork多个 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.tsserialization: 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.tstest/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 端仍然使用NodeClusterEngineIPCRedis 连接只存在于主进程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 msnoopUpgradeInterval客户端升级 WebSocket 期间向轮询通道发送 noop 心跳包的间隔毫秒数200 msdelayedConnectionTimeout等待升级成功的最大毫秒数超时则在当前节点完成连接建立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.mdThis engine extends the one provided by theengine.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归属方 → 请求方授予或拒绝锁请求successDRAIN归属方 → 远端传输 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.verifylib/engine.ts会触发_acquireLock请求方发布ACQUIRE_LOCK含sid、transportName、锁类型read/write其中 GET 为read、POST 为write持有该会话的 worker 收到后用isClientLockablelib/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 落到归属 workerEngine.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)促使客户端立即发起升级探测若超过upgradeTimeoutengine.io 原生选项仍未完成则重置升级状态并调用_doConnect完成连接lib/engine.ts_doConnectlib/engine.ts是兜底逻辑延迟期满仍未升级成功时就地完成connection派发并回放缓冲包保证客户端无论如何都能连上。升级探测阶段由_tryUpgradelib/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.yamldocker compose up -d排查问题时可开启 debug 日志观察集群消息流转DEBUGengine:* node server.js结语socket.io/cluster-engine通过分布式锁仲裁会话归属 延迟连接等待升级 跨节点数据包转发三件套把 Engine.IO 从单进程模型改造成了多进程友好模型。对于已经在用 Socket.IO 且受限于粘性会话的团队它提供了一条平滑的横向扩展路径单机多核用NodeClusterEnginesetupPrimary多机部署用RedisEngine或组合模式再搭配对应的socket.io/cluster-adapter/socket.io/redis-adapter完成消息广播层即可构建完整的 Socket.IO 集群。LicenseMIT【免费下载链接】socket.ioBidirectional and low-latency communication for every platform项目地址: https://gitcode.com/gh_mirrors/so/socket.io创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻