Node.js Cluster:如何等待所有工作进程的特定响应并汇总结果?
如何在Node.js Cluster中一次请求获取所有工作进程状态
这个需求很常见,核心是要把每个工作进程的响应转换成可等待的异步任务,而不是单纯靠事件监听器瞎等。我们可以用Promise来封装每个worker的回复逻辑,再结合Promise.all来等待所有结果完成,最后汇总返回给客户端。
下面直接上可运行的代码示例,我会一步步解释关键部分:
主进程代码(带Express请求处理)
const cluster = require('cluster'); const express = require('express'); const app = express(); const numCPUs = require('os').cpus().length; if (cluster.isPrimary) { console.log(`主进程 ${process.pid} 正在运行`); // 启动所有工作进程 for (let i = 0; i < numCPUs; i++) { cluster.fork(); } cluster.on('exit', (worker, code, signal) => { console.log(`工作进程 ${worker.process.pid} 已退出`); }); // 处理获取所有worker状态的请求 app.get('/worker-stats', async (req, res) => { try { // 获取所有在线的工作进程 const workers = Object.values(cluster.workers).filter(worker => worker.isConnected()); if (workers.length === 0) { return res.json({ status: 'success', data: [], message: '无在线工作进程' }); } // 给每个worker发消息,并创建对应的Promise const workerPromises = workers.map(worker => { return new Promise((resolve, reject) => { // 设置超时,避免某个worker挂了导致请求一直挂着 const timeoutId = setTimeout(() => { reject(new Error(`工作进程 ${worker.process.pid} 未及时响应`)); }, 5000); // 监听worker的回复,注意要匹配对应的worker id const handleMessage = (msg) => { if (msg.type === 'worker-status' && msg.workerId === worker.id) { clearTimeout(timeoutId); worker.off('message', handleMessage); // 移除监听器,避免内存泄漏 resolve({ workerId: worker.id, pid: worker.process.pid, status: msg.status }); } }; worker.on('message', handleMessage); // 向worker发送获取状态的指令 worker.send({ type: 'get-status' }); }); }); // 等待所有worker的响应 const allStats = await Promise.all(workerPromises); res.json({ status: 'success', data: allStats }); } catch (err) { res.status(500).json({ status: 'error', message: err.message }); } }); app.listen(3000, () => { console.log('服务器运行在 http://localhost:3000'); }); } else { // 工作进程代码 console.log(`工作进程 ${process.pid} 已启动`); // 模拟工作进程的状态数据 let workerStatus = { uptime: process.uptime(), memoryUsage: process.memoryUsage() }; // 定时更新状态(可选,模拟状态变化) setInterval(() => { workerStatus = { uptime: process.uptime(), memoryUsage: process.memoryUsage() }; }, 1000); // 监听主进程的消息 process.on('message', (msg) => { if (msg.type === 'get-status') { // 向主进程返回状态,带上worker id以便主进程匹配 process.send({ type: 'worker-status', workerId: cluster.worker.id, status: workerStatus }); } }); }
关键逻辑解释
Promise封装worker响应:
给每个worker发送消息时,创建一个Promise。当收到该worker的对应回复时,resolve这个Promise;如果超时(比如5秒),则reject,避免请求无限挂起。匹配worker回复:
每个回复都带上workerId,主进程通过这个id来区分是哪个worker的回复,确保不会把A worker的回复算到B worker头上。清理监听器:
一旦收到对应worker的回复,就移除message事件监听器,避免内存泄漏——因为每个请求都会创建新的监听器,不用的话要及时清理。过滤离线worker:
用worker.isConnected()过滤掉已经断开连接的worker,避免给已经退出的worker发消息。超时处理:
给每个Promise设置超时时间,防止个别worker异常导致整个请求卡住。
测试方式
启动服务后,访问http://localhost:3000/worker-stats,就能看到所有在线工作进程的状态汇总结果了。
内容的提问来源于stack exchange,提问作者Bobu
相关产品推荐
相关产品推荐

