websocket需在应用层实现频道订阅模型,通过subscribe/unsubscribe控制消息、服务端map维护channel→set订阅关系、客户端封装事件监听,再扩展权限校验、redis集群支持与心跳重连。

WebSocket 本身不内置频道(channel)或主题(topic)概念,要实现“订阅/取消订阅频道”的通信模型,需在应用层设计协议约定和状态管理。核心思路是:客户端发送带频道标识的控制消息(如 SUBSCRIBE),服务端维护每个连接的订阅关系,并在广播时按频道过滤目标连接。
定义轻量级通信协议
双方约定简单 JSON 消息格式,区分控制指令与业务数据:
-
SUBSCRIBE:客户端请求加入频道,例如
{"type":"SUBSCRIBE","channel":"chat/general"} -
UNSUBSCRIBE:退出频道,例如
{"type":"UNSUBSCRIBE","channel":"chat/general"} -
PUBLISH:向某频道发消息,例如
{"type":"PUBLISH","channel":"chat/general","data":{"user":"Alice","text":"Hi!"}} -
MESSAGE:服务端推送的频道消息(单播或广播),例如
{"type":"MESSAGE","channel":"chat/general","data":{...}}
服务端维护订阅关系(Node.js + ws 示例)
使用 Map 结构管理:`channel → Set
const channels = new Map(); // 'chat/general' → Set<websocket>
const clientSubs = new WeakMap(); // ws → Set<string>
ws.on('message', (data) => {
try {
const msg = JSON.parse(data);
if (msg.type === 'SUBSCRIBE' && msg.channel) {
const channel = msg.channel;
if (!channels.has(channel)) channels.set(channel, new Set());
channels.get(channel).add(ws);
let subs = clientSubs.get(ws);
if (!subs) {
subs = new Set();
clientSubs.set(ws, subs);
}
subs.add(channel);
} else if (msg.type === 'UNSUBSCRIBE' && msg.channel) {
const channel = msg.channel;
channels.get(channel)?.delete(ws);
clientSubs.get(ws)?.delete(channel);
} else if (msg.type === 'PUBLISH' && msg.channel && msg.data !== undefined) {
const channel = msg.channel;
const payload = { type: 'MESSAGE', channel, data: msg.data };
const clients = channels.get(channel);
if (clients) {
clients.forEach(client => {
if (client.readyState === WebSocket.OPEN) {
client.send(JSON.stringify(payload));
}
});
}
}
} catch (e) {
ws.send(JSON.stringify({ type: 'ERROR', reason: 'Invalid message' }));
}
});
// 断连时自动清理
ws.on('close', () => {
const subs = clientSubs.get(ws);
if (subs) {
for (const channel of subs) {
channels.get(channel)?.delete(ws);
}
}
});
</string></websocket>
客户端封装订阅逻辑(浏览器端)
避免每次手动拼 JSON,封装 subscribe()、unsubscribe() 和 publish() 方法:
class ChannelSocket {
constructor(url) {
this.ws = new WebSocket(url);
this.ws.onmessage = (e) => {
const msg = JSON.parse(e.data);
if (msg.type === 'MESSAGE' && msg.channel && msg.data !== undefined) {
const handlers = this.handlers[msg.channel] || [];
handlers.forEach(cb => cb(msg.data));
}
};
}
subscribe(channel, callback) {
this.ws.send(JSON.stringify({ type: 'SUBSCRIBE', channel }));
if (!this.handlers) this.handlers = {};
if (!this.handlers[channel]) this.handlers[channel] = [];
this.handlers[channel].push(callback);
}
unsubscribe(channel, callback) {
this.ws.send(JSON.stringify({ type: 'UNSUBSCRIBE', channel }));
const handlers = this.handlers[channel];
if (handlers) {
const idx = handlers.indexOf(callback);
if (idx > -1) handlers.splice(idx, 1);
}
}
publish(channel, data) {
this.ws.send(JSON.stringify({ type: 'PUBLISH', channel, data }));
}
}
// 使用示例
const socket = new ChannelSocket('ws://localhost:3000');
socket.subscribe('chat/general', (data) => console.log('收到:', data));
socket.publish('chat/general', { user: 'Bob', text: 'Hello world' });
进阶考虑:权限、持久化与扩展性
生产环境还需补充:
-
频道权限校验:在
SUBSCRIBE时检查用户身份(如从 JWT 或 session 中提取),拒绝未授权频道 - 服务端集群支持:单机内存存储订阅关系无法跨进程共享,需接入 Redis Pub/Sub 或使用一致性哈希 + 全局订阅表
- 心跳与重连:客户端检测断连后自动重连,并重新发送 SUBSCRIBE 消息恢复状态
-
频道命名规范:建议用斜杠分隔层级(如
room/123/chat、user/456/notify),便于路由和权限控制
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











