Socket.IO 的 socket.io-adapter 深入解析:内存适配器的双索引结构、广播流水线与自研多节点适配器的地基
socket.io-adapter 是 Socket.IO 服务端默认的内存适配器(in-memory adapter),负责管理"房间-连接"映射并执行广播、断线恢复会话持久化等核心动作。本篇以 packages/socket.io-adapter/Readme.md 为核心,结合 lib/in-memory-adapter.ts、lib/cluster-adapter.ts 的源码实现,讲清楚它的数据结构、广播匹配算法、SessionAwareAdapter 的断线恢复机制,以及如何在它之上构建自己的多节点适配器。
socket.io-adapter 是什么:一句话定位与兼容性
原始文档对它的定义非常明确:
"Default socket.io in-memory adapter class." —— Socket.IO 默认的内存适配器类。
它同时声明了两个关键事实:
- 该模块不面向最终用户直接使用("This module is not intended for end-user usage")——业务代码不会直接
import它来发消息,而是由 Socket.IO 服务端在内部持有并调用它; - 它是为"被继承"而设计的接口("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#serverSideEmit(namespace.ts#L547) |
serverSideEmit(packet) |
服务端节点间互相发事件 |
fetchSockets/socketsJoin/...(broadcast-operator.ts#L381-L472) |
fetchSockets/addSockets/delSockets/disconnectSockets |
集群管理操作 |
这张表本质上就是"自研适配器必须关心的方法清单"——它比原文档只有一句"可以继承"具体得多。
数据结构:rooms 与 sids 双索引
Adapter 类继承自 Node.js 的 EventEmitter(in-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条目。
由于 Adapter 是 EventEmitter,这四类事件(create-room、join-room、leave-room、delete-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)的完整流程是:
- 给
packet补上nsp名; - 通过命名空间编码器把包编码一次(
preEncoded: true),随后对所有命中连接复用同一份编码结果; - 调用私有的
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)断言只有不在 r1 的 s2 收到了消息。
两个性能细节
- 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
基础 Adapter 的 persistSession() 是空方法、restoreSession() 直接返回 null(L372-L397)——即默认不做断线恢复。真正干活的是同文件的 SessionAwareAdapter(L409-L499),它维护两个私有状态:
private sessions: Map<PrivateSessionId, SessionWithTimestamp> = new Map();
private packets: PersistedPacket[] = [];
工作机制拆解:
- 保存会话:
persistSession()记录disconnectedAt时间戳并以pid(私有会话 ID)为键存会话。调用方在服务端是Socket断开时(socket.ts#L656-L665),且只有断线原因属于"可恢复"集合时才触发; - 给事件包附加 offset:重写后的
broadcast()会判断——事件类型包(type === 2)、无确认、非 volatile 三个条件同时满足时,用yeast()生成一个短 ID 追加到packet.data末尾(客户端能据此知道"最后处理的包",且格式向后兼容),并把{id, opts, data, emittedAt}压入packets队列; - 恢复会话:
restoreSession(pid, offset)先查会话是否存在、是否超过maxDisconnectionDuration(该值取自服务端配置nsp.server.opts.connectionStateRecovery.maxDisconnectionDuration),再按 offset 在packets中定位索引,之后逐个包用shouldIncludePacket()(L501-L509)判断"断线期间该客户端是否本来就能收到"——即包的rooms与其所在房间有交集、且没有命中except——最后返回missedPackets; - 定期清理:构造函数里有一个 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
doPublishanddoPublishResponse- call
onMessageandonResponse"
即自研集群适配器只需实现两个抽象方法:doPublish(message): Promise<Offset>(L700)把消息推给集群其他成员(Redis 频道、消息总线等),doPublishResponse(requesterUid, response)(L730-L733)把响应点对点发回请求方;并在收到消息时分别调用 onMessage/onResponse。
ClusterAdapter 替子类完成的通用逻辑包括:
- 节点身份:每个实例用
randomBytes(8)生成uid(L14-L16),onMessage会忽略自己和其他命名空间的消息; - 消息/响应协议:
ClusterMessage/ClusterResponse类型与MessageType枚举(L45-L59),涵盖心跳(INITIAL_HEARTBEAT/HEARTBEAT/ADAPTER_CLOSE)、BROADCAST、SOCKETS_JOIN/LEAVE、DISCONNECT_SOCKETS、FETCH_SOCKETS(_RESPONSE)、SERVER_SIDE_EMIT(_RESPONSE)、BROADCAST_CLIENT_COUNT/ACK等。注意BROADCAST消息带requestId即表示需要确认; - 广播(L427-L446):非
local时先publishAndReturnOffset把{packet, opts}发布出去,addOffsetIfNecessary()会把发布返回的offset追加到事件包尾部(同样要求事件包、无确认、非 volatile,且服务端开启了connectionStateRecovery),最后super.broadcast()处理本地连接; - fetchSockets(L570-L614):本地结果 + 向其余
serverCount() - 1个节点发FETCH_SOCKETS,收到齐全部响应才 resolve,超时(flags.timeout,默认 5000ms)则 reject 并报告"只收到 x/y 个响应"; - serverSideEmit(L616-L673):末参是函数时视为需要 ack,向其他节点广播
SERVER_SIDE_EMIT,各节点执行nsp._onServerSideEmit(packet)并把返回值作为SERVER_SIDE_EMIT_RESPONSE回传。
ClusterAdapterWithHeartbeat:心跳与节点发现
再往上还有一层 ClusterAdapterWithHeartbeat(L744-L771),解决"不知道集群里有多少节点"的问题,它正是原文档所说"继承本模块构建其他适配器"的完整范本。两个可配置项及默认值(L29-L43):
| 选项 | 含义 | 默认值 |
|---|---|---|
heartbeatInterval |
两次心跳的间隔(毫秒) | 5_000 |
heartbeatTimeout |
超过多久没收到心跳即视为节点下线(毫秒) | 10_000 |
工作机制:
init()时发布INITIAL_HEARTBEAT,收到它的其他节点会回一条HEARTBEAT,从而互相"认识";- 每次
publish都会scheduleHeartbeat(),通过setTimeout的refresh()实现"有消息就不额外发心跳"; nodesMap(uid -> 最后见到时间)驱动serverCount() = 1 + nodesMap.size(L828-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-redis、examples/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)、fetchSockets、addSockets、delSockets、disconnectSockets—— 对应io.fetchSockets()、io.in("r").fetchSockets()、socketsJoin/Leave、disconnectSockets; - 服务端通信:
serverSideEmit(基础类默认只打一条console.warn提示不支持,L376-L380); - 断线恢复(可选):
persistSession/restoreSession,若实现则客户端带pid重连时可补齐missedPackets; - 生命周期:
init/close/serverCount。
若你的需求只是"多节点共享广播",不必从头实现:继承 ClusterAdapter(或 ClusterAdapterWithHeartbeat)并实现 doPublish/doPublishResponse 即可,消息协议、请求-响应聚合、超时处理全部由基类完成。
测试与验证方式
该包的测试基于 mocha + tsx 直接跑 TypeScript(package.json 的 test 脚本为 npm run format:check && npm run compile && nyc mocha --import=tsx test/*.ts)。测试组织:
- test/index.ts:内存适配器的房间增删、
sockets()列表、except排除广播、fetchSockets、serverSideEmit告警行为等,均用假nsp对象(提供server.encoder与socketsMap)构造,例如new Adapter({ server: { encoder: null } }); - test/cluster-adapter.ts:用一个内存版
doPublish模拟多节点,验证广播、fetchSockets响应聚合、心跳与节点下线等路径; - test/util.ts:共享工具。
小结
socket.io-adapter 在 Socket.IO 架构中承担的是"房间索引 + 广播执行器"的角色,它的原始文档虽短,但给出的两条信息——默认内存适配器、面向继承的接口设计——恰恰是整个包的核心:
Adapter用rooms/sids双索引实现 O(1) 级的房间增删与广播匹配,apply()统一处理rooms/except过滤,broadcastWithAck支撑多节点确认广播;SessionAwareAdapter在内存中保留断线窗口内的会话与事件包(带 offset),支撑连接状态恢复,且仅在服务端开启connectionStateRecovery时成为默认实现;ClusterAdapter/ClusterAdapterWithHeartbeat把"发布/响应"抽象为两个钩子,让 Redis 等后端的集群适配器只需实现doPublish/doPublishResponse;- 自研适配器时,本文列出的服务端调用点对照表与覆写检查清单,可直接作为实现与验收依据。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00