首页
/ Socket.IO 的 socket.io-adapter 深入解析:内存适配器的双索引结构、广播流水线与自研多节点适配器的地基

Socket.IO 的 socket.io-adapter 深入解析:内存适配器的双索引结构、广播流水线与自研多节点适配器的地基

2026-09-04 11:14:16作者:咎竹峻Karen

socket.io-adapter 是 Socket.IO 服务端默认的内存适配器(in-memory adapter),负责管理"房间-连接"映射并执行广播、断线恢复会话持久化等核心动作。本篇以 packages/socket.io-adapter/Readme.md 为核心,结合 lib/in-memory-adapter.tslib/cluster-adapter.ts 的源码实现,讲清楚它的数据结构、广播匹配算法、SessionAwareAdapter 的断线恢复机制,以及如何在它之上构建自己的多节点适配器。

socket.io-adapter 是什么:一句话定位与兼容性

原始文档对它的定义非常明确:

"Default socket.io in-memory adapter class." —— Socket.IO 默认的内存适配器类。

它同时声明了两个关键事实:

  1. 该模块不面向最终用户直接使用("This module is not intended for end-user usage")——业务代码不会直接 import 它来发消息,而是由 Socket.IO 服务端在内部持有并调用它;
  2. 它是为"被继承"而设计的接口("can be used as an interface to inherit from other adapters you might want to build")——官方文档以独立的 socket.io-redis 包作为在其之上构建适配器的示例,本仓库内同样提供了多个构建于该接口之上的实现(见后文集群场景)。

与 Socket.IO 服务端的版本兼容

原文档给出的兼容表如下:

Adapter version Socket.IO server version
1.x.x 1.x.x / 2.x.x
2.x.x 3.x.x

需要注意:上表沿袭自原始文档;当前仓库中该包的版本为 2.5.8(见 package.json),属于 2.x 系列,运行在 MIT 许可之下,依赖仅两个:debug(约 4.4.1)和 ws(约 8.21.0,用于 WebSocket 帧预编码优化)。

Socket.IO 服务端如何自动选择适配器

