如何用RxJS实现Web Worker队列并行处理请求且不阻塞?
问题分析与解决方案
嘿,我瞅了下你的代码,问题根源其实出在两个关键地方——请求的串行执行和Worker的批量调度方式,这才导致用1个和4个Worker没啥性能差异。
先说说你的代码为啥没发挥多Worker的作用:
- concatMap导致请求串行:你用了
concatMap来发起API请求,这意味着必须等前一个请求完全返回,才会发下一个请求。哪怕你有4个Worker,它们也得等着响应一个接一个过来,根本没法同时干活,自然和单Worker效率差不多。 - bufferCount+switchMap的批量处理浪费Worker空闲时间:
bufferCount(4)会攒够4个响应才一起分给Worker,然后combineAll()要等这4个Worker都干完才会处理下一批。哪怕某个Worker提前完成了,它也得闲着等另外3个,完全没做到“一空闲就分配新任务”的预期。
修正后的代码方案
下面是调整后的代码,完美贴合你要的“请求不阻塞、Worker一空闲就接新任务”的需求:
const octrees$ = from(spatialKeys) .pipe( // 把concatMap换成mergeMap,让API请求并发执行 mergeMap((spatialKey: string) => { console.log("REQUEST"); return this.sgsApi .getCellById(sessionId, volumetricModelQueryId, spatialKey) .then(resp => resp.arrayBuffer()) .then(datBuffer => ({ datBuffer, spatialKey })); }, // 可选:如果你的API允许更高并发,可以把这个数值调大(比如10),让更多请求同时发出 // 10 ), // 用mergeMap调度Worker,设置并发数为4,自动复用空闲Worker mergeMap(({ spatialKey, datBuffer }) => { // 这里假设你的unpackDat返回Observable/Promise,直接返回即可 // 如果需要指定空闲Worker,可以自己加逻辑找空闲的,比如: // const idleWorker = this.unpackWorkers.find(worker => !worker.isBusy); // return idleWorker.unpackDat(spatialKey, datBuffer); // 要是没做Worker状态检测,mergeMap的concurrency参数也会帮你控制最多4个任务同时跑 return this.unpackWorkers[0].unpackDat(spatialKey, datBuffer); }, 4), // 核心:并发数设为4,对应你的4个Worker // 每个Worker任务完成就会触发这里,不用等所有任务都结束 tap(val => { const t1 = performance.now(); console.log("COMPLETED at :", (t1 - t0) / 1000, val); }) ) .subscribe();
关键改动说明:
- mergeMap替代concatMap:让API请求并发发送,不用等前一个请求完成再发下一个,这样多个响应能同时到达,给Worker提供持续的任务源。
- mergeMap调度Worker并设置concurrency:4:RxJS会自动管理4个并发的Worker任务,当某个Worker完成当前任务,立刻就会给它分配下一个刚到的响应,完全实现你示意图里的“Worker一空闲就接新活”的逻辑。
- 去掉bufferCount和combineAll:不再批量攒任务,而是逐个分配,彻底避免Worker空闲等待的情况。
额外小提示:
如果你的Worker的unpackDat方法返回的是Promise,记得用from()把它转换成Observable(比如from(this.unpackWorkers[0].unpackDat(...))),这样RxJS的mergeMap能更好地控制并发。
内容的提问来源于stack exchange,提问作者Luca Marangon
相关产品推荐
相关产品推荐

