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

RxJS:如何避免race操作符取消"落败"的HTTP请求Observable?

问题描述

我正尝试用RxJS实现一种竞争场景:需调用两个API接口,但仅当数据返回耗时超过150ms时才显示加载动画。UI流程为:等待150ms,若数据已返回则立即展示;若未返回,则显示加载动画至少约250ms。

当前逻辑及UI表现均正常,但存在一个问题:若其中一个API请求在150ms内返回,另一个延迟返回,race操作符会取消延迟的HTTP请求,导致加载完成后需重新发起该请求。

请问是否有办法让HTTP请求即使在竞争中落败也继续执行,以便通过shareReplay复现请求结果,确保每个接口仅发起一次HTTP调用?

现有代码

class GameDetailsService {
// more methods and fields are hidden as not relevant to this example.

  private state = initialState;

  private store = new BehaviorSubject<GameDetailsState>(this.state);
  private store$ = this.store.asObservable();
  selectedGame$ = this.store$.pipe(
    map((state) => state.selectedGameId),
    distinctUntilChanged(),
    switchMap((id): Observable<null | Game> => {
      if (id == null) {
        return this.clearGame();
      } else {
        return this.findGameDetails(id);
      }
    })
  );

private findGameDetails(id: string): Observable<Game> {
    const gameDetails$ = this.findGameById(id).pipe(
      catchError((err) => {
        console.error(err);
        return of(null);
      }),
      retry(2)
    );
    const screenshots$ = this.findScreenshotsById(id).pipe(
      catchError((err) => {
        console.error(err);
        return of(null);
      }),
      retry(2)
    );

    const data$ = forkJoin([gameDetails$, screenshots$]).pipe(
      filter(
        (data): data is [Game, { image: string }[]] =>
          data[0] != null && data[1] != null
      ),
      map(([gameDetails, screenshots]) => {
        // enhance game object with screenshots
        gameDetails.screenshots = screenshots;

        return gameDetails;
      }),
      tap((gameDetails) => {
        // create new state
        const newState = produce(this.state, (draft) => {
          draft.selectedGame = gameDetails;
          draft.loading = false;
        });
        this.state = newState;
        this.store.next(newState);
      }),
      shareReplay({ bufferSize: 1, refCount: true })
    );

    /**
     * We want to display the loading indicator only if the requests takes more than `${INITIAL_WAITING_TIME}
     * if it does, we want to wait at least `${MINIMUM_TIME_TO_DISPLAY_LOADER} before emitting
     */

    const startLoading$ = of({}).pipe(
      tap(() => console.log('start loading ' + new Date().getSeconds())),
      tap(() => this.startLoadingGame()),
      delay(MINIMUM_TIME_TO_DISPLAY_LOADER),
      tap(() => console.log('emit finish loading ' + new Date().getSeconds())),
      switchMap(() => EMPTY)
    );

    const hideLoading$ = of({}).pipe(
      tap(() => console.log('stop loading ' + new Date().getSeconds())),
      tap(() => this.stopLoadingGame()),
      switchMap(() => EMPTY)
    );

    const timer$ = timer(INITIAL_WAITING_TIME).pipe(
      tap(() => console.log('timer emitted ' + new Date().getSeconds()))
    );

    /**
     * We want to race two streams:
     *
     * - initial waiting time: the time we want to hold on any UI updates
     * to wait for the API to get back to us
     *
     * - data: the response from the API.
     *
     * Scenario A: API comes back before the initial waiting time
     *
     * We avoid displaying the loading spinner altogether, and instead we directly update
     * the state with the new data.
     *
     * Scenario B: API doesn't come back before initial waiting time.
     *
     * We want to display the loading spinner, and to avoid awkward flash (for example the response comes back 10ms after the initial waiting time) we extend the delay to 250ms
     * to give the user the time to understand the actions happening on the screen.
     */

     const race$ = race(timer$, data$).pipe(
      mergeMap((winner) => (typeof winner === 'number' ? startLoading$ : EMPTY))
    );

    return concat(race$, hideLoading$, data$);
  }

