首页
/ Socket.IO @socket.io/cluster-adapter 版本演进全解析:从 0.1.0 首发到 Monorepo 架构重构

Socket.IO @socket.io/cluster-adapter 版本演进全解析:从 0.1.0 首发到 Monorepo 架构重构

2026-09-04 14:04:26作者:丁柯新Fawn

本文以 @socket.io/cluster-adapter 包的官方变更日志为主线,逐版本复盘该包从 2021 年首发到 2025 年完成架构重构的完整历程,并结合当前 monorepo 中的源码(NodeClusterAdapterClusterAdapter、心跳机制与请求-响应路由)说明每个版本变更对应的实现细节。读完后你将理解:为什么多 ack 广播、IPC 错误处理、peerDependencies 与逻辑下沉这四次关键变更分别解决了什么问题,以及 0.3.0 之后"瘦包 + 共享基类"的架构是如何工作的。

Socket.IO cluster adapter 架构图:Worker 进程通过 Primary 进程中继消息

包的定位:在 Node.js cluster Worker 之间广播包

在进入版本时间线之前,先明确这个包解决的核心问题:在 Node.js cluster 模块(README 中示例写法为 cluster)架构下,每个 Worker 进程各持有一份独立的 Socket.IO Server 实例,内存适配器天然无法跨进程通信。@socket.io/cluster-adapter 通过主进程(Primary)作为中继,把广播、房间变更、断连、fetchSocketsserverSideEmit 等操作转发给所有其他 Worker,使多台 Socket.IO 服务器表现得像一个逻辑整体。

从当前源码结构看,整个通信链路只有两条原语:

  • Worker 侧:NodeClusterAdapter.doPublish() 调用 process.send(message, null, {}, ignoreError) 把消息发给主进程,消息统一打上 MESSAGE_SOURCE = "_sio_adapter" 标记(见 lib/index.tsdoPublish);
  • 主进程侧:setupPrimary() 注册 cluster.on("message") 监听器,按消息类型执行两种路由策略——响应类消息(FETCH_SOCKETS_RESPONSESERVER_SIDE_EMIT_RESPONSE)按 requesterUid 定向回发给发起请求的 Worker;其余广播类消息则转发给"除发送者之外的所有 Worker"(见 setupPrimary)。

这正是 0.3.0 重构后仍保留在本包中的全部核心逻辑:消息来源校验、命名空间隔离(message.nsp !== this.nsp.name 时忽略)、以及星型拓扑的中继转发。

版本时间线

以下为 CHANGELOG.md 中记录的完整版本表:

版本 发布日期 核心变更
0.3.0 2025-10-16 重要重构:大部分逻辑下沉至 socket.io-adapter 包的 ClusterAdapter 类;包并入 socket.io monorepo
0.2.2 2023-03-24 socket.io-adapter 加入 peerDependencies,修复与 socket.io 包引入的 adapter 版本不同步的问题
0.2.1 2022-10-13 正确修复 ERR_IPC_CHANNEL_CLOSED 错误
0.2.0 2022-04-28 支持广播并收集多个客户端 ack(对应 socket.io@4.5.0 新增特性)
0.1.0 2021-06-22 初始版本

下面逐版本展开,并对照当前仓库源码验证每条变更的落地形态。

0.1.0(2021-06-22):初始版本

变更日志对首版的描述只有一行:"Initial commit"。它是作为独立仓库(standalone package)发布的,提供 createAdapter() 工厂函数与 setupPrimary() 主进程中继函数两个公开 API——这两个 API 的签名一直延续到当前版本。

0.2.0(2022-04-28):广播与多 Ack 支持

这是该包功能演进上最重要的一版。变更日志明确写道:

broadcast and expect multiple acks — This feature was added in socket.io@4.5.0

io.timeout(1000).emit("some-event", (err, responses) => {
  // ...
});

Thanks to this change, it will now work within a Node.js cluster.

