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

Node.js为何提前退出未执行完代码?如何确保执行到指定日志?

Worker线程导致"DONE"随机不打印的原因与解决方法

执行下方代码时,程序会随机打印或不打印"DONE",请问这是什么原因?如何确保每次都能执行到console.log("DONE");语句?

const {Worker, isMainThread, parentPort} = require('node:worker_threads');

async function main() {
  if (isMainThread) {
    const worker = new Worker(__filename);
    let resultResolve = null;
    let resultPromise = new Promise(resolve => resultResolve = resolve);
    worker.on('message', (msg) => resultResolve(msg));
    while (await resultPromise != null) {
      resultPromise = new Promise(resolve => resultResolve = resolve);
    }
    console.log("DONE");
  } else {
    for (let i = 0; i < 10000; i++) {
      parentPort.postMessage(i);
    }
    parentPort.postMessage(null);
  }
}

main();

我猜测这是因为Worker线程退出时,主线程事件循环是否能先执行到await resultPromise语句存在竞态条件导致的。

更新1

我尝试实现一个异步生成器,产出Worker线程生成的值。更具代表性的示例如下:

const {Worker, isMainThread, parentPort} = require('node:worker_threads');

async function* fooGenerator() {
  const worker = new Worker(__filename);
  let resultResolve = null;
  let resultPromise = new Promise(resolve => resultResolve = resolve);
  worker.on('message', (msg) => resultResolve(msg));
  let result = null;
  while ((result = await resultPromise) != null) {
    resultPromise = new Promise(resolve => resultResolve = resolve);
    yield result;
  }
}

async function main() {
  if (isMainThread) {
    for await (let value of fooGenerator());
    console.log("DONE");
  } else {
    for (let i = 0; i < 10000; i++) {
      parentPort.postMessage(i);
    }
    parentPort.postMessage(null);
  }
}

main();

更新2

添加setInterval并未解决问题,程序仍然无法打印"DONE"。

const {Worker, isMainThread, parentPort} = require('node:worker_threads');

async function main() {
  if (isMainThread) {
    setInterval(() => {}, 1000);
    const worker = new Worker(__filename);
    let resultResolve = null;
    let resultPromise = new Promise(resolve => resultResolve = resolve);
    worker.on('message', (msg) => resultResolve(msg));
    while ((await resultPromise) != null) {
      resultPromise = new Promise(resolve => resultResolve = resolve);
    }
    console.log("DONE");
  } else {
    for (let i = 0; i < 10000; i++) {
      parentPort.postMessage(i);
    }
    parentPort.postMessage(null);
  }
}

main();

问题原因

你的猜测完全正确,核心问题是竞态条件:

  • Worker线程发送完最后一条null消息后会立即退出,此时主线程的事件循环可能还没来得及处理这条null消息,Worker的message监听就已被移除。
  • 未被处理的null消息对应的resultResolve永远不会被调用,主线程会一直卡在await resultPromise处,无法执行到console.log("DONE")。

解决方法

要确保主线程处理完所有消息,包括最后一条null,需要保证message事件的处理优先级高于Worker退出的清理逻辑,以下是两种可靠实现:

方法1:消息队列+Exit事件监听

通过维护消息队列,优先处理已接收的消息,同时监听Worker退出事件,避免遗漏最后一条消息:

const {Worker, isMainThread, parentPort} = require('node:worker_threads');

async function main() {
  if (isMainThread) {
    const worker = new Worker(__filename);
    const messageQueue = [];
    let resolveNext = null;

    worker.on('message', (msg) => {
      if (resolveNext) {
        resolveNext(msg);
        resolveNext = null;
      } else {
        messageQueue.push(msg);
      }
    });

    const exitPromise = new Promise(resolve => worker.on('exit', resolve));

    let result;
    while (true) {
      result = messageQueue.shift() || await Promise.race([
        new Promise(resolve => resolveNext = resolve),
        exitPromise
      ]);

      if (result === null || (result === undefined && messageQueue.length === 0)) {
        break;
      }
    }
    console.log("DONE");
  } else {
    for (let i = 0; i < 10000; i++) {
      parentPort.postMessage(i);
    }
    parentPort.postMessage(null);
  }
}

main();

方法2:利用Worker的异步迭代器

Node.js原生支持Worker的异步迭代器,内部已处理竞态问题,代码更简洁:

const {Worker, isMainThread, parentPort} = require('node:worker_threads');

async function main() {
  if (isMainThread) {
    const worker = new Worker(__filename);
    // 遍历Worker的消息异步迭代器
    for await (const msg of worker) {
      if (msg === null) break;
      // 可在此处理每个消息
    }
    console.log("DONE");
  } else {
    for (let i = 0; i < 10000; i++) {
      parentPort.postMessage(i);
    }
    parentPort.postMessage(null);
    // 必须关闭端口,否则主线程迭代器不会结束
    parentPort.close();
  }
}

main();

内容的提问来源于stack exchange,提问作者Marcin Król

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 19:53:19