Node.js集群开发疑问:Master节点通信及Worker故障转移实现
Node.js集群Master节点通信与Worker故障转移实现方案
一、Master节点间通信实现
方案1:基于共享PostgreSQL数据库的消息传递
利用已有的共享数据库实现跨机器/本地Master间通信,步骤如下:
- 创建消息存储表
在数据库中新增master_messages表,用于存储Master间的消息:
CREATE TABLE master_messages ( id SERIAL PRIMARY KEY, sender_pid INT NOT NULL, receiver_pid INT, -- 为NULL时表示广播消息 content TEXT NOT NULL, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, is_read BOOLEAN DEFAULT FALSE );
同时新增master表维护Master状态:
CREATE TABLE master ( master_id INT PRIMARY KEY, status VARCHAR(20) DEFAULT 'Online', heartbeat TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
- Master注册与消息处理逻辑
在现有Master代码中添加注册、心跳更新、消息轮询与发送逻辑:
if (cluster.isMaster) { // 注册当前Master到数据库 pool.query( 'INSERT INTO master (master_id) VALUES ($1) ON CONFLICT (master_id) DO UPDATE SET status = $2', [process.pid, 'Online'], (err) => err && console.error('Master注册失败:', err) ); // 每5秒更新心跳 setInterval(() => { pool.query('UPDATE master SET heartbeat = CURRENT_TIMESTAMP WHERE master_id = $1', [process.pid]); }, 5000); // 每2秒轮询未读消息 setInterval(() => { pool.query( 'SELECT * FROM master_messages WHERE (receiver_pid = $1 OR receiver_pid IS NULL) AND is_read = FALSE', [process.pid], (err, result) => { if (err) return console.error('读取消息失败:', err); result.rows.forEach(msg => { console.log(`[Master ${process.pid}] 收到Master ${msg.sender_pid} 的消息: ${msg.content}`); // 标记消息为已读 pool.query('UPDATE master_messages SET is_read = TRUE WHERE id = $1', [msg.id]); }); } ); }, 2000); // 发送消息方法(支持定向/广播) function sendMasterMsg(receiverPid, content) { pool.query( 'INSERT INTO master_messages (sender_pid, receiver_pid, content) VALUES ($1, $2, $3)', [process.pid, receiverPid, content], (err) => err && console.error('发送消息失败:', err) ); } // 示例:启动后广播消息 setTimeout(() => sendMasterMsg(null, `Master ${process.pid} 已启动`), 3000); // 原有Master逻辑... }
方案2:基于UDP广播的本地通信
如果集群部署在同一机器,UDP广播更高效,无需依赖数据库:
const dgram = require('dgram'); const BROADCAST_PORT = 3000; const BROADCAST_ADDR = '255.255.255.255'; if (cluster.isMaster) { const socket = dgram.createSocket('udp4'); socket.bind(BROADCAST_PORT, () => socket.setBroadcast(true)); // 接收消息 socket.on('message', (msg, rinfo) => { if (rinfo.port !== BROADCAST_PORT) return; // 过滤自身发送的消息 console.log(`[Master ${process.pid}] 收到消息: ${msg.toString()}`); }); // 发送广播消息 function broadcastMsg(content) { const msg = Buffer.from(`Master ${process.pid}: ${content}`); socket.send(msg, 0, msg.length, BROADCAST_PORT, BROADCAST_ADDR); } // 每5秒发送心跳广播 setInterval(() => broadcastMsg('心跳'), 5000); // 原有Master逻辑... }
二、Worker检测Master存活并转为Master
实现逻辑
- Worker启动时通过环境变量获取对应Master的PID
- 定期查询数据库中Master的心跳时间,超过阈值判定为离线
- 检测到Master离线后,以Master身份重新启动自身进程
代码整合修改
在现有Worker逻辑中添加检测与切换逻辑:
else { const masterPid = process.env.MASTER_PID; const CHECK_INTERVAL = 5000; const TIMEOUT_THRESHOLD = 10000; // 10秒超时 // 定期检测Master状态 const checkTimer = setInterval(async () => { try { const result = await pool.query('SELECT heartbeat FROM master WHERE master_id = $1', [masterPid]); if (result.rows.length === 0) { switchToMaster(); return; } const lastHeartbeat = new Date(result.rows[0].heartbeat); if (new Date() - lastHeartbeat > TIMEOUT_THRESHOLD) { switchToMaster(); } } catch (err) { console.error(`[Worker ${process.pid}] 检测Master状态失败:`, err); } }, CHECK_INTERVAL); // 切换为Master进程 function switchToMaster() { console.log(`[Worker ${process.pid}] Master ${masterPid} 离线,切换为Master模式`); clearInterval(checkTimer); require('child_process').spawn(process.execPath, process.argv, { stdio: 'inherit', env: { ...process.env, IS_MASTER: 'true' } }); process.exit(0); } // 原有Worker逻辑 process.send(`Worker ${process.pid} 已启动`); }
同时修改Master启动Worker时传递环境变量:
// 替换原有cluster.fork() cluster.fork({ MASTER_PID: process.pid });
完整修改后代码示例
const { Pool } = require('pg'); const cluster = require('cluster'); const numCPUs = require('os').cpus().length; const dgram = require('dgram'); const { spawn } = require('child_process'); const dbParams = { user: 'postgres', host: 'localhost', database: 'clustersdb', password: 'root', port: 5433, }; const pool = new Pool(dbParams); const workerPool = new Pool(dbParams); // UDP广播配置(本地Master通信) const BROADCAST_PORT = 3000; const BROADCAST_ADDR = '255.255.255.255'; // 强制启动为Master的标识 const forceMaster = process.env.IS_MASTER === 'true'; if (cluster.isMaster || forceMaster) { // 重置cluster身份(针对Worker切换的场景) if (forceMaster) { cluster.isMaster = true; cluster.isWorker = false; } // Master注册与心跳维护 pool.query( 'INSERT INTO master (master_id) VALUES ($1) ON CONFLICT (master_id) DO UPDATE SET status = $2', [process.pid, 'Online'], (err) => err && console.error('Master注册失败:', err) ); setInterval(() => { pool.query('UPDATE master SET heartbeat = CURRENT_TIMESTAMP WHERE master_id = $1', [process.pid]); }, 5000); // UDP广播通信 const socket = dgram.createSocket('udp4'); socket.bind(BROADCAST_PORT, () => socket.setBroadcast(true)); socket.on('message', (msg, rinfo) => { if (rinfo.port !== BROADCAST_PORT) return; console.log(`[Master ${process.pid}] 收到消息: ${msg.toString()}`); }); const broadcastMsg = (content) => { const msg = Buffer.from(`Master ${process.pid}: ${content}`); socket.send(msg, 0, msg.length, BROADCAST_PORT, BROADCAST_ADDR); }; // 启动Worker for (let i = 0; i < numCPUs - 2; i++) { cluster.fork({ MASTER_PID: process.pid }); } // 原有Master事件监听 cluster.on('fork', (worker) => { insert(worker.process.pid, process.pid); broadcastMsg(`启动Worker ${worker.process.pid}`); }); cluster.on('exit', (worker) => { update(worker.process.pid); const newWorker = cluster.fork({ MASTER_PID: process.pid }); broadcastMsg(`Worker ${worker.process.pid} 离线,重启为 ${newWorker.process.pid}`); }); cluster.on('message', (worker, message) => { console.log(`[Master ${process.pid}] 收到Worker ${worker.process.pid} 消息: ${message}`); }); } else { const masterPid = process.env.MASTER_PID; const CHECK_INTERVAL = 5000; const TIMEOUT_THRESHOLD = 10000; // Master存活检测 const checkTimer = setInterval(async () => { try { const result = await pool.query('SELECT heartbeat FROM master WHERE master_id = $1', [masterPid]); if (result.rows.length === 0 || new Date() - new Date(result.rows[0].heartbeat) > TIMEOUT_THRESHOLD) { switchToMaster(); } } catch (err) { console.error(`[Worker ${process.pid}] 检测Master状态失败:`, err); } }, CHECK_INTERVAL); // 切换为Master const switchToMaster = () => { console.log(`[Worker ${process.pid}] Master ${masterPid} 离线,切换为Master模式`); clearInterval(checkTimer); spawn(process.execPath, process.argv, { stdio: 'inherit', env: { ...process.env, IS_MASTER: 'true' } }); process.exit(0); }; process.send(`Worker ${process.pid} 已启动`); } // 原有数据库操作函数 function insert(worker, master) { workerPool.query( 'INSERT INTO worker (worker_id,master_id,status) VALUES ($1, $2, $3)', [worker, master, 'Online'], (err) => err && console.error('插入Worker数据失败:', err) ); } function update(worker) { workerPool.query( 'UPDATE worker SET status = $1 WHERE worker_id = $2', ['Offline', worker], (err) => err && console.error('更新Worker状态失败:', err) ); } function fetch(master) { workerPool.query('SELECT * FROM worker WHERE master_id = $1', [master], (err, result) => { if (err) console.error('查询Worker数据失败:', err); else console.log(result.rows.length > 0 ? result.rows : `Master ${master} 无关联Worker`); }); }
内容的提问来源于stack exchange,提问作者azerty
相关产品推荐
相关产品推荐