也就是说,socket.io@4.5.0 在服务端新增了"一次广播、等待所有客户端回执、超时兜底"的 API,但只有在集群适配器里补齐跨 Worker 的响应聚合逻辑后,它在 Node.js cluster 下才能真正工作。

从当前源码看,这条能力由 ClusterAdapter 中的三部分组成(0.3.0 重构前它位于本包内部,重构后移至 cluster-adapter.ts):

  1. 请求注册与消息发布broadcastWithAck() 生成随机 requestId,把 clientCountCallbackack 存入 ackRequests Map,然后携带 requestIdBROADCAST 消息发布到集群;requestId 的存在与否就是"是否需要 ack"的判定依据(见 broadcastWithAck);
  2. 响应消息类型:消息类型枚举中专门为此定义了 BROADCAST_CLIENT_COUNT(各 Worker 上报本地目标客户端数量)与 BROADCAST_ACK(客户端回执透传)两种 ClusterResponse(见 MessageTypeClusterResponse);
  3. 超时清理:由于 adapter 层无法得知"是否已从每个客户端收到回执",源码用 setTimeoutopts.flags.timeout 到期后直接从 ackRequests 中删除该请求,避免 Map 无限增长。

测试用例 test/index.ts 中的 "broadcasts with multiple acknowledgements" 系列(含二进制内容、无客户端、超时四个变体)逐一对应这些行为:三个客户端分别 cb(1)/cb(2)/cb(3) 后发起端能聚合全部响应;某客户端不回时收到 Error 且已收响应仍保留(见 test/index.tsworker.js)。

0.2.1(2022-10-13):修复 ERR_IPC_CHANNEL_CLOSED

cluster 模块的一个典型隐患是:Worker 进程退出后,向其发送 IPC 消息会触发 ERR_IPC_CHANNEL_CLOSED 异常,若未捕获会导致整个进程崩溃——在 Worker 频繁重启的集群场景下尤其致命。

当前源码中可以看到该修复的落地形态:doPublish()doPublishResponse()setupPrimary() 中的每一次 process.send() / worker.send() 都显式传入了一个空函数 ignoreError 作为错误回调(function ignoreError() {},见 lib/index.ts)。即 IPC 通道已关闭这类错误被静默吞掉,消息丢失交由心跳机制去感知节点掉线,而不是让发送方崩溃。

0.2.2(2023-03-24):将 socket.io-adapter 纳入 peerDependencies

变更日志说明:socket.io-adapter 被加入 peerDependencies,"in order to fix sync issues with the version imported by the socket.io package"——即修复本包与 socket.io 包各自引入的 adapter 版本不一致导致的同步问题。

当前 package.json 中可以确认这一约束仍然有效:

"peerDependencies": {
  "socket.io-adapter": "~2.5.5"
}

这个约束有实际意义:NodeClusterAdapter 继承自 socket.io-adapter 导出的 ClusterAdapterWithHeartbeat(见 lib/index.ts),若两处解析出的 adapter 版本不同,类层次与消息类型枚举可能对不上,跨 Worker 消息就会解析失败。使用 npm/yarn 安装时,peer 依赖会强制安装方对齐版本。

0.3.0(2025-10-16):重要重构与 Monorepo 化

变更日志对 0.3.0 的原文是:

This release contains an important refactor of the adapter, as most of the logic has been moved in the ClusterAdapter class of the socket.io-adapter package. Besides, the @socket.io/cluster-adapter package is now part of the socket.io monorepo.

这一版有两层含义,仓库现状可以逐条印证:

其一,包并入 socket.io monorepo。 本包现位于仓库的 packages/socket.io-cluster-adapter 目录下,与 socket.iosocket.io-adapterengine.io 等包同仓开发、同版本节奏发布,README 也明确指向 monorepo 内的地址。

其二,逻辑下沉到 socket.io-adapter 重构后本包的 lib/index.ts 仅剩约 117 行,NodeClusterAdapter 的职责收窄为三件"Node.js cluster 专属"的事:

