Node集群中HTTP Upgrade请求转发至指定WebSocket实例的问题
解决Node Cluster下WebSocket请求转发问题
核心问题原因
直接用worker.send()传递request、socket这类对象时,Node会对对象进行JSON序列化,而这类对象包含大量不可序列化的底层句柄、流方法和内部状态,导致传递后丢失关键内容,无法正常完成WebSocket握手。
可行解决方案
方案1:Master作为反向代理直接转发Socket
不需要向Worker传递请求对象,Master解析请求参数后,直接建立到目标Worker的WebSocket服务的连接,将客户端Socket与目标服务的Socket双向转发,实现无缝握手。
示例代码:
// Master进程代码 const cluster = require('cluster'); const http = require('http'); const net = require('net'); const url = require('url'); if (cluster.isPrimary) { // 启动Worker,每个Worker对应一个WebSocket服务端口 const workerCount = 5; for (let i = 1; i <= workerCount; i++) { const worker = cluster.fork({ PORT: 8000 + i }); worker.port = 8000 + i; } // Master的HTTP服务器,作为单一入口 const server = http.createServer((req, res) => { res.end('Use WebSocket connection'); }); server.on('upgrade', (req, socket, head) => { const query = url.parse(req.url, true).query; const targetPort = 8000 + parseInt(query.id); if (isNaN(targetPort)) { socket.write('HTTP/1.1 400 Bad Request\r\n\r\n'); socket.destroy(); return; } // 连接到目标Worker的WebSocket服务 const proxySocket = net.connect(targetPort, 'localhost', () => { // 转发客户端的Upgrade请求头 socket.write(`GET ${req.url} HTTP/1.1\r\n`); Object.entries(req.headers).forEach(([key, value]) => { socket.write(`${key}: ${value}\r\n`); }); socket.write('\r\n'); if (head.length > 0) socket.write(head); // 双向转发数据 socket.pipe(proxySocket); proxySocket.pipe(socket); }); // 处理连接错误 proxySocket.on('error', (err) => { console.error(`Failed to connect to port ${targetPort}:`, err); socket.write('HTTP/1.1 503 Service Unavailable\r\n\r\n'); socket.destroy(); }); socket.on('error', (err) => { console.error('Client socket error:', err); proxySocket.destroy(); }); }); server.listen(8000, () => { console.log('Master HTTP server listening on port 8000'); }); } else { // Worker进程代码,运行WebSocket服务 const WebSocket = require('ws'); const port = process.env.PORT; const wss = new WebSocket.Server({ port }); wss.on('connection', (ws) => { console.log(`Client connected to Worker on port ${port}`); ws.on('message', (data) => { ws.send(`Received from port ${port}: ${data}`); }); }); console.log(`Worker WebSocket server running on port ${port}`); }
方案2:通过Cluster传递Socket句柄(不序列化对象)
Node的worker.send()支持传递Socket句柄作为第二个参数,此时不会对Socket进行序列化,而是直接传递底层文件描述符,Worker收到后可以重新构建Socket并触发WebSocket升级。
示例代码:
// Master进程代码 const cluster = require('cluster'); const http = require('http'); const url = require('url'); if (cluster.isPrimary) { const workers = new Map(); // 启动Worker,记录Worker与ID的映射 for (let i = 1; i <= 5; i++) { const worker = cluster.fork({ WORKER_ID: i }); workers.set(i, worker); } const server = http.createServer((req, res) => { res.end('Use WebSocket connection'); }); server.on('upgrade', (req, socket, head) => { const query = url.parse(req.url, true).query; const workerId = parseInt(query.id); const targetWorker = workers.get(workerId); if (!targetWorker) { socket.write('HTTP/1.1 404 Not Found\r\n\r\n'); socket.destroy(); return; } // 传递Upgrade事件和Socket句柄,不序列化对象 targetWorker.send( { type: 'websocket-upgrade', url: req.url, headers: req.headers }, socket // 传递Socket句柄 ); }); server.listen(8000, () => { console.log('Master HTTP server listening on port 8000'); }); } else { // Worker进程代码 const WebSocket = require('ws'); const http = require('http'); const workerId = parseInt(process.env.WORKER_ID); // 创建一个内部HTTP服务器用于处理升级(不需要监听端口) const internalServer = http.createServer(); const wss = new WebSocket.Server({ noServer: true }); wss.on('connection', (ws) => { console.log(`Client connected to Worker ${workerId}`); ws.on('message', (data) => { ws.send(`Received from Worker ${workerId}: ${data}`); }); }); // 处理Master发来的Upgrade请求 process.on('message', (msg, socket) => { if (msg.type !== 'websocket-upgrade' || !socket) return; // 手动触发内部HTTP服务器的upgrade事件 internalServer.emit('upgrade', { url: msg.url, headers: msg.headers, method: 'GET' }, socket, Buffer.alloc(0)); }); // 绑定WebSocket升级逻辑到内部服务器 internalServer.on('upgrade', (req, socket, head) => { wss.handleUpgrade(req, socket, head, (ws) => { wss.emit('connection', ws, req); }); }); console.log(`Worker ${workerId} ready to handle WebSocket connections`); }
注意事项
- 要处理端口不存在、Worker离线等异常情况,避免客户端Socket挂起或内存泄漏。
- 方案1的反向代理方式更直观,不需要Worker做额外的事件监听,适合大多数场景。
- 方案2利用Cluster的句柄传递特性,减少了一次Socket转发,但需要Worker配合处理消息事件。
内容的提问来源于stack exchange,提问作者Some Dinosaur
相关产品推荐
相关产品推荐

