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

如何用RxJS实现Web Worker队列并行处理请求且不阻塞?

问题分析与解决方案

嘿,我瞅了下你的代码,问题根源其实出在两个关键地方——请求的串行执行和Worker的批量调度方式,这才导致用1个和4个Worker没啥性能差异。

先说说你的代码为啥没发挥多Worker的作用:

  1. concatMap导致请求串行:你用了concatMap来发起API请求,这意味着必须等前一个请求完全返回,才会发下一个请求。哪怕你有4个Worker,它们也得等着响应一个接一个过来,根本没法同时干活,自然和单Worker效率差不多。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:41:36