【问题标题】:RxJS parallel queue with concurrent workers and handling of each request具有并发工作人员的 RxJS 并行队列并处理每个请求
【发布时间】:2019-06-30 19:24:47
【问题描述】:

我在使用 RxJS 和处理请求数组的正确方法时遇到了问题。 假设我有一个包含大约 50 个请求的数组,如下所示:

let requestCounter = 0;
function makeRequest(timeToDelay) {
  return of('Request Complete!').pipe(delay(timeToDelay));
}

const requestArray = []
for(let i=0;i<25;i++){
  requestArray.push(makeRequest(3000)); //3 seconds request
  requestArray.push(makeRequest(1000)); //1 second request
}

我的目标是:

  • 并行启动请求
  • 同一时刻只能运行 5 个
  • 当一个请求完成后,数组中的下一个开始
  • 完成请求(成功或错误)后,我需要将变量“requestCounter”加一 (requestCounter++)
  • 当我在队列中的最后一个请求完成时,我需要订阅此事件并处理每个请求结果的数组

到目前为止,我最接近的做法是遵循这篇文章中的回复:

RxJS parallel queue with concurrent workers?

问题是我正在发现 RxJS,这个例子对我来说太复杂了,我找不到如何处理每个请求的计数器。

希望你能帮助我。 (抱歉英语不好,这不是我的母语)

编辑: 最终解决方案如下所示:

forkJoinConcurrent<T>(
    observables: Observable<T>[],
    concurrent: number
  ): Observable<T[]> {
    return from(observables).pipe(
      mergeMap((outerValue, outerIndex) => outerValue.pipe(
        tap(// my code ),
        last(),
        catchError(error => of(error)),
        map((innerValue, innerIndex) => ({index: outerIndex, value: innerValue})),
      ), concurrent),
      toArray(),
      map(a => (a.sort((l, r) => l.index - r.index).map(e => e.value))),
    );
  }

【问题讨论】:

  • 您发布的链接中的答案符合您的要求。您可以在操作员队列中使用tap 来获得诸如增加计数器之类的副作用。在您发布的链接的答案中,在 last 运算符之后完成了一个请求。 toArray 运算符将所有请求组合到一个数组中,因此您必须在队列中的某处添加 tap(_ =&gt; requestCounter++) last() 之后但在 toArray() 之前。
  • 太棒了,它完成了工作!我唯一的问题是,如果其中一个呼叫失败,所有其他呼叫都将被取消。有没有办法让它继续运行并处理最后返回的数组中的错误?
  • 您必须使用 catchError 运算符捕获每个请求的错误并返回替代值流,例如catchError(error =&gt; of(error))。在您创建请求的管道末尾添加catchError,即在您的makeRequest 函数中(首选),或者在last 之后的内部管道中添加它。

标签: javascript parallel-processing rxjs concurrent-queue


【解决方案1】:

首先您应该使用 subject 来存储请求队列,并查看 mergeMap 运算符,您可以为最大并发设置一个 concurrency 参数以及一个 index 变量来跟踪通话次数

https://www.learnrxjs.io/operators/transformation/mergemap.html

const requestArray=new Subject()
for(let i=0;i<25;i++){
  requestArray.next(makeRequest(3000)); //3 seconds request
  requestArray.next(makeRequest(1000)); //1 second request
}

requestArray.pipe(
mergeMap((res,index)=>of([res,index]),
res=>res,5),
map((res,index)=>{if(index===25) .... do your thing ; return res;})
).subscribe(console.log)

【讨论】:

    猜你喜欢
    • 2019-06-12
    • 1970-01-01
    • 1970-01-01
    • 2018-10-24
    • 2019-03-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-02-16
    相关资源
    最近更新 更多