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
相关产品推荐
相关产品推荐