虽然文档说它"不面向最终用户",但服务端内部对适配器的选择逻辑直接决定了你会拿到哪个实现。在 Server 构造函数中(packages/socket.io/lib/index.ts#L337-L340):

// 若启用了 connectionStateRecovery,默认用 SessionAwareAdapter;否则用基础 Adapter
if (opts.connectionStateRecovery) {
  this.adapter(opts.adapter || SessionAwareAdapter);
} else {
  this.adapter(opts.adapter || Adapter);
}

也就是说,适配器有三条确定路径:

  • 显式传入 adapter 选项(ServerOptions 中声明为 adapter: AdapterConstructor,见 index.ts#L92-L96);
  • 启用了 connectionStateRecovery 配置 → 自动落到 SessionAwareAdapter
  • 什么都没配 → 落到基础 Adapter

Server#adapter(v) 方法(index.ts#L451-L462)除了保存类引用,还会遍历所有命名空间触发 _initAdapter()。而每个 Namespace 创建适配器实例的方式是(packages/socket.io/lib/namespace.ts#L197-L204):

_initAdapter(): void {
  this.adapter = new (this.server.adapter()!)(this); // 构造函数只接收一个 nsp 参数
  Promise.resolve(this.adapter.init()).catch((err) => {
    debug("error while initializing adapter: %s", err);
  });
}

这里有两个对自研适配器很重要的契约:构造函数只接收命名空间对象 nsp实例化后服务端会调用 init(),你可以在 init() 里建立 Redis 连接池等资源(基础类中 init()/close() 均为空实现,见 in-memory-adapter.ts#L61-L69)。

服务端各调用点与适配器方法的对应关系

从源码调用链可以确认,Socket.IO 与适配器之间的全部交互点如下(每一行都是仓库内的真实调用):

服务端调用点 适配器方法 触发场景
Socket#join()socket.ts#L464-L467 addAll(id, rooms) 连接加入房间
Socket#leave()socket.ts#L488 del(id, room) 连接离开房间
Socket#leaveAll()socket.ts#L497 delAll(id) 连接断开时退出所有房间
BroadcastOperator#emit()broadcast-operator.ts#L226 broadcast(packet, opts) 向客户端发事件
带确认的广播(broadcast-operator.ts#L283 broadcastWithAck(...) emit 最后一个参数为回调函数
BroadcastOperator#serverCount()broadcast-operator.ts#L303 serverCount() 带确认广播等待完整性校验
Socket#disconnect 断线恢复(socket.ts#L660-L665 persistSession(session) 可恢复断线时持久化会话
Namespace 重连握手(namespace.ts#L396 restoreSession(pid, offset) 客户端带 pid/offset 重连
socket.rooms 读取(socket.ts#L934-L936 socketRooms(id) 读取某连接所在房间
Namespace#serverSideEmitnamespace.ts#L547 serverSideEmit(packet) 服务端节点间互相发事件
fetchSockets/socketsJoin/...broadcast-operator.ts#L381-L472 fetchSockets/addSockets/delSockets/disconnectSockets 集群管理操作

这张表本质上就是"自研适配器必须关心的方法清单"——它比原文档只有一句"可以继承"具体得多。

数据结构:roomssids 双索引

Adapter 类继承自 Node.js 的 EventEmitterin-memory-adapter.ts#L46-L59),内部维护两张互为反向索引的内存表:

export class Adapter extends EventEmitter {
  public rooms: Map<Room, Set<SocketId>> = new Map();   // 房间 -> 连接集合
  public sids: Map<SocketId, Set<Room>> = new Map();    // 连接 -> 房间集合
  private readonly encoder;

  constructor(readonly nsp: any) {
    super();
    this.encoder = nsp.server.encoder; // nsp is a Namespace object
  }
}

先说类型契约(in-memory-adapter.ts#L8-L20):

  • SocketId:服务端在 Socket.IO 会话开始时下发给客户端的公开 ID,用于私聊等场景;
  • PrivateSessionId:会话开始时下发的私有 ID,专用于断线重连时的连接状态恢复;
  • Room:固定为 string。源码注释解释了原因——"可以把它扩展成 string | number,但那是破坏性变更",并关联了社区的兼容性问题。

房间生命周期:addAll / del / delAll

  • addAll(id, rooms)L87-L104):为每个房间双向写索引;首次创建房间时发出 create-room 事件,连接首次进入房间时发出 join-room 事件
  • del(id, room) / 私有 _del(room, id)L112-L131):删除双向索引;连接离开房间发 leave-room房间变空并被移除时发 delete-room
  • delAll(id)L138-L148):遍历该连接的所有房间逐个 _del,最后清掉 sids 条目。

由于 AdapterEventEmitter,这四类事件(create-roomjoin-roomleave-roomdelete-room)可以被外部订阅,用于在线人数统计等场景。

单元行为由 test/index.ts 直接验证,例如"加入/删除连接与房间"的用例(L5-L27):

const adapter = new Adapter({ server: { encoder: null } });
adapter.addAll("s1", new Set(["r1", "r2"]));
adapter.addAll("s2", new Set(["r2", "r3"]));
expect(adapter.rooms.has("r1")).to.be(true);
// ...
adapter.delAll("s2");
expect(adapter.rooms.has("r3")).to.be(false);

广播流水线:BroadcastOptions、匹配算法与编码优化

广播选项的两个类型

所有广播/房间操作共享同一组过滤参数(in-memory-adapter.ts#L22-L35):

export interface BroadcastFlags {
  volatile?: boolean;   // 允许丢弃(客户端未就绪时)
  compress?: boolean;   // 压缩传输
  local?: boolean;     // 仅当前节点(集群适配器会据此跳过发布)
  broadcast?: boolean;
  binary?: boolean;
  timeout?: number;    // 等待确认的超时(毫秒)
}

export interface BroadcastOptions {
  rooms: Set<Room>;    // 目标房间集合;空集合 = 全体连接
  except?: Set<Room>; // 排除的房间集合
  flags?: BroadcastFlags;
}

这些 flag 与你在服务端写的 io.to("room").volatile.compress(1).timeout(1000).emit(...) 一一对应——BroadcastOperator 会把链式调用累积进 flags 并原样传给适配器(broadcast-operator.ts#L117-L186)。

匹配与分发算法:apply()

broadcast()L162-L180)的完整流程是:

  1. packet 补上 nsp 名;
  2. 通过命名空间编码器把包编码一次preEncoded: true),随后对所有命中连接复用同一份编码结果;
  3. 调用私有的 apply(opts, callback) 遍历目标连接,逐个 socket.client.writeToEngine(encodedPackets, packetOpts) 写入底层 Engine.IO 连接。

apply()L333-L358)的匹配规则分两支:

  • 指定了 rooms:逐个房间取出连接集合,用临时 Set 去重(一个连接可能在多个目标房间),再排除 except 命中者;
  • rooms 为空:直接遍历 sids 的全部 key(即全体连接),只排除 except

except 参数在 computeExceptSids()L360-L370)中从"房间集合"展开为"连接 ID 集合"。这套逻辑有专门测试覆盖,比如"排除特定房间中的连接后广播"(test/index.ts#L61-L99)断言只有不在 r1s2 收到了消息。

两个性能细节

  • WebSocket 帧预编码_encode()L234-L256)在编码结果只有单个字符串帧时,利用 ws 库的 WebSocket.Sender.frame() 提前把 Engine.IO 的 "4"(message 类型前缀)+ 数据打包成完整的 WebSocket 帧挂到 wsPreEncodedFrame,后续写出时可跳过一次帧封装;
  • 多确认广播 broadcastWithAck()L197-L232):为所有命中连接分配同一个 packet.id(来自 nsp._ids 公共计数器,保证不重复),把回调存入每个连接的 acks,并统计 clientCount 回调给上层——配合 serverCount()BroadcastOperator 据此判断"所有节点上的所有客户端是否都已确认"(broadcast-operator.ts#L299-L306)。

连接状态恢复:SessionAwareAdapter

基础 AdapterpersistSession() 是空方法、restoreSession() 直接返回 nullL372-L397)——即默认不做断线恢复。真正干活的是同文件的 SessionAwareAdapterL409-L499),它维护两个私有状态:

