You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现进程间通信的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:15:23