如何在Express+ws集群环境中向不同Worker的所有客户端广播消息?
跨Worker的WebSocket广播解决方案
问题场景
在Express应用中使用cluster模块和ws包搭建WebSocket服务时,客户端连接会被分配到不同的Worker进程处理。当前实现的广播功能只能将消息发送给同一Worker进程下的客户端,无法覆盖其他Worker上的连接,而单进程模式下功能正常。
核心原因
cluster模式下每个Worker都是独立的Node.js进程,各自维护自己的wss.clients客户端集合,进程间内存完全隔离,所以单个Worker的广播无法触达其他Worker的客户端。
实现方案
通过Master进程作为消息中转站,基于cluster的IPC(进程间通信)实现跨Worker广播:
- 当某个Worker收到客户端的广播消息时,先将消息发送给Master进程
- Master进程收到消息后,转发给所有在线的Worker进程
- 每个Worker收到Master的转发指令后,向自己管理的所有客户端广播消息
修改后的完整代码
主运行函数
const run = () => { if (cluster.isPrimary) { // Master进程:创建Worker并处理消息转发 for (let index = 0; index < numberCPUs; index++) { cluster.fork(); } // 接收Worker发来的广播请求,转发给所有其他Worker cluster.on('message', (worker, msg) => { if (msg.type === 'broadcast') { Object.values(cluster.workers).forEach(targetWorker => { if (targetWorker.id !== worker.id) { targetWorker.send({ type: 'broadcast', data: msg.data }); } }); } }); cluster.on("exit", (worker, code, signal) => { console.error(`worker ${worker.process.pid} died (${signal || code}). restarting it in a sec`); setTimeout(() => cluster.fork(), 1000); }); } else { console.log(`Worker ${cluster.worker.id} : started`); const http = require("node:http"); const WebSocket = require("ws"); const server = http.createServer(application); const wss = new WebSocket.Server({ server }); // 监听Master发来的广播指令,执行本地客户端广播 process.on('message', (msg) => { if (msg.type === 'broadcast') { broadcast(WebSocket, wss, msg.data); } }); wss.on("connection", (ws, request, client) => { console.log(`Worker ${cluster.worker.id} handled new connection !`); ws.on("message", (message) => { console.log(`Worker ${cluster.worker.id} : Received message ${message}`); // 先给当前Worker下的客户端广播 broadcast(WebSocket, wss, message); // 将消息发给Master,让其转发给其他Worker process.send({ type: 'broadcast', data: message.toString() }); }); }); server.listen(port, () => { console.log(`Server is running on ${port} ...`); }); } };
广播函数(无需修改)
const broadcast = (webSocket, wss, message) => { for (const client of wss.clients) { if (client.readyState === webSocket.OPEN) { client.send(message, { binary: false }); } } };
方案说明
- Master进程:作为消息枢纽,接收任意Worker的广播请求后,转发给所有其他Worker,避免重复向发送消息的Worker转发。
- Worker进程:
- 收到客户端消息时,先给本地客户端广播一次(减少延迟)
- 同时将消息发送给Master,由Master完成跨Worker转发
- 监听Master的广播指令,收到后向自己管理的客户端同步消息
- 注意将
message转为字符串传输,避免IPC处理二进制数据可能出现的异常(若需二进制传输,可调整为Buffer格式)
内容的提问来源于stack exchange,提问作者Alirezakvr
相关产品推荐
相关产品推荐

