Socket.IO Postgres Emitter:从独立 Node.js 进程广播与管控整个 Socket.IO 服务集群
@socket.io/postgres-emitter 包让你在一个与 Socket.IO 服务器集群分离的 Node.js 进程中,通过 PostgreSQL 的 LISTEN/NOTIFY 机制向整组 Socket.IO 服务器广播事件、批量管理房间、断开指定连接,甚至向服务器实例本身发送事件。本篇基于仓库内该包的 README 与源码,完整讲解其安装、API 用法、底层"小载荷走 NOTIFY、大载荷走附表"的双路径发布机制,并结合 examples/postgres-adapter-example 示例与 packages/socket.io-postgres-emitter/test/index.ts 测试用例,说明如何在真实多节点集群中落地与验证这些能力。
适用场景:为什么不直接在服务器进程里 emit?
Socket.IO 的官方广播能力局限在同一个服务器进程(或通过 adapter 扩展到集群内其他进程)。但许多业务逻辑跑在独立的服务里:定时任务、消息队列消费者、爬虫、CI 系统等。这些进程并不持有 http.Server,无法直接创建 Server 实例,却需要触达所有在线客户端。
Emitter 就是为这个场景设计的:它只需要一个 pg 包的连接池对象,就能把事件"写进" PostgreSQL 的通道里;集群中每个 Socket.IO 服务器(配置了配套的 @socket.io/postgres-adapter)会监听同一通道并执行相应动作。README 中明确说明:本包必须与 @socket.io/postgres-adapter 配合使用,二者构成"发射端 + 接收端"的完整闭环。
从 包定义 可以看到,该包当前版本为 0.1.1,运行时仅依赖 @msgpack/msgpack(用于二进制载荷编码)与 debug(用于日志输出)。
README 列出的支持功能为:
- 广播(broadcasting)
- 实用方法(utility methods):
socketsJoin、socketsLeave、disconnectSockets、serverSideEmit
安装
npm install @socket.io/postgres-emitter pg
TypeScript 项目额外需要 @types/pg。pg 是 PostgreSQL 的 Node.js 官方客户端,Emitter 不自己管理数据库连接,而是复用你传入的 pg.Pool 对象。
快速上手
README 给出的最小用法如下:
const { Emitter } = require("@socket.io/postgres-emitter");
const { Pool } = require("pg");
const pool = new Pool({
user: "postgres",
host: "localhost",
database: "postgres",
password: "changeit",
port: 5432,
});
const io = new Emitter(pool);
setInterval(() => {
io.emit("ping", new Date());
}, 1000);
这段代码的含义是:连接本地 PostgreSQL(与 测试配置 及 postgres-adapter 示例 相同的连接参数),每 1 秒向默认命名空间 / 的所有客户端推送一条 ping 事件,参数为当前时间。
前提是集群中每台 Socket.IO 服务器都挂载了 postgres-adapter。仓库内置的 postgres-adapter 示例服务器 展示了接收端写法:
import { Server } from "socket.io";
import { createAdapter } from "@socket.io/postgres-adapter";
import pg from "pg";
const pool = new pg.Pool({
user: "postgres",
host: "localhost",
database: "postgres",
password: "changeit",
port: 5432,
});
// 大载荷/二进制数据需要的附表
await pool.query(`
CREATE TABLE IF NOT EXISTS socket_io_attachments (
id bigserial UNIQUE,
created_at timestamptz DEFAULT NOW(),
payload bytea
);
`);
const io = new Server({
adapter: createAdapter(pool)
});
注意 socket_io_attachments 这张表:它是 Emitter 存放"超额载荷"的地方,两端都必须能访问同一张表(表名可通过选项定制,见下文)。
示例目录还提供了 cluster.js,用 Node.js 的 cluster 模块 fork 出 3 个 server.js 进程(分别监听 3000–3002 端口),配合 compose.yaml 中定义的 postgres:14 容器即可复现一个三节点集群:
# 启动 postgres
docker compose up -d
# 运行 3 节点集群
node cluster.js
构造器与配置选项:new Emitter(pool[, nsp][, opts])
README 只说明第一个参数 pool 是 pg 包的 Pool 对象。完整的三参数签名可以从 源码构造函数 得到:
constructor(
readonly pool: any,
readonly nsp: string = "/",
opts: Partial<PostgresEmitterOptions> = {},
)
各参数与选项的默认值如下表(默认值来自 PostgresEmitterOptions 接口定义 与构造函数赋值逻辑):
| 参数 / 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
pool |
pg.Pool |
必填 | 数据库连接池,Emitter 的所有读写都经由它完成 |
nsp |
string |
"/" |
目标命名空间;后续可用 of() 切换 |
opts.channelPrefix |
string |
"socket.io" |
通知通道前缀。最终通道名为 ${channelPrefix}#${nsp},例如默认的 socket.io#/ |
opts.tableName |
string |
"socket_io_attachments" |
存放超过阈值或含二进制数据的载荷的表名 |
opts.payloadThreshold |
number |
8000 |
载荷字节阈值,超过则改走附表。PostgreSQL 的 NOTIFY 载荷上限即为 8000 字节,这是该默认值的由来 |
一个完整的构造示例:
const io = new Emitter(pool, "/admin", {
channelPrefix: "my-app", // 通道变为 my-app#/admin,便于与别的系统隔离
tableName: "io_attachments",
payloadThreshold: 4000,
});
广播 API:to / in / except / of
Emitter#to(room) / Emitter#in(room)
向指定房间广播(in 是 to 的别名):
io.to("room1").emit("hello");
从 源码 可以看到 to() 与 in() 都返回一个不可变的 BroadcastOperator 新实例(内部用 Set 累积房间名),因此支持链式追加:io.to("a").to("b").emit(...) 会同时定向到 a、b 两个房间,且房间参数既可以是字符串也可以是字符串数组。
Emitter#except(room)
将某房间排除在广播之外:
io.except("room2").emit("hello");
测试用例 验证了这一行为:三个客户端中仅加入 room1 的那个收不到 test 事件,另外两个正常收到。
Emitter#of(namespace)
切换目标命名空间,返回一个绑定该命名空间的新 Emitter:
const customNamespace = io.of("/custom");
customNamespace.emit("hello");
of() 的 实现 会自动为不带前缀的命名空间补上 /。注意它只复制了连接池(channelPrefix、tableName 等构造选项不会带过去,与 Redis emitter 的行为保持一致时请注意这一点)。
emit 的返回值与保留事件名
emit() 始终同步返回 true(表示事件已排队发布),真正的数据库写入是异步的;发布失败时错误通过 error 事件抛出,例如 publish 的 catch 分支:
io.on("error", (err) => { /* 发布到 PostgreSQL 失败 */ });
另外,connect、connect_error、disconnect、disconnecting、newListener、removeListener 属于 保留事件名,对 emit 使用这些名字会直接抛出 '"xxx" is a reserved event name' 错误。
volatile 与 compress 修饰符
README 未展开,但 源码 同时支持这两个修饰符,语义与服务端 BroadcastOperator 一致:
// 允许丢弃:客户端暂不可接收时(如 long polling 请求-响应周期中)直接放弃
io.volatile.emit("fast-update");
// 压缩发送
io.compress(true).emit("big-payload", data);
标志位最终随 opts.flags 一起写入广播文档,由集群中的服务器执行。
实用方法:房间管理与批量断开
这四个方法在 README 中都有示例,它们同样以 BroadcastOperator 形式工作,因此可以组合 of / in / except 精确圈定作用范围。
Emitter#socketsJoin(rooms)
让匹配的 socket 实例加入指定房间:
// 让所有 Socket 实例加入 "room1" 房间
io.socketsJoin("room1");
// 让 "admin" 命名空间中位于 "room1" 房间的所有 Socket 实例加入 "room2"
io.of("/admin").in("room1").socketsJoin("room2");
源码 将其编码为 SOCKETS_JOIN 类型文档,其中 opts.rooms / opts.except 记录筛选条件,rooms 是目标房间数组(字符串会被规整为数组)。测试 进一步验证了两个边界:按房间筛选(in("room1").socketsJoin("room2") 只影响已加入 room1 的 socket)以及按单个 socket ID 定位(emitter.in(serverSocket.id).socketsJoin("room3") 只让该实例加入房间)。
Emitter#socketsLeave(rooms)
让匹配的 socket 实例离开指定房间:
// 让所有 Socket 实例离开 "room1" 房间
io.socketsLeave("room1");
// 让 "admin" 命名空间中位于 "room1" 房间的所有 Socket 实例离开 "room2"
io.of("/admin").in("room1").socketsLeave("room2");
测试 覆盖了"全部离开""按房间筛选离开"与"按 socket ID 离开"三种形态,例如 emitter.in(serverSockets[1].id).socketsLeave("room3") 后只有目标实例的 rooms 不再包含 room3。
Emitter#disconnectSockets(close)
让匹配的 socket 实例断开连接:
// 让所有 Socket 实例断开
io.disconnectSockets();
// 让 "admin" 命名空间中位于 "room1" 房间的所有 Socket 实例断开
io.of("/admin").in("room1").disconnectSockets();
// 用单个 socket ID 也可以
io.of("/admin").in(theSocketId).disconnectSockets();
close 参数默认 false,表示是否关闭底层传输连接(对应 源码 中的 DISCONNECT_SOCKETS 文档字段)。测试 验证了客户端收到的断开原因为 "io server disconnect",说明走的是服务端主动断开路径。
Emitter#serverSideEmit(ev, ...args)
向集群中的每个 Socket.IO 服务器实例本身发事件(发给服务器,而不是发给浏览器客户端):
io.serverSideEmit("ping");
服务器端可以用 io.of("/").on("ping", cb) 或 io.engine.on(...) 之外的常规 EventEmitter 方式监听。测试 中,三个服务器实例都收到了 hello 事件及完整参数。一个重要限制:源码 检测到最后一个参数是函数时会直接抛出 "Acknowledgements are not supported"——serverSideEmit 不支持回调确认。
底层原理:pg_notify 内联与附表双路径发布
这是理解整个包最关键的部分。BroadcastOperator 的所有方法最终都汇入 publish 私有方法,其决策逻辑为:
- 为文档打上
uid: "emitter"标记(用于接收端区分来源); - 若文档类型是
BROADCAST/SERVER_SIDE_EMIT/SERVER_SIDE_EMIT_RESPONSE且参数中含有二进制数据(ArrayBuffer、ArrayBufferView,经 hasBinary 递归判断,支持嵌套对象、数组与toJSON),走附表路径; - 否则先
JSON.stringify,若字节长度超过payloadThreshold(默认 8000,正是 PostgreSQLNOTIFY的载荷上限),也走附表路径; - 两种情况都不命中时,直接执行
SELECT pg_notify($1, $2),把完整 JSON 文档作为通知载荷推入通道。
附表路径(publishWithAttachment)则分两步:
- 用
@msgpack/msgpack的encode将文档(含二进制)编码为二进制,INSERT INTO ${tableName} (payload) VALUES ($1) RETURNING id,取得附件id; - 再通过
pg_notify发送一个轻量 JSON 头{ uid, type, attachmentId },接收端拿到头之后回表读取完整载荷。
小载荷: Emitter --JSON 文档--> pg_notify(通道) --> 各服务器 adapter --> 客户端
大/二进制: Emitter --msgpack 写入 socket_io_attachments--> INSERT
Emitter --{type, attachmentId} 头--> pg_notify(通道) --> 服务器回表取数据 --> 客户端
文档类型由 EventType 枚举 定义:INITIAL_HEARTBEAT、HEARTBEAT、BROADCAST、SOCKETS_JOIN、SOCKETS_LEAVE、DISCONNECT_SOCKETS、FETCH_SOCKETS、FETCH_SOCKETS_RESPONSE、SERVER_SIDE_EMIT、SERVER_SIDE_EMIT_RESPONSE。其中心跳相关类型用于 emitter 与 adapter 之间的存活感知。
这种设计的直接推论:
- 普通 JSON 事件不会污染附表,附表只在真正需要时增长(生产环境可基于
created_at定期清理); - 二进制广播(如推送图片块、
Buffer参数)被透明地路由到附表,调用方无需感知; - 调试时设置
DEBUG=socket.io-postgres-emitter(调试前缀来自 源码)即可看到每次发布的通道与类型日志。
端到端验证:多节点集群测试
test/index.ts 本身就是一份可运行的验收脚本:beforeEach 中创建同一 Pool,建好 socket_io_attachments 表,随后 fork 出 NODES_COUNT = 3 台 Server,每台都通过 io.adapter(createAdapter(pool)) 挂上 postgres-adapter,并各自连接一个浏览器端客户端。它覆盖了全部公开 API:
| 测试 | 断言 |
|---|---|
emitter.emit("test", 1, "2", Buffer.from([3, 4])) |
3 个客户端全部收到,且第 3 个参数 Buffer.isBuffer 为 true(验证二进制附表路径) |
emitter.of("/custom").emit("test") |
仅 /custom 命名空间的客户端收到 |
emitter.to("room1").emit("test") |
只有加入 room1 的客户端收到 |
emitter.of("/").except("room1").emit("test") |
room1 内的客户端收不到,其余收到 |
emitter.socketsJoin / socketsLeave |
按房间或按 socket ID 精确生效 |
emitter.disconnectSockets() |
3 个客户端均收到 disconnect,原因为 "io server disconnect" |
emitter.serverSideEmit("hello", "world", 1, "2") |
3 台服务器实例均收到且参数无损 |
运行方式见 package.json 的 test 脚本:nyc mocha --import=tsx --timeout 5000 test/index.ts,需先按 compose.yaml 启动一个 POSTGRES_PASSWORD=changeit 的 postgres:14 实例。
使用限制与注意事项
结合 README 与源码,落地时需要注意以下边界:
- 必须搭配 postgres-adapter:Emitter 只负责"发送",若 Socket.IO 服务器没有配置
@socket.io/postgres-adapter,通知不会被任何进程消费; - 通道与表名两端一致:
channelPrefix、tableName等选项决定 emitter 写到哪里,adapter 侧读取的通道/表必须与之匹配,否则消息"发出去即石沉大海"; - 8000 字节阈值:这是 PostgreSQL
NOTIFY载荷的硬上限(源码注释中即引用了 PostgreSQL 官方文档说明),超过的部分自动落表,属于正常机制而非错误; serverSideEmit不支持确认回调:参数末尾传函数会抛错,需要"确认"语义时应改用其他通信机制;- 发布是尽力而为的异步操作:
emit同步返回true只代表入队成功;数据库错误通过error事件异步暴露,生产代码应挂上监听。
小结
@socket.io/postgres-emitter 用极小的 API 面(一个 Emitter 类 + BroadcastOperator 链式方法)解决了"集群外进程触达 Socket.IO 客户端"的问题:小载荷直接 pg_notify,大载荷与二进制数据经 msgpack 编码写入 socket_io_attachments 附表再发轻量头部,从而绕开 PostgreSQL NOTIFY 的 8000 字节限制。与 @socket.io/postgres-adapter 配对后,广播、命名空间、房间筛选/排除、批量进出房间、批量断连与服务器侧事件都能在独立进程中完成,并有完整的三节点集群测试作为行为依据。
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
