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

Node.js Worker Threads:任务未完成却触发回调问题排查

Node.js Worker Pool 异步任务回调提前触发问题排查与解决

问题描述

使用Node.js 14的Worker Pool处理任务时,首次运行正常,但第二次调用时,主线程的任务回调先触发(打印"This run took X ms"),而Worker中的任务尚未执行完毕(后续才打印"Done")。怀疑问题与Promise使用、WorkerPool实现有关,同时不确定Worker中是否支持异步逻辑。

代码复现

workerPool.js

const { AsyncResource } = require("async_hooks");
const { EventEmitter } = require("events");
const path = require("path");
const { Worker } = require("worker_threads");

const kTaskInfo = Symbol("kTaskInfo");
const kWorkerFreedEvent = Symbol("kWorkerFreedEvent");

const { MONGODB_URI } = process.env;

class WorkerPoolTaskInfo extends AsyncResource {
  constructor(callback) {
    super("WorkerPoolTaskInfo");
    this.callback = callback;
  }

  done(err, result) {
    console.log("<<<<<<<<<<<");
    this.runInAsyncScope(this.callback, null, err, result);
    this.emitDestroy(); // TaskInfos are used only once.
  }
}

class WorkerPool extends EventEmitter {
  constructor(numThreads, workerFile) {
    super();
    this.numThreads = numThreads;
    this.workerFile = workerFile;
    this.workers = [];
    this.freeWorkers = [];

    for (let i = 0; i < numThreads; i++) this.addNewWorker();
  }

  addNewWorker() {
    const worker = new Worker(path.resolve(__dirname, this.workerFile), {
      workerData: { MONGODB_URI },
    });
    worker.on("message", (result) => {
      worker[kTaskInfo].done(null, result);
      worker[kTaskInfo] = null;
      this.freeWorkers.push(worker);
      this.emit(kWorkerFreedEvent);
    });
    worker.on("error", (err) => {
      if (worker[kTaskInfo]) worker[kTaskInfo].done(err, null);
      else this.emit("error", err);
      this.workers.splice(this.workers.indexOf(worker), 1);
      this.addNewWorker();
    });
    this.workers.push(worker);
    this.freeWorkers.push(worker);
    this.emit(kWorkerFreedEvent);
  }

  runTest(data, callback) {
    if (this.freeWorkers.length === 0) {
      console.log("No free threads. Process queued.");
      this.once(kWorkerFreedEvent, () => this.runTest(data, callback));
      return;
    }

    const worker = this.freeWorkers.pop();
    worker[kTaskInfo] = new WorkerPoolTaskInfo(callback);
    worker.postMessage(data);
  }

  close() {
    for (const worker of this.workers) worker.terminate();
  }
}

module.exports = WorkerPool;

index.js

return new Promise(async (resolve, reject) => {
  const startTime = new Date();
  threadPool.runTest(
    { ...some data },
    (err, result) => {
      if (err) {
        console.error("Bad error from test:", err);
        reject(err)
      }

      const endTime = new Date();

      console.log("Run Test Complete");
      console.log(`This run took ${endTime - startTime} ms`);

      resolve(true);
    }
  );
});

Worker代码

const { parentPort, threadId, workerData } = require("worker_threads");

parentPort.on("message", async (data) => {
  // await func1(...)
  // await func2(...)
  
  console.log("Done");
  parentPort.postMessage('Done stuff');
});

问题原因

  1. 任务信号与上下文绑定错位:
    Worker的message事件回调为async函数,若第一次任务的异步操作出现延迟,导致postMessage发送时机晚于主线程将Worker标记为空闲并分配第二次任务的时机,此时worker[kTaskInfo]已被替换为第二次任务的回调,第一次任务的postMessage会错误触发第二次任务的回调,造成“回调提前触发”的假象。

  2. Worker未严格串行处理任务:
    现有WorkerPool仅通过postMessage的接收标记Worker为空闲,但Worker在await期间事件循环处于空闲状态,理论上可接收新消息,若边缘情况下主线程误判Worker状态,会导致任务并发处理。

  3. Worker支持异步逻辑:
    Node.js Worker完全兼容异步/Promise逻辑,问题并非来自异步本身,而是任务完成信号与上下文的绑定机制存在漏洞。

解决方案

方案1:Worker内添加串行处理标志

修改Worker代码,确保同一时间只处理一个任务:

const { parentPort, threadId, workerData } = require("worker_threads");

let isBusy = false;

parentPort.on("message", async (data) => {
  if (isBusy) return;
  isBusy = true;
  
  try {
    // await func1(...)
    // await func2(...)
    
    console.log("Done");
    parentPort.postMessage('Done stuff');
  } finally {
    isBusy = false;
  }
});

方案2:添加任务ID绑定上下文

通过唯一任务ID确保信号与任务严格匹配:

修改workerPool.js的runTest方法

runTest(data, callback) {
  if (this.freeWorkers.length === 0) {
    console.log("No free threads. Process queued.");
    this.once(kWorkerFreedEvent, () => this.runTest(data, callback));
    return;
  }

  const worker = this.freeWorkers.pop();
  const taskId = `${Date.now()}-${Math.random().toString(36).slice(2)}`;
  worker[kTaskInfo] = { taskId, callback: new WorkerPoolTaskInfo(callback) };
  worker.postMessage({ taskId, data });
}

修改Worker代码

parentPort.on("message", async ({ taskId, data }) => {
  try {
    // await func1(data)
    // await func2(data)
    
    console.log("Done");
    parentPort.postMessage({ taskId, result: 'Done stuff' });
  } catch (err) {
    parentPort.postMessage({ taskId, error: err.message });
  }
});

修改workerPool.js的message事件监听

worker.on("message", ({ taskId, result, error }) => {
  const taskInfo = worker[kTaskInfo];
  if (taskInfo && taskInfo.taskId === taskId) {
    if (error) {
      taskInfo.callback.done(new Error(error), null);
    } else {
      taskInfo.callback.done(null, result);
    }
    worker[kTaskInfo] = null;
    this.freeWorkers.push(worker);
    this.emit(kWorkerFreedEvent);
  }
});

方案3:Promise化WorkerPool调用

将runTest改为返回Promise的方法,简化任务管理:

runTest(data) {
  return new Promise((resolve, reject) => {
    if (this.freeWorkers.length === 0) {
      console.log("No free threads. Process queued.");
      const onFreed = () => {
        this.runTest(data).then(resolve).catch(reject);
        this.off(kWorkerFreedEvent, onFreed);
      };
      this.on(kWorkerFreedEvent, onFreed);
      return;
    }

    const worker = this.freeWorkers.pop();
    const taskId = `${Date.now()}-${Math.random().toString(36).slice(2)}`;
    worker[kTaskInfo] = { taskId, resolve, reject };
    worker.postMessage({ taskId, data });
  });
}

// 修改message事件监听
worker.on("message", ({ taskId, result, error }) => {
  const taskInfo = worker[kTaskInfo];
  if (taskInfo && taskInfo.taskId === taskId) {
    error ? taskInfo.reject(new Error(error)) : taskInfo.resolve(result);
    worker[kTaskInfo] = null;
    this.freeWorkers.push(worker);
    this.emit(kWorkerFreedEvent);
  }
});

index.js调用示例

const startTime = new Date();
try {
  await threadPool.runTest({ ...some data });
  const endTime = new Date();
  console.log("Run Test Complete");
  console.log(`This run took ${endTime - startTime} ms`);
  return true;
} catch (err) {
  console.error("Bad error from test:", err);
  throw err;
}

总结

Worker中完全可以使用异步/Promise逻辑,问题核心是任务完成信号与上下文的绑定错位。通过上述方案,可有效解决回调提前触发的问题,提升WorkerPool的稳定性。

内容的提问来源于stack exchange,提问作者jlh77

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 00:35:26