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

RxJava1中concatMap引发MissingBackpressureException的原因与解决疑问

关于RxJava concatMap背压与队列大小的疑问

问题场景

我在业务中需要保证Observable的处理顺序,因此使用了concatMap操作符,测试代码如下:

@Test fun load_data() { 
    val sub = TestSubscriber<Long>() 
    var s = BehaviorSubject.create<Long>() 
    s.concatMap { Observable.timer(it, TimeUnit.MILLISECONDS) } 
        .take(4) 
        .subscribe(sub) 
    s.onNext(5) 
    s.onNext(6) 
    s.onNext(7) 
    s.onNext(8) 
    // 这里抛出rx.exceptions.MissingBackpressureException
    sub.awaitTerminalEvent(500, TimeUnit.MILLISECONDS) 
    sub.assertNoErrors() 
}

注:为简化复现,我将实际数据加载替换为Observable.timer();应用中使用BehaviorSubject关联UI操作与Rx流。

根据文档的弹珠图,我预期concatMap会将元素存入队列并逐个处理,但实际它的队列大小似乎仅为2,添加更多元素会引发MissingBackpressureException。现提出两个问题:

  1. 为何concatMap的队列大小为2而非其他算子使用的RxRingBuffer.SIZE?
  2. 是否必须在调用concatMap前添加onBackpressure*算子来避免MissingBackpressureException?

回答

1. 为什么concatMap的队列大小是2而非RxRingBuffer.SIZE?

这是由concatMap的核心设计和背压处理逻辑决定的:

concatMap的核心是严格按顺序处理每个上游事件生成的Observable——它必须等前一个Observable完全完成后,才会订阅下一个由上游事件生成的Observable。为了保证顺序且避免过度内存占用,它的内部背压策略非常保守:

  • 当它正在处理某个上游事件对应的Observable时,只会缓冲1个后续的上游事件(而不是使用RxRingBuffer默认的128大小缓冲区)。
  • 加上当前正在处理的那个事件,看起来就像是队列大小为2。

在你的测试场景中:

  • 发送5后,concatMap开始处理Observable.timer(5, ...);
  • 发送6时,它会把6缓冲起来;
  • 发送7时,缓冲区已经满了(当前处理+缓冲共2个),此时上游BehaviorSubject是Hot Observable,不会等待下游请求,继续发送8就直接触发了MissingBackpressureException。

这种设计是为了在顺序处理的场景下,避免因为上游快速发送事件导致内存暴涨——毕竟concatMap无法并行处理,缓冲太多事件反而会增加内存压力和延迟。

2. 是否必须添加onBackpressure*算子?

答案是取决于你的上游Observable类型和业务场景:

  • 如果你的上游是Hot Observable(比如BehaviorSubject、PublishSubject这类),它们不会响应下游的背压请求,会持续发送事件。这种情况下,你必须在上游(concatMap之前)添加背压处理算子:

    • 用onBackpressureBuffer():可以自定义缓冲区大小,把暂时处理不过来的事件全部缓冲起来,适合不能丢失事件的场景;
    • 用onBackpressureDrop():丢弃处理不过来的事件,适合允许丢失部分非关键事件的场景;
    • 用onBackpressureLatest():只保留最新的未处理事件,适合只关心最新状态的场景。

    调整你的测试代码示例,添加onBackpressureBuffer()后就能避免异常:

    s.onBackpressureBuffer()
      .concatMap { Observable.timer(it, TimeUnit.MILLISECONDS) }
      .take(4)
      .subscribe(sub)
    
  • 如果你的上游是Cold Observable(比如Observable.just()、Observable.fromIterable()这类),它们会响应下游的背压请求,只有当下游请求时才会发送下一个事件,这种情况下不需要额外添加onBackpressure*算子,concatMap可以正常处理。

另外,如果你使用的是RxJava 2.x及以上版本,部分操作符的背压逻辑有调整,但核心思路是一致的:Hot Observable需要手动处理背压,Cold Observable则天然支持背压。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:02:24