private sessions: Map<PrivateSessionId, SessionWithTimestamp> = new Map();
private packets: PersistedPacket[] = [];

工作机制拆解:

  1. 保存会话persistSession() 记录 disconnectedAt 时间戳并以 pid(私有会话 ID)为键存会话。调用方在服务端是 Socket 断开时(socket.ts#L656-L665),且只有断线原因属于"可恢复"集合时才触发;
  2. 给事件包附加 offset:重写后的 broadcast() 会判断——事件类型包(type === 2)、无确认、非 volatile 三个条件同时满足时,用 yeast() 生成一个短 ID 追加到 packet.data 末尾(客户端能据此知道"最后处理的包",且格式向后兼容),并把 {id, opts, data, emittedAt} 压入 packets 队列;
  3. 恢复会话restoreSession(pid, offset) 先查会话是否存在、是否超过 maxDisconnectionDuration(该值取自服务端配置 nsp.server.opts.connectionStateRecovery.maxDisconnectionDuration),再按 offset 在 packets 中定位索引,之后逐个包用 shouldIncludePacket()L501-L509)判断"断线期间该客户端是否本来就能收到"——即包的 rooms 与其所在房间有交集、且没有命中 except——最后返回 missedPackets
  4. 定期清理:构造函数里有一个 60 秒的 setInterval 定时器(unref() 防止其阻止进程退出),按 maxDisconnectionDuration 阈值剔除过期会话,并从 packets 头部批量丢弃过期包。

调用链上,Namespace 在重连握手时解析出 sessionId/offset 后调用 restoreSession()namespace.ts#L392-L400),拿到会话后新的 Socket 会沿用原 sid/pid 并重新加入原房间(socket.ts#L162-L170)。

多节点集群:ClusterAdapter 抽象

原文档强调本包的价值在于"作为接口被继承"。仓库内对这一点最强的佐证,就是同包的 lib/cluster-adapter.ts 提供的 ClusterAdapter 抽象类。其类注释直接给出了实现契约:

"Any extending class must:

  • implement doPublish and doPublishResponse
  • call onMessage and onResponse"