  private findGameById(id: string): Observable<Game> {
    return this.http.get<Game>(`${env.BASE_URL}/games/${id}`);
  }

  private findScreenshotsById(id: string): Observable<{ image: string }[]> {
    return this.http
      .get<APIResponse<{ image: string }>>(
        `${env.BASE_URL}/games/${id}/screenshots`
      )
      .pipe(map(({ results: screenshots }) => screenshots));
  }


}
解决方案

问题核心是race操作符会在其中一个流发出值后取消其他流的订阅,导致data$关联的HTTP请求被中断。要让请求持续执行,需将race的竞争逻辑与data$的主订阅解耦,具体修改如下:

1. 分离竞争信号与数据流

创建一个仅用于告知数据就绪的信号流dataSignal$,让它参与race竞争,而非直接使用data$。这样race只会监听信号,不会干扰data$的订阅生命周期。

2. 确保data$的订阅独立

由于data$使用了shareReplay,只要存在至少一个订阅,它就会保持活跃直到请求完成。后续的订阅(比如concat中的data$)会直接复用缓存的结果,无需重新发起请求。

修改后的findGameDetails方法关键代码:

private findGameDetails(id: string): Observable<Game> {
    const gameDetails$ = this.findGameById(id).pipe(
      catchError((err) => {
        console.error(err);
        return of(null);
      }),
      retry(2)
    );
    const screenshots$ = this.findScreenshotsById(id).pipe(
      catchError((err) => {
        console.error(err);
        return of(null);
      }),
      retry(2)
    );

    const data$ = forkJoin([gameDetails$, screenshots$]).pipe(
      filter(
        (data): data is [Game, { image: string }[]] =>
          data[0] != null && data[1] != null
      ),
      map(([gameDetails, screenshots]) => {
        gameDetails.screenshots = screenshots;
        return gameDetails;
      }),
      tap((gameDetails) => {
        const newState = produce(this.state, (draft) => {
          draft.selectedGame = gameDetails;
          draft.loading = false;
        });
        this.state = newState;
        this.store.next(newState);
      }),
      shareReplay({ bufferSize: 1, refCount: true })
    );

    const startLoading$ = of({}).pipe(
      tap(() => console.log('start loading ' + new Date().getSeconds())),
      tap(() => this.startLoadingGame()),
      delay(MINIMUM_TIME_TO_DISPLAY_LOADER),
      tap(() => console.log('emit finish loading ' + new Date().getSeconds())),
      switchMap(() => EMPTY)
    );

    const hideLoading$ = of({}).pipe(
      tap(() => console.log('stop loading ' + new Date().getSeconds())),
      tap(() => this.stopLoadingGame()),
      switchMap(() => EMPTY)
    );

    const timer$ = timer(INITIAL_WAITING_TIME).pipe(
      tap(() => console.log('timer emitted ' + new Date().getSeconds()))
    );

    // 新增:创建仅用于竞争的信号流,不影响data$的订阅
    const dataSignal$ = data$.pipe(mapTo('data-ready'));

    // 修改race的竞争对象为timer$和dataSignal$
    const race$ = race(timer$, dataSignal$).pipe(
      mergeMap((winner) => (typeof winner === 'number' ? startLoading$ : EMPTY))
    );

    return concat(race$, hideLoading$, data$);
  }

修改说明

  • dataSignal$通过mapTo将data$的结果转换为简单信号,仅用于通知race数据已就绪,不会传递实际数据,也不会中断data$的请求。
  • data$的订阅由concat中的data$和dataSignal$共同维持,确保HTTP请求会完整执行,结果被shareReplay缓存。
  • 原有的UI逻辑(加载动画的显示时机、最短显示时长)完全保留,不会受到影响。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 14:40:38