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

如何在Ballerina中动态创建Workers并等待其完成?

动态创建Worker并等待全部执行完成

Web Worker 场景(浏览器环境)

核心思路是用Promise数组跟踪每个Worker的执行状态,再通过Promise.all等待所有任务完成。步骤如下:

  1. 运行时确定Worker数量(比如从用户输入、接口返回获取)
  2. 循环创建Worker实例,将每个Worker的执行逻辑包装为Promise
  3. 用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 00:40:17