即自研集群适配器只需实现两个抽象方法doPublish(message): Promise<Offset>L700)把消息推给集群其他成员(Redis 频道、消息总线等),doPublishResponse(requesterUid, response)L730-L733)把响应点对点发回请求方;并在收到消息时分别调用 onMessage/onResponse

ClusterAdapter 替子类完成的通用逻辑包括:

  • 节点身份:每个实例用 randomBytes(8) 生成 uidL14-L16),onMessage 会忽略自己和其他命名空间的消息;
  • 消息/响应协议ClusterMessage/ClusterResponse 类型与 MessageType 枚举(L45-L59),涵盖心跳(INITIAL_HEARTBEAT/HEARTBEAT/ADAPTER_CLOSE)、BROADCASTSOCKETS_JOIN/LEAVEDISCONNECT_SOCKETSFETCH_SOCKETS(_RESPONSE)SERVER_SIDE_EMIT(_RESPONSE)BROADCAST_CLIENT_COUNT/ACK 等。注意 BROADCAST 消息带 requestId 即表示需要确认;
  • 广播L427-L446):非 local 时先 publishAndReturnOffset{packet, opts} 发布出去,addOffsetIfNecessary() 会把发布返回的 offset 追加到事件包尾部(同样要求事件包、无确认、非 volatile,且服务端开启了 connectionStateRecovery),最后 super.broadcast() 处理本地连接;
  • fetchSocketsL570-L614):本地结果 + 向其余 serverCount() - 1 个节点发 FETCH_SOCKETS,收到齐全部响应才 resolve,超时(flags.timeout,默认 5000ms)则 reject 并报告"只收到 x/y 个响应";
  • serverSideEmitL616-L673):末参是函数时视为需要 ack,向其他节点广播 SERVER_SIDE_EMIT,各节点执行 nsp._onServerSideEmit(packet) 并把返回值作为 SERVER_SIDE_EMIT_RESPONSE 回传。

ClusterAdapterWithHeartbeat:心跳与节点发现

再往上还有一层 ClusterAdapterWithHeartbeatL744-L771),解决"不知道集群里有多少节点"的问题,它正是原文档所说"继承本模块构建其他适配器"的完整范本。两个可配置项及默认值(L29-L43):

选项 含义 默认值
heartbeatInterval 两次心跳的间隔(毫秒) 5_000
heartbeatTimeout 超过多久没收到心跳即视为节点下线(毫秒) 10_000

工作机制:

  • init() 时发布 INITIAL_HEARTBEAT,收到它的其他节点会回一条 HEARTBEAT,从而互相"认识";
  • 每次 publish 都会 scheduleHeartbeat(),通过 setTimeoutrefresh() 实现"有消息就不额外发心跳";
  • nodesMap(uid -> 最后见到时间)驱动 serverCount() = 1 + nodesMap.sizeL828-L830);
  • 每秒一次的清理定时器把超过 heartbeatTimeout 的节点移除;ADAPTER_CLOSE 消息则触发即时移除。removeNode() 还会提前结算等待该节点的 FETCH_SOCKETS/SERVER_SIDE_EMIT 请求,避免调用方无谓等到超时;
  • 与基础 ClusterAdapter 的区别在于请求完成条件:心跳版用 missingUids 集合按"UID 到齐"判定(L947-L999),因为基础版只能按"响应数 == serverCount-1"判定,而节点数只有心跳版才实时掌握。

本仓库中基于这套接口落地的完整示例包括:packages/socket.io-cluster-adapter(配合 Redis 的集群适配器)、packages/socket.io-cluster-engine(集群连接模式)以及示例工程 examples/cluster-engine-redisexamples/cluster-node-cluster,它们的测试(如 test/cluster-adapter.ts)覆盖了上述消息流程。

自研自定义适配器:最小模板与检查清单

结合原文档"可作为接口继承"的定位与前述源码调用点,一个自定义适配器的最小骨架如下(以注入外部存储为示意):

import { Adapter, BroadcastOptions, Room, SocketId } from "socket.io-adapter";

/**
 * 最小模板:仅标注必须/建议覆写的方法。
 * 存储逻辑按你的后端(Redis/Postgres 等)填充。
 */
class MyAdapter extends Adapter {
  private store: any; // 外部存储句柄

  public init() {
    // Namespace 实例化本类后会调用 init(),在此建立连接/连接池
  }

  public close() {
    // 服务端关闭时会被调用(packages/socket.io/lib/index.ts 的 close 流程)
  }

  public async addAll(id: SocketId, rooms: Set<Room>): Promise<void> {
    // 写入 "连接 -> 房间" 关系
  }

  public async del(id: SocketId, room: Room): Promise<void> {
    // 删除单个关系
  }

  public async delAll(id: SocketId): Promise<void> {
    // 删除该连接的全部关系
  }

  public async broadcast(packet: any, opts: BroadcastOptions): Promise<void> {
    // 1. 按 opts.rooms / opts.except 从存储中解析目标连接
    // 2. 本地命中的连接直接写;其余节点通过发布/订阅转发
    // 3. opts.flags.local === true 时只处理本地
  }

  public serverCount(): Promise<number> {
    // 基础类固定返回 1;集群实现需返回真实节点数
  }
}

在服务端切换(ServerOptions.adapter 或链式 io.adapter(),见 index.ts#L451-L462):

const { Server } = require("socket.io");
const io = new Server(httpServer);

io.adapter(MyAdapter); // 传的是类而非实例,各 Namespace 会各自 new MyAdapter(nsp)

按源码接口整理的覆写检查清单:

  • 房间索引addAll / del / delAll / socketRooms —— Socket#join#leave、断开清理、socket.rooms 都依赖它们;
  • 广播broadcast(必须处理 rooms/except/flags)、broadcastWithAck(跨节点场景需自行聚合 clientCount 与各节点 ack);
  • 查询与管理sockets(rooms)fetchSocketsaddSocketsdelSocketsdisconnectSockets —— 对应 io.fetchSockets()io.in("r").fetchSockets()socketsJoin/LeavedisconnectSockets
  • 服务端通信serverSideEmit(基础类默认只打一条 console.warn 提示不支持,L376-L380);
  • 断线恢复(可选):persistSession / restoreSession,若实现则客户端带 pid 重连时可补齐 missedPackets
  • 生命周期init / close / serverCount

若你的需求只是"多节点共享广播",不必从头实现:继承 ClusterAdapter(或 ClusterAdapterWithHeartbeat)并实现 doPublish/doPublishResponse 即可,消息协议、请求-响应聚合、超时处理全部由基类完成。

测试与验证方式

该包的测试基于 mocha + tsx 直接跑 TypeScript(package.jsontest 脚本为 npm run format:check && npm run compile && nyc mocha --import=tsx test/*.ts)。测试组织:

  • test/index.ts:内存适配器的房间增删、sockets() 列表、except 排除广播、fetchSocketsserverSideEmit 告警行为等,均用假 nsp 对象(提供 server.encodersockets Map)构造,例如 new Adapter({ server: { encoder: null } })
  • test/cluster-adapter.ts:用一个内存版 doPublish 模拟多节点,验证广播、fetchSockets 响应聚合、心跳与节点下线等路径;
  • test/util.ts:共享工具。

小结

socket.io-adapter 在 Socket.IO 架构中承担的是"房间索引 + 广播执行器"的角色,它的原始文档虽短,但给出的两条信息——默认内存适配器面向继承的接口设计——恰恰是整个包的核心:

  • Adapterrooms/sids 双索引实现 O(1) 级的房间增删与广播匹配,apply() 统一处理 rooms/except 过滤,broadcastWithAck 支撑多节点确认广播;
  • SessionAwareAdapter 在内存中保留断线窗口内的会话与事件包(带 offset),支撑连接状态恢复,且仅在服务端开启 connectionStateRecovery 时成为默认实现;
  • ClusterAdapter/ClusterAdapterWithHeartbeat 把"发布/响应"抽象为两个钩子,让 Redis 等后端的集群适配器只需实现 doPublish/doPublishResponse
  • 自研适配器时,本文列出的服务端调用点对照表与覆写检查清单,可直接作为实现与验收依据。
登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
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.79 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
988
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384