首页
/ Socket.IO Postgres Emitter:从独立 Node.js 进程广播与管控整个 Socket.IO 服务集群

Socket.IO Postgres Emitter:从独立 Node.js 进程广播与管控整个 Socket.IO 服务集群

2026-09-04 18:40:39作者:鲍丁臣Ursa

@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 测试用例,说明如何在真实多节点集群中落地与验证这些能力。

Socket.IO 数据包经 PostgreSQL 转发到多个服务器节点的示意图

适用场景:为什么不直接在服务器进程里 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):socketsJoinsocketsLeavedisconnectSocketsserverSideEmit

安装

npm install @socket.io/postgres-emitter pg

TypeScript 项目额外需要 @types/pgpg 是 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 只说明第一个参数 poolpg 包的 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)

向指定房间广播(into 的别名):

io.to("room1").emit("hello");

源码 可以看到 to()in() 都返回一个不可变的 BroadcastOperator 新实例(内部用 Set 累积房间名),因此支持链式追加:io.to("a").to("b").emit(...) 会同时定向到 ab 两个房间,且房间参数既可以是字符串也可以是字符串数组。

Emitter#except(room)

将某房间排除在广播之外:

io.except("room2").emit("hello");

测试用例 验证了这一行为:三个客户端中仅加入 room1 的那个收不到 test 事件,另外两个正常收到。

Emitter#of(namespace)

切换目标命名空间,返回一个绑定该命名空间的新 Emitter

const customNamespace = io.of("/custom");

customNamespace.emit("hello");

of()实现 会自动为不带前缀的命名空间补上 /。注意它只复制了连接池(channelPrefixtableName 等构造选项不会带过去,与 Redis emitter 的行为保持一致时请注意这一点)。

emit 的返回值与保留事件名

emit() 始终同步返回 true(表示事件已排队发布),真正的数据库写入是异步的;发布失败时错误通过 error 事件抛出,例如 publish 的 catch 分支

io.on("error", (err) => { /* 发布到 PostgreSQL 失败 */ });

另外,connectconnect_errordisconnectdisconnectingnewListenerremoveListener 属于 保留事件名,对 emit 使用这些名字会直接抛出 '"xxx" is a reserved event name' 错误。

volatilecompress 修饰符

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 私有方法,其决策逻辑为:

  1. 为文档打上 uid: "emitter" 标记(用于接收端区分来源);
  2. 若文档类型是 BROADCAST / SERVER_SIDE_EMIT / SERVER_SIDE_EMIT_RESPONSE 且参数中含有二进制数据ArrayBufferArrayBufferView,经 hasBinary 递归判断,支持嵌套对象、数组与 toJSON),走附表路径;
  3. 否则先 JSON.stringify,若字节长度超过 payloadThreshold(默认 8000,正是 PostgreSQL NOTIFY 的载荷上限),也走附表路径;
  4. 两种情况都不命中时,直接执行 SELECT pg_notify($1, $2),把完整 JSON 文档作为通知载荷推入通道。

附表路径(publishWithAttachment)则分两步:

  1. @msgpack/msgpackencode 将文档(含二进制)编码为二进制,INSERT INTO ${tableName} (payload) VALUES ($1) RETURNING id,取得附件 id
  2. 再通过 pg_notify 发送一个轻量 JSON 头 { uid, type, attachmentId },接收端拿到头之后回表读取完整载荷。
小载荷:  Emitter --JSON 文档--> pg_notify(通道) --> 各服务器 adapter --> 客户端
大/二进制: Emitter --msgpack 写入 socket_io_attachments--> INSERT
          Emitter --{type, attachmentId} 头--> pg_notify(通道) --> 服务器回表取数据 --> 客户端

文档类型由 EventType 枚举 定义:INITIAL_HEARTBEATHEARTBEATBROADCASTSOCKETS_JOINSOCKETS_LEAVEDISCONNECT_SOCKETSFETCH_SOCKETSFETCH_SOCKETS_RESPONSESERVER_SIDE_EMITSERVER_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 = 3Server,每台都通过 io.adapter(createAdapter(pool)) 挂上 postgres-adapter,并各自连接一个浏览器端客户端。它覆盖了全部公开 API:

测试 断言
emitter.emit("test", 1, "2", Buffer.from([3, 4])) 3 个客户端全部收到,且第 3 个参数 Buffer.isBuffertrue(验证二进制附表路径)
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.jsontest 脚本:nyc mocha --import=tsx --timeout 5000 test/index.ts,需先按 compose.yaml 启动一个 POSTGRES_PASSWORD=changeitpostgres:14 实例。

使用限制与注意事项

结合 README 与源码,落地时需要注意以下边界:

  1. 必须搭配 postgres-adapter:Emitter 只负责"发送",若 Socket.IO 服务器没有配置 @socket.io/postgres-adapter,通知不会被任何进程消费;
  2. 通道与表名两端一致channelPrefixtableName 等选项决定 emitter 写到哪里,adapter 侧读取的通道/表必须与之匹配,否则消息"发出去即石沉大海";
  3. 8000 字节阈值:这是 PostgreSQL NOTIFY 载荷的硬上限(源码注释中即引用了 PostgreSQL 官方文档说明),超过的部分自动落表,属于正常机制而非错误;
  4. serverSideEmit 不支持确认回调:参数末尾传函数会抛错,需要"确认"语义时应改用其他通信机制;
  5. 发布是尽力而为的异步操作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 配对后,广播、命名空间、房间筛选/排除、批量进出房间、批量断连与服务器侧事件都能在独立进程中完成,并有完整的三节点集群测试作为行为依据。

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