WebSocket发布订阅总线需自定义消息格式与路由:服务端用Map管理主题-客户端映射,支持subscribe/unsubscribe/publish;客户端封装PubSubBus类,统一处理连接、回调和消息分发,并注意内存泄漏与重连。

用 WebSocket 实现简易的发布订阅总线,核心是把 WebSocket 连接当作通信通道,自己在客户端和服务器端分别封装 publish(发布)和 subscribe(订阅)逻辑,而不是直接依赖 WebSocket 原生的 send/onmessage。
服务端:用简单 Node.js + ws 库搭建消息中转中心
WebSocket 协议本身不内置“主题”或“频道”概念,需要我们约定消息格式并做路由分发。
- 安装:
npm install ws - 服务端监听连接,维护一个全局的
topics = new Map(),每个 topic 存储订阅它的 client 列表 - 收到消息时解析为
{ type: 'publish', topic: 'chat', data: {...} },然后广播给该 topic 的所有 client - 支持
subscribe和unsubscribe消息类型,动态增删 client 到对应 topic 的列表中
示例服务端片段:
const WebSocket = require('ws');
const wss = new WebSocket.Server({ port: 8080 });
const topics = new Map(); // 'chat' → [client1, client2]
<p>wss.on('connection', (ws) => {
ws.on('message', (data) => {
try {
const msg = JSON.parse(data);
if (msg.type === 'subscribe' && msg.topic) {
if (!topics.has(msg.topic)) topics.set(msg.topic, new Set());
topics.get(msg.topic).add(ws);
} else if (msg.type === 'unsubscribe' && msg.topic) {
topics.get(msg.topic)?.delete(ws);
} else if (msg.type === 'publish' && msg.topic) {
const clients = topics.get(msg.topic);
if (clients) {
const broadcastMsg = JSON.stringify({
type: 'message',
topic: msg.topic,
data: msg.data,
timestamp: Date.now()
});
clients.forEach(client => {
if (client.readyState === WebSocket.OPEN) {
client.send(broadcastMsg);
}
});
}
}
} catch (e) { /<em> 忽略解析错误 </em>/ }
});</p><p>ws.on('close', () => {
// 清理所有 topic 中的该 client
topics.forEach((clients, topic) => clients.delete(ws));
});
});</p>客户端:封装 PubSubBus 类,屏蔽底层 WebSocket 细节
让业务代码只关心 bus.subscribe('news', cb) 和 bus.publish('news', {title: '...'}),不暴露连接、重连、序列化等逻辑。
立即学习“Java免费学习笔记(深入)”;
- 内部持有一个
WebSocket实例和一个handlers = new Map(),key 是 topic,value 是回调数组 -
subscribe(topic, cb)先发 subscribe 消息给服务端,再把 cb 加入 handlers[topic] -
publish(topic, data)打包成 publish 消息发送 -
onmessage收到服务端广播后,按 topic 查找并触发所有匹配的回调 - 建议加基础重连逻辑(如断开后 3 秒自动重试)和连接状态管理
示例客户端类(简化版):
class PubSubBus {
constructor(url) {
this.url = url;
this.handlers = new Map();
this.connect();
}
<p>connect() {
this.ws = new WebSocket(this.url);
this.ws.onopen = () => console.log('Connected');
this.ws.onmessage = (e) => {
const msg = JSON.parse(e.data);
if (msg.type === 'message' && msg.topic) {
const cbs = this.handlers.get(msg.topic);
if (cbs) cbs.forEach(cb => cb(msg.data));
}
};
this.ws.onclose = () => setTimeout(() => this.connect(), 3000);
}</p><p>subscribe(topic, cb) {
if (!this.handlers.has(topic)) this.handlers.set(topic, []);
this.handlers.get(topic).push(cb);
this.ws.send(JSON.stringify({ type: 'subscribe', topic }));
}</p><p>publish(topic, data) {
this.ws.send(JSON.stringify({ type: 'publish', topic, data }));
}</p><p>unsubscribe(topic, cb) {
const cbs = this.handlers.get(topic);
if (cbs) {
const idx = cbs.indexOf(cb);
if (idx > -1) cbs.splice(idx, 1);
if (cbs.length === 0) this.handlers.delete(topic);
}
this.ws.send(JSON.stringify({ type: 'unsubscribe', topic }));
}
}</p><p>// 使用
const bus = new PubSubBus('ws://localhost:8080');
bus.subscribe('chat', (msg) => console.log('收到聊天:', msg));
bus.publish('chat', { from: 'Alice', text: 'Hi!' });</p>关键细节与避坑提醒
-
消息格式必须统一:前后端对
type、topic、data字段命名和结构要严格一致,否则路由失效 -
避免内存泄漏:每次
subscribe都存回调,务必提供unsubscribe并在组件卸载(如 React useEffect cleanup)时调用 - 服务端不校验身份:简易版默认所有 client 权限相同;生产环境需加入 token 鉴权、topic 白名单、消息大小限制等
- 不保证消息顺序和送达:WebSocket 是可靠传输,但若 client 网络闪断,中间消息会丢失;如需可靠性,得引入 ACK、重传或结合 MQTT 等协议
扩展思路(按需选配)
- 支持通配符 topic(如
user.*),用 minimatch 匹配 - 增加本地事件总线 fallback:当 WebSocket 断开时,先缓存 publish 消息,恢复后再补发(适合弱网场景)
- 服务端记录每个 client 的最后活跃时间,超时自动清理,防僵尸连接
- 前端 bus 实例支持多实例隔离(不同 URL 或 namespace),避免冲突
不复杂但容易忽略。WebSocket 提供的是管道,发布订阅是建在上面的一层语义——你定义规则,它就听话转发。


