export class NodeClusterAdapter extends ClusterAdapterWithHeartbeat {
  constructor(nsp: any, opts: ClusterAdapterOptions = {}) {
    super(nsp, opts);
    // 1. 监听 process "message" 事件,校验来源与命名空间后交给 onMessage()
    process.on("message", (message) => {
      if (message?.source !== MESSAGE_SOURCE) return;      // 忽略非 adapter 消息
      if (message.nsp !== this.nsp.name) return;           // 忽略其他命名空间
      this.onMessage(message);
    });
    this.init();
  }

  // 2. 通过 IPC 把消息发给 Primary(source 标记为 _sio_adapter)
  protected override doPublish(message) {
    message.source = MESSAGE_SOURCE;
    process.send(message, null, {}, ignoreError);
    return Promise.resolve(""); // connection state recovery is not supported
  }

  // 3. 通过 IPC 定向回发响应(携带 requesterUid)
  protected override doPublishResponse(requesterUid, response) {
    response.source = MESSAGE_SOURCE;
    response.requesterUid = requesterUid;
    process.send(response, null, {}, ignoreError);
    return Promise.resolve();
  }
}

而广播、房间管理、fetchSocketsserverSideEmit、心跳与节点存活追踪等"集群语义"全部由基类承担,位于 packages/socket.io-adapter/lib/cluster-adapter.ts。这种"传输抽象 + 集群语义基类"的分层,使得同一套 ClusterAdapter 可以直接派生出 Redis、Postgres 等其他传输的集群适配器(仓库中已有 socket.io-redis-streams-emittersocket.io-postgres-emitter 等配套包)。

两个值得注意的实现细节:

  • doPublish 返回空 offsetaddOffsetIfNecessary() 依赖 doPublish 返回的 offset 来支撑 Connection State Recovery(连接状态恢复,socket.io@4.6.0 引入)。由于 Node.js IPC 无法提供全局消息偏移量,doPublish 固定返回 Promise.resolve(""),注释写明 "connection state recovery is not supported"。这意味着即便 0.2.2 日志曾预告"下一个版本支持连接状态恢复",从当前源码看该特性在 cluster 适配器中仍未启用,配置 connectionStateRecovery 后跨 Worker 的重放仍不可用——这是使用该包时需要明确的限制;
  • 心跳参数默认值ClusterAdapterWithHeartbeat 构造时默认 heartbeatInterval: 5_000heartbeatTimeout: 10_000(毫秒),即每 5 秒发一次心跳,超过 10 秒未收到某节点消息即判定下线并从 nodesMap 中移除(见 ClusterAdapterOptionsClusterAdapterWithHeartbeat 构造函数)。节点被移除时会触发 removeNode(),把仍在等待该节点响应的 fetchSockets / serverSideEmit 请求直接 resolve,避免请求因"少一个响应"而空转至超时(见 removeNode)。

标准用法:配合 @socket.io/sticky 的完整示例

README 给出的标准接入方式(@socket.io/cluster-adapter 常与粘性会话中间件 @socket.io/sticky 联用)如下,主进程负责粘性路由与 IPC 中继,每个 Worker 各挂一个带 adapter 的 Server

const cluster = require("cluster");
const http = require("http");
const { Server } = require("socket.io");
const numCPUs = require("os").cpus().length;
const { setupMaster, setupWorker } = require("@socket.io/sticky");
const { createAdapter, setupPrimary } = require("@socket.io/cluster-adapter");

if (cluster.isMaster) {
  console.log(`Master ${process.pid} is running`);

  const httpServer = http.createServer();

  // setup sticky sessions
  setupMaster(httpServer, {
    loadBalancingMethod: "least-connection",
  });

  // setup connections between the workers
  setupPrimary();

  // needed for packets containing buffers (you can ignore it if you only send plaintext objects)
  // Node.js < 16.0.0
  cluster.setupMaster({
    serialization: "advanced",
  });
  // Node.js > 16.0.0
  // cluster.setupPrimary({
  //   serialization: "advanced",
  // });

  httpServer.listen(3000);

  for (let i = 0; i < numCPUs; i++) {
    cluster.fork();
  }

  cluster.on("exit", (worker) => {
    console.log(`Worker ${worker.process.pid} died`);
    cluster.fork();
  });
} else {
  console.log(`Worker ${process.pid} started`);

  const httpServer = http.createServer();
  const io = new Server(httpServer);

  // use the cluster adapter
  io.adapter(createAdapter());

  // setup connection with the primary process
  setupWorker(io);

  io.on("connection", (socket) => {
    /* ... */
  });
}

