Promise控制Workers并发数失效:无法按批次启动工作线程求助
问题:Promise控制Worker并发数量失效,所有Worker被一次性启动
我用Promise控制同时运行的Worker最大数量为MAX_WORKERS=3,想实现每批启动3个Worker,完成一批后再启动下一批的效果,但当前代码里所有Worker被一次性启动,达不到分批执行的目的。
原核心代码
const MAX_WORKERS = 3; let promises = []; let chunk = 3; let totalNodes = keyCountMap.length / 2; let doneCount = 0; for (let m = 0; m < keyCountMap.length; ) { let remaining = totalNodes - doneCount; let numbWorker = Math.min(chunk, remaining); for (let i = 0; i < numbWorker; i++) { promises.push(createWorker(keyCountMap[m], keyCountMap[m + 1])); doneCount++; m += 2; console.log(doneCount); if (doneCount % MAX_WORKERS == 0) { Promise.all(promises).then((response) => { console.log("one chunk finish"); }); } } }
原createWorker函数
function createWorker(data1, data2) { return new Promise((resolve) => { let worker = new Worker(); worker.onmessage = (event) => { postMessageRes = event.data; if (postMessageRes == 200) { worker.postMessage([ nodePagesString, pagesString, copcString, data1, data2, ]); } else { workerCount += 1; let position = postMessageRes[0]; let color = postMessageRes[1]; for (let i = 0; i < position.length; i++) { positions.push(position[i]); colors.push(colors[i]); } if (workerCount == MAX_WORKERS) { worker.terminate(); workerCount = 0; promises = []; } resolve(true); } }; }); }
问题根源
- 同步循环不等待Promise完成:外层
for循环是同步执行的,不会因为Promise.all的回调而阻塞,会一次性创建并启动所有Worker。 Promise.all仅注册回调,不阻塞执行:内部判断doneCount % MAX_WORKERS ==0时调用Promise.all,只是添加了完成后的回调,不会阻止循环继续创建新Worker。- 全局变量控制Worker生命周期不可靠:用全局
workerCount来管理Worker终止和数组清空,逻辑耦合性高,容易出现并发状态错误。
修复方案
要实现分批执行,必须让每批任务的Promise.all完成后再启动下一批,改用**异步循环(async/await)**实现阻塞等待,同时简化createWorker的逻辑,去掉全局变量依赖。
修改后的代码
主逻辑(异步分批执行)
const MAX_WORKERS = 3; const totalNodes = keyCountMap.length / 2; // 封装异步函数实现分批执行 async function runWorkersBatch() { let m = 0; while (m < keyCountMap.length) { const batchPromises = []; // 计算当前批次可启动的Worker数量(最多MAX_WORKERS个) const batchSize = Math.min(MAX_WORKERS, Math.floor((keyCountMap.length - m) / 2)); // 创建当前批次的Worker任务 for (let i = 0; i < batchSize; i++) { batchPromises.push(createWorker(keyCountMap[m], keyCountMap[m + 1])); m += 2; } // 等待当前批次所有Worker完成,再进入下一批 await Promise.all(batchPromises); console.log("one chunk finish"); } } // 启动分批任务 runWorkersBatch();
修改后的createWorker函数
function createWorker(data1, data2) { return new Promise((resolve) => { const worker = new Worker(); worker.onmessage = (event) => { const postMessageRes = event.data; if (postMessageRes === 200) { // Worker初始化完成,发送任务数据 worker.postMessage([ nodePagesString, pagesString, copcString, data1, data2, ]); } else { // 处理Worker返回的结果 const [newPositions, newColors] = postMessageRes; for (let i = 0; i < newPositions.length; i++) { positions.push(newPositions[i]); colors.push(newColors[i]); } // 任务完成,终止Worker并resolve worker.terminate(); resolve(true); } }; }); }
Worker代码(保持原有逻辑即可)
import { Copc, Key } from "copc"; import * as THREE from "three"; const color = new THREE.Color(); const colors = []; let firstTime = true; var nodePages, pages, receivedData, copc; let x_min, y_min, z_min, x_max, y_max, z_max, width; let positions = []; let filename = "https://s3.amazonaws.com/data.entwine.io/millsite.copc.laz"; const readPoints = (id, getters) => { let returnPoint = getXyzi(id, getters); positions.push( returnPoint[0] - x_min - 0.5 * width, returnPoint[1] - y_min - 0.5 * width, returnPoint[2] - z_min - 0.5 * width ); const vx = (returnPoint[3] / 65535) * 255; color.setRGB(vx, vx, vx); colors.push(color.r, color.g, color.b); firstTime = false; }; function getXyzi(index, getters) { return getters.map((get) => get(index)); } async function load() { // copc = await Copc.create(filename); // let scale = copc.header.scale[0]; // [x_min, y_min, z_min, x_max, y_max, z_max] = copc.info.cube; // width = Math.abs(x_max - x_min); // // let center_x = (x_min + x_max) / 2; // // let center_y = (y_min + y_max) / 2; // // let center_z = (z_min + z_max) / 2; // receivedData = await Copc.loadHierarchyPage( // filename, // copc.info.rootHierarchyPage // ); // nodePages = receivedData.nodes; // pages = receivedData.pages; postMessage(200); } async function loadData(nodes, pages, copc, myRoot, pointCount) { // console.log(copc, myRoot); const view = await Copc.loadPointDataView(filename, copc, myRoot); let getters = ["X", "Y", "Z", "Intensity"].map(view.getter); for (let j = 0; j < pointCount; j += 1) { readPoints(j, getters); } postMessage([positions, colors]); } load(); onmessage = function (message) { let nodePages = message.data[0]; let nodes = JSON.parse(nodePages); let pagesStr = message.data[1]; let pages = JSON.parse(pagesStr); let copcStr = message.data[2]; let copc = JSON.parse(copcStr); let mapIndex = message.data[3]; let pointCount = message.data[4]; let myRoot = nodes[mapIndex]; // console.log(mapIndex); loadData(nodes, pages, copc, myRoot, pointCount); };
内容的提问来源于stack exchange,提问作者Pravin Poudel
相关产品推荐
相关产品推荐

