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

如何用RxJS批量发送请求并等待全部响应后执行后续逻辑

用RxJS实现批量请求并等待全部完成的解决方案

你的核心需求是批量发送请求,等待所有请求响应后再执行后续逻辑,原来的forEach+独立subscribe无法统一监听所有请求的完成状态,下面用RxJS的两种常用方案解决:

方案1:并行发送所有请求(最快完成)

使用forkJoin组合所有请求,它会等待所有请求全部完成后触发complete回调,适合请求之间无依赖的场景:

private getValuesForEditTariff(parameter: Parameter) {
  if (!parameter.dependParams || parameter.dependParams.length === 0) return;

  // 收集所有参数到数组
  const pars: Parameter[] = [];
  this.actionEditTariff.groups.forEach(group => pars.push(...group.parameters));

  const locator = (p: Parameter, id: number) => p.id === id;
  const service = this.actionService;
  const action = this.actionEditTariff;

  // 生成请求可观察对象数组
  const requestObservables = parameter.dependParams
    .map(depParamId => {
      const par = pars.find(p => locator(p, depParamId));
      if (!par) return null;

      par.loading = true;
      // 包装请求,处理结果和错误
      return service.getValues(depParamId, pars, action.action, action.actionScheme).pipe(
        tap({
          next: (data) => {
            par.values = data;
            par.loading = false;
            console.log("Значения получены");
          },
          error: (error) => {
            this.message = error.error.message || error.statusText;
            par.loading = false;
            this.showError();
          }
        }),
        // 可选:如果某个请求出错,不终止整个批量任务,返回EMPTY跳过
        catchError(() => EMPTY)
      );
    })
    .filter(Boolean); // 过滤无效请求

  // 等待所有请求完成
  forkJoin(requestObservables).subscribe({
    complete: () => {
      // 所有请求完成后执行这里的逻辑
      console.log("所有参数请求已处理完毕");
      // 在这里添加你需要延迟执行的代码
    }
  });
}

关键说明:

  • forkJoin会并行发送所有请求,只要有一个请求未完成就会等待;
  • 如果不需要容错(某个请求出错则整个任务终止),可以去掉catchError(() => EMPTY);
  • 需要导入forkJoin, tap, catchError, EMPTY(从rxjs或rxjs/operators)。

方案2:顺序发送请求(一个接一个)

如果请求之间有依赖,或者需要严格按顺序执行,使用from+concatMap组合,concatMap会等待前一个请求完成后再发送下一个:

private getValuesForEditTariff(parameter: Parameter) {
  if (!parameter.dependParams || parameter.dependParams.length === 0) return;

  const pars: Parameter[] = [];
  this.actionEditTariff.groups.forEach(group => pars.push(...group.parameters));

  const locator = (p: Parameter, id: number) => p.id === id;
  const service = this.actionService;
  const action = this.actionEditTariff;

  // 把参数数组转成可观察流,用concatMap顺序处理每个请求
  from(parameter.dependParams).pipe(
    concatMap(depParamId => {
      const par = pars.find(p => locator(p, depParamId));
      if (!par) return EMPTY;

      par.loading = true;
      return service.getValues(depParamId, pars, action.action, action.actionScheme).pipe(
        tap({
          next: (data) => {
            par.values = data;
            par.loading = false;
            console.log("Значения получены");
          },
          error: (error) => {
            this.message = error.error.message || error.statusText;
            par.loading = false;
            this.showError();
          }
        }),
        // 可选:出错时不中断后续请求
        catchError(() => EMPTY)
      );
    })
  ).subscribe({
    complete: () => {
      console.log("所有参数请求已处理完毕");
      // 后续逻辑
    }
  });
}

关键说明:

  • from将普通数组转换成RxJS流,每个元素依次进入concatMap;
  • concatMap的核心是顺序执行,前一个请求的可观察对象完成后,才会处理下一个元素;
  • 需要导入from, concatMap, tap, catchError, EMPTY。

为什么原来的concatMap没起作用?

你之前可能错误地在forEach里使用concatMap,但concatMap是RxJS流的操作符,必须作用于可观察流(比如from生成的流),而不是普通数组的forEach循环。

内容的提问来源于stack exchange,提问作者denikotin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 09:03:12