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

如何正确使用exhaustMap:等待服务调用返回后再发送新请求

解决RxJS请求并发/顺序控制问题:exhaustMap与concatMap的实践

嗨,Giannis!你遇到的按钮重复触发导致请求混乱的问题,用RxJS的算子完全可以优雅解决。你提到的exhaustMap确实是对应第一种需求的合适选择,另外如果要实现按顺序处理所有请求,concatMap会更适用。我来帮你一步步改造代码:

一、用exhaustMap实现「忽略后续请求直到当前请求完成」

这个方案的核心是把按钮触发的请求转化为一个Observable流,通过exhaustMap确保同一时间只有一个请求在处理,后续的触发会被直接忽略直到当前请求结束。

1. 改造组件代码(APs.component.ts)

我们需要创建一个Subject来统一接收按钮触发的信号,然后在流里用exhaustMap控制请求:

import { Subject } from 'rxjs';
import { exhaustMap, tap, finalize } from 'rxjs/operators';

export class APsComponent implements OnInit {
  // 创建一个Subject作为请求触发源
  private gridInitialiseTrigger$ = new Subject<void>();

  constructor(private APsService: APsService) {}

  ngOnInit() {
    // 配置流的处理逻辑
    this.gridInitialiseTrigger$.pipe(
      // exhaustMap:当前请求未完成时,忽略新的触发
      exhaustMap(() => {
        this.APsService.isLoading = true;
        // 返回fetch请求的Observable,exhaustMap会等待它完成再处理下一个信号
        return this.APsService.fetch(this.token, this.tableName, this.state).pipe(
          tap((responseData) => {
            // 处理请求成功的逻辑
            this.APsService.next(responseData);
            this.APsService.mediatorService.sendMessage("APsRefreshed");
          }),
          // 不管成功失败,最终都把loading设为false
          finalize(() => {
            this.APsService.isLoading = false;
          })
        );
      })
    ).subscribe({
      error: (err) => {
        console.log("err", err);
      }
    });
  }

  public initialiseGrid() {
    // 不再直接调用query,而是通过Subject发送触发信号
    this.gridInitialiseTrigger$.next();
  }
}

2. 调整服务代码(APs.service.ts)

原来的query方法可以简化或者移除,因为我们已经把订阅逻辑移到了组件的流里,fetch方法保持基本不变,只需要确保它返回干净的Observable:

public fetch(token: string, tableName: string, state: any): Observable<any> {
  let queryStr = `${toODataString(state)}&$count=true`;
  queryStr = queryStr.replace(/('[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}')/ig, function (x) {
      return x.substring(1, x.length - 1);
  });
  queryStr = queryStr.replace(/substringof\((.+),(.*?)\)/, "contains($2,$1)");
  const regex = /T00:00:00\.000Z/gi;
  const noTimeZoneQueryStr = queryStr.replace(regex, '');

  return this.http
      .get(`${this.BASE_URL}/${token}/${tableName}&${noTimeZoneQueryStr}`)
      .pipe(
          map((response: any) => {
              // 这里可以保持原来的data赋值逻辑
              this.data = {
                  data: response['value'],
                  total: parseInt(response['@odata.count'], 10)
              };
              return this.data;
          })
      );
}

二、用concatMap实现「严格按请求发起顺序处理响应」

如果你需要的是即使连续点击,所有请求都按顺序发起,响应也按顺序处理,只需要把上面代码里的exhaustMap换成concatMap即可:

// 把exhaustMap替换成concatMap
concatMap(() => {
  this.APsService.isLoading = true;
  return this.APsService.fetch(this.token, this.tableName, this.state).pipe(
    tap((responseData) => {
      this.APsService.next(responseData);
      this.APsService.mediatorService.sendMessage("APsRefreshed");
    }),
    finalize(() => {
      this.APsService.isLoading = false;
    })
  );
})

concatMap会把所有触发的请求按顺序排队,前一个请求完成后才会发起下一个,响应也会严格按照请求发起的顺序返回处理。

关键说明

  • 为什么原来的写法会出问题:原来每次调用query都会直接订阅一个新的Observable,多个请求之间没有任何并发控制,所以会同时发起,响应回来的顺序可能混乱,或者同时修改状态。
  • finalize算子的作用:它会在Observable完成或报错时都执行回调,比原来在subscribe的next和error里分别设置isLoading更简洁可靠。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 10:03:11