如何实现进程间通信的Promise化?解决响应消息的Promise解析问题
实现基于Promise的进程间通信
当然可以实现这种发送消息后等待响应再执行后续操作的Promise化进程通信!你的思路方向是对的,核心问题就是要把每个请求对应的resolve/reject函数和job_id关联起来,这样收到响应时就能精准找到对应的Promise去完成它。
下面是修改后的完整可运行代码,同时补充了错误处理和内存泄漏防护:
const child_process = require('child_process'); let worker; let job_id = 0; // 用对象存储每个待处理请求的resolve/reject函数,key为job_id const pendingJobs = {}; // 创建子进程 worker = child_process.fork(__dirname + '/w.js'); // 监听子进程的响应消息 worker.on('message', (message) => { const { job_id } = message; // 找到当前job对应的处理函数 if (pendingJobs[job_id]) { // 用响应消息resolve对应的Promise pendingJobs[job_id].resolve(message); // 处理完成后删除条目,避免内存泄漏 delete pendingJobs[job_id]; } }); // 处理子进程出错的情况,reject所有待处理的Promise worker.on('error', (err) => { Object.values(pendingJobs).forEach(({ reject }) => { reject(new Error(`子进程出错: ${err.message}`)); }); // 清空所有待处理条目 Object.keys(pendingJobs).forEach(key => delete pendingJobs[key]); }); // Promise化的send函数 function send(data) { job_id++; const currentJobId = job_id; // 存局部变量,避免闭包引用最新job_id的问题 data.job_id = currentJobId; const promise = new Promise((resolve, reject) => { // 将当前Promise的resolve/reject函数存入pendingJobs pendingJobs[currentJobId] = { resolve, reject }; }); worker.send(data); return promise; } // 你的期望使用方式完全可以正常工作 send({ content: '测试请求' }).then((response_message) => { console.log('收到响应:', response_message); }).catch((err) => { console.error('请求失败:', err); });
关键说明
- 关联请求与处理函数:我们用
pendingJobs对象存储每个job_id对应的resolve和reject函数,而不是存储Promise实例本身。这样在收到子进程的响应时,能直接调用对应的resolve来完成Promise。 - 闭包陷阱规避:在
send函数中把当前job_id存入局部变量currentJobId,避免后续job_id自增后,闭包中的job_id指向最新值,导致请求和响应不匹配。 - 内存泄漏防护:处理完响应后立刻删除
pendingJobs中的对应条目,防止无用的函数引用长期占用内存。 - 错误处理增强:监听子进程的
error事件,当子进程出错时自动reject所有还在等待的Promise,避免请求一直处于pending状态。
额外优化建议:添加超时机制
为了避免某些请求因为网络或子进程问题无限等待,可以给每个请求添加超时处理:
function send(data, timeout = 5000) { job_id++; const currentJobId = job_id; data.job_id = currentJobId; const promise = new Promise((resolve, reject) => { pendingJobs[currentJobId] = { resolve, reject }; // 设置超时定时器 const timeoutTimer = setTimeout(() => { reject(new Error(`请求超时(job_id: ${currentJobId})`)); delete pendingJobs[currentJobId]; }, timeout); // 存下定时器,收到响应时清除 pendingJobs[currentJobId].timeoutTimer = timeoutTimer; }); worker.send(data); return promise; } // 同时修改message监听函数,清除超时定时器 worker.on('message', (message) => { const { job_id } = message; if (pendingJobs[job_id]) { clearTimeout(pendingJobs[job_id].timeoutTimer); pendingJobs[job_id].resolve(message); delete pendingJobs[job_id]; } });
内容的提问来源于stack exchange,提问作者Dawn
相关产品推荐
相关产品推荐

