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

基于RxJS实现多并发观察者共享HTTP请求的原理疑问

问题

在Angular应用中实现以下功能:同一数据的多并发请求仅触发一次HTTP请求,所有并发订阅者共享该请求;第一个订阅者立即发送请求,请求未完成时新订阅者订阅同一请求,但不做缓存——请求完成后,后续请求会触发新的HTTP请求以获取最新数据(缓存由应用其他模块管理)。

当前简化实现代码如下:

class DataLoader<T> {

    private readonly requestInitiator = new Subject<void>();
    private readonly requests$: Observable<T>;

    constructor(requestCreator: () => Observable<T>) {
        this.requests$ = this.requestInitiator
            .pipe(
                mergeMap(requestCreator, 1), // concurrent = 1
                share(),
            );
    }

    load(): Observable<T> {
        return new Observable<T>(observer => {
            this.requests$.subscribe(observer);
            this.requestInitiator.next();
        });
    }
}

在Angleton单例服务中的使用示例:

@Injectable({ providedIn: 'root' })
export class SomeService {

    private readonly loader = new DataLoader<Something>(
        () => this.getHttp().get('/api/something'),
    );

    getSomething(): Observable<Something> {
        if (this.cache.isSomethingCached()) {
            return this.cache.getSomething();
        }

        return this.loader.load()
            .pipe(/* Cache if necessary */);
    }
}

当前实现可正常工作,但存在以下原理困惑:

  • mergeMap设置并发数为1为何能实现需求(原以为会像concatMap那样每个观察者触发一次请求)
  • 最初使用exhaustMap时,部分场景下观察者会陷入无输出、无完成的状态
  • 不清楚share操作符如何改变了Observable的行为

需要解释mergeMap(1)生效的原因,以及exhaustMap失效的问题。


原理解释

mergeMap(1)的生效逻辑

mergeMap的核心是将上游每个值转换为内部Observable,再合并输出到下游。设置concurrent=1后,同时最多只能有1个内部Observable处于运行状态,新的上游值会被暂存,直到当前内部Observable完成后再处理。

结合代码逻辑:

  1. 第一个订阅者调用load()时,先订阅requests$,此时因share的作用,上游requestInitiator被首次订阅;随后requestInitiator.next()发出信号,mergeMap触发对应的HTTP请求Observable,请求开始执行。
  2. 若此时有第二个订阅者调用load(),会先订阅requests$——因share已建立上游订阅,新订阅者直接加入现有订阅列表,共享当前运行的HTTP请求输出;随后发送的next()信号会被mergeMap(1)暂存,不会触发新请求。
  3. 当前HTTP请求完成后,mergeMap会处理暂存的next()信号,触发新的HTTP请求,满足“请求完成后后续请求获取最新数据”的需求。

这里的关键是,mergeMap(1)限制了同时运行的请求数量,避免并发请求触发重复HTTP调用;同时结合share的共享订阅特性,让所有并发订阅者复用同一个请求的输出,而非各自触发请求。

concatMap本质是mergeMap(1)的别名,行为完全一致。之前的误解源于混淆了“订阅者”和“上游信号”的关系:concatMap/mergeMap(1)处理的是上游发出的每个信号,而非每个订阅者,因此不会为每个订阅者单独触发请求。

share操作符的作用

默认Observable是冷Observable,每个订阅者都会触发完整的数据流执行。若没有share,每个订阅者订阅requests$时,都会重新订阅requestInitiator,导致每个load()调用都触发独立的HTTP请求,无法实现共享请求的需求。

share将冷Observable转换为连接Observable:

  1. 第一个订阅者订阅时,share会订阅上游数据流(即requestInitiator.pipe(mergeMap(...))),启动上游逻辑。
  2. 后续订阅者不会触发新的上游订阅,直接接收当前运行的上游数据流输出。
  3. 所有订阅者取消订阅后,share会取消上游订阅;若之后有新订阅者,会重新启动上游逻辑,触发新的数据流执行。

正是share让所有load()调用的订阅者共享同一个HTTP请求的输出,这是实现并发请求复用的核心基础。

exhaustMap失效的原因

exhaustMap的逻辑是:上游发出信号时,若当前无运行中的内部Observable,则执行对应的内部Observable;若已有运行中的内部Observable,则直接忽略该上游信号,不会暂存。

在业务场景中,exhaustMap会导致以下问题:

  1. 第一个订阅者调用load()触发HTTP请求后,第二个订阅者在请求未完成时调用load(),发送的next()信号会被exhaustMap忽略。
  2. 若第一个订阅者在请求完成前取消订阅,share会取消上游订阅,导致exhaustMap终止当前运行的HTTP请求;此时第二个订阅者仍处于订阅状态,但因之前的next()信号被忽略,没有新的请求被触发,最终陷入无输出、无完成的等待状态。

而mergeMap(1)会暂存上游信号,即使最早的订阅者取消订阅,只要还有其他订阅者存在,share会保持上游订阅,当前请求完成后,mergeMap会处理暂存的信号触发新请求,避免了无输出的情况。


内容的提问来源于stack exchange,提问作者Markus Yrjölä

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 20:45:10