如何在Ballerina中动态创建Workers并等待其完成?
动态创建Worker并等待全部执行完成
Web Worker 场景(浏览器环境)
核心思路是用Promise数组跟踪每个Worker的执行状态,再通过Promise.all等待所有任务完成。步骤如下:
- 运行时确定Worker数量(比如从用户输入、接口返回获取)
- 循环创建Worker实例,将每个Worker的执行逻辑包装为Promise
- 用
Promise.all等待所有Promise完成,统一处理结果或错误
示例代码
主线程代码:
// 运行时动态确定Worker数量(这里示例为5,可替换为实际逻辑) const workerCount = 5; const workerPromises = []; for (let i = 0; i < workerCount; i++) { // 创建Worker实例,指定Worker脚本路径 const worker = new Worker('worker.js'); // 向Worker发送任务数据 worker.postMessage({ taskId: i, data: `任务${i}的原始数据` }); // 将Worker的执行流程包装为Promise const workerPromise = new Promise((resolve, reject) => { worker.onmessage = (e) => { worker.terminate(); // 任务完成后销毁Worker,释放资源 resolve(e.data); }; worker.onerror = (error) => { worker.terminate(); reject(new Error(`Worker ${i}执行失败: ${error.message}`)); }; }); workerPromises.push(workerPromise); } // 等待所有Worker执行完成 async function waitForAllWorkers() { try { const allResults = await Promise.all(workerPromises); console.log('所有Worker任务完成,结果:', allResults); // 后续业务逻辑处理 } catch (err) { console.error('任务执行出错:', err); } } waitForAllWorkers();
Worker脚本(worker.js):
self.onmessage = (e) => { const { taskId, data } = e.data; // 模拟耗时任务处理(比如数据计算、文件解析等) const processedData = `任务${taskId}处理完成:${data} -> 加工后数据`; // 将结果返回主线程 self.postMessage(processedData); };
Node.js Worker线程场景
Node.js 提供worker_threads模块实现多线程,实现逻辑和Web Worker类似,同样用Promise数组+Promise.all来管理。
示例代码
主线程代码:
const { Worker } = require('worker_threads'); // 运行时动态确定Worker数量 const workerCount = 3; const workerPromises = []; for (let i = 0; i < workerCount; i++) { const workerPromise = new Promise((resolve, reject) => { const worker = new Worker('./worker-node.js', { workerData: { taskId: i, data: `Node任务${i}的原始数据` } }); worker.on('message', (result) => { worker.terminate(); resolve(result); }); worker.on('error', (err) => { worker.terminate(); reject(new Error(`Worker ${i}执行失败: ${err.message}`)); }); worker.on('exit', (code) => { if (code !== 0) { reject(new Error(`Worker ${i}异常退出,退出码: ${code}`)); } }); }); workerPromises.push(workerPromise); } // 等待所有Worker完成 (async () => { try { const allResults = await Promise.all(workerPromises); console.log('所有Node Worker任务完成,结果:', allResults); } catch (err) { console.error('任务执行出错:', err); } })();
Worker脚本(worker-node.js):
const { parentPort, workerData } = require('worker_threads'); const { taskId, data } = workerData; // 模拟任务处理 const processedResult = `Node任务${taskId}处理完成:${data} -> 加工后数据`; // 返回结果到主线程 parentPort.postMessage(processedResult);
额外注意事项
- 动态数量来源:
workerCount可以从任何运行时逻辑获取,比如用户输入、接口响应、计算得出的数值,直接赋值即可。 - 资源清理:任务完成后务必调用
terminate()销毁Worker,避免内存泄漏。 - 容错处理:如果需要等待所有Worker执行完毕(无论成功失败),可以用
Promise.allSettled替代Promise.all,遍历结果分别处理成功和失败的情况:
const results = await Promise.allSettled(workerPromises); results.forEach((result, index) => { if (result.status === 'fulfilled') { console.log(`Worker ${index}执行成功:`, result.value); } else { console.error(`Worker ${index}执行失败:`, result.reason); } });
内容的提问来源于stack exchange,提问作者ThisaruG
相关产品推荐
相关产品推荐