两点实操细节:

  • serialization: "advanced"(Node.js 16 起更名为 cluster.setupPrimary)是发送含 Buffer 的二进制包所必需的,只发纯文本对象可省略——测试工程在 test/index.ts 中也是以同样方式配置 cluster.setupMaster,并用"含 Buffer 的广播"与"多 ack 二进制内容"两个用例覆盖这条路径;
  • 仓库内的 examples/cluster-engine-node-cluster/server.js 展示了 0.3.0 重构后的新写法:Worker 通过 new Server({ adapter: createAdapter() }) 在构造参数中传入 adapter,主进程则同时调用 setupPrimaryEngine()(Engine.IO 层连接在 Worker 间迁移)与 setupPrimaryAdapter()(应用层消息中继),这是当前版本更贴近实际部署的组合形态。

测试如何验证各版本能力

包的测试套件 test/index.ts 以"3 个真实 cluster Worker + 3 个真实客户端"的端到端方式组织,覆盖了 CHANGELOG 各版本承诺的行为边界:

  • 0.1.0 基础能力:全量广播、命名空间内广播、房间定向(io.to("room1"))、房间排除(except)、local 只发本 Worker 客户端;
  • 0.2.0 多 ack:含二进制内容、无目标客户端(responses 为空数组)、部分客户端不回(收到 Error 且返回已收响应)三个边界场景;
  • 实用方法socketsJoin / socketsLeave(含 io.in("room1").socketsJoin("room2") 的选择性变体)、disconnectSockets(断连原因应为 io server disconnect)、fetchSockets(聚合 3 个 Worker 共 3 个 socket,并断言 adapter.requests 清理完毕)、serverSideEmit(带 ack 时等待其余 2 个 Worker 响应,任一节点失联则收到 "timeout reached: missing 1 responses" 错误,与 removeNode/DEFAULT_TIMEOUT 逻辑呼应)。

Worker 侧的 test/worker.js 通过 io.adapter(createAdapter()) 装配适配器后,用 process.on("message") 把测试指令映射为具体的 Server API 调用,是整个集群行为的执行端。

小结:如何选择与使用

  • 如果你要在同一台机器的 Node.js cluster Worker 之间广播事件,@socket.io/cluster-adapter + @socket.io/sticky 是零外部依赖的组合,无需 Redis 等中间件;跨机器的多实例部署则应考虑 Redis/Postgres 等其他适配器(README 中列有相关包)。
  • 版本上,0.2.0 是"多 ack 广播"的可用性分界线(依赖 socket.io@4.5.0),0.2.2 是"peer 依赖对齐"的稳定性分界线,0.3.0 是"逻辑下沉到 socket.io-adapter"的架构分界线——当前仓库版本即为 0.3.0,NodeClusterAdapter 已是一个只封装 IPC 传输的瘦类。
  • 已知限制需心里有数:连接状态恢复(Connection State Recovery)在该包中不受支持(doPublish 固定返回空 offset),且 serverSideEmit 带 ack 场景的超时时长为源码内固定值(DEFAULT_TIMEOUT = 5000 毫秒,测试注释也确认 "currently not possible to configure the timeout delay")。

以上所有结论均可在当前仓库中复核:变更历史见 packages/socket.io-cluster-adapter/CHANGELOG.md,实现见 packages/socket.io-cluster-adapter/lib/index.tspackages/socket.io-adapter/lib/cluster-adapter.ts,行为验证见 packages/socket.io-cluster-adapter/test/index.ts

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
987
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384