Node.js Worker Threads:任务未完成却触发回调问题排查
问题描述
使用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'); });
问题原因
任务信号与上下文绑定错位:
Worker的message事件回调为async函数,若第一次任务的异步操作出现延迟,导致postMessage发送时机晚于主线程将Worker标记为空闲并分配第二次任务的时机,此时worker[kTaskInfo]已被替换为第二次任务的回调,第一次任务的postMessage会错误触发第二次任务的回调,造成“回调提前触发”的假象。Worker未严格串行处理任务:
现有WorkerPool仅通过postMessage的接收标记Worker为空闲,但Worker在await期间事件循环处于空闲状态,理论上可接收新消息,若边缘情况下主线程误判Worker状态,会导致任务并发处理。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

