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

为何Cold Observable类型的可观察对象也需使用backpressure?

为什么Cold Observable也需要背压策略?

很多人误以为只有Hot Observable才需要处理背压,但实际上只要生产者发射事件的速度远快于消费者处理速度,不管是Hot还是Cold Observable,都会触发背压问题,甚至导致程序崩溃。下面结合你给出的两个例子具体分析:

示例1:Flowable.interval导致的崩溃

// 若不使用backpressure,程序将会崩溃
Flowable.interval(1, TimeUnit.MILLISECONDS, Schedulers.io())
        .observeOn(AndroidSchedulers.mainThread())
        .subscribe(new Consumer<Long>() {
            @Override
            public void accept(Long aLong) throws Exception {
                // 执行操作
            }
        });

Flowable.interval确实是Cold Observable,但它在Schedulers.io()线程以每毫秒1次的高频发射事件,而消费者在主线程处理任务。主线程本身要承担UI渲染、用户交互等大量工作,处理事件的速度远远跟不上每秒1000次的发射频率。

observeOn操作符内部有一个默认大小的缓冲区(Android平台通常为16),当生产者快速填满这个缓冲区后,后续的事件没有存储空间,就会直接抛出MissingBackpressureException,导致程序崩溃。

示例2:Flowable.range的背压需求

Flowable.range(0, 1000000)
    .onBackpressureBuffer()
    .observeOn(Schedulers.computation())
    .subscribe(new FlowableSubscriber<Integer>() {
        @Override
        public void onSubscribe(Subscription s) {
        }
        @Override
        public void onNext(Integer integer) {
            Log.d(TAG, "onNext: " + integer);
        }
        @Override
        public void onError(Throwable t) {
            Log.e(TAG, "onError: ", t);
        }
        @Override
        public void onComplete() {
        }
    });

Flowable.range是典型的Cold Observable,它会同步、快速地发射100万条数据。虽然消费者运行在computation线程,但打印Log的操作速度完全赶不上range的发射速度。如果去掉onBackpressureBuffer,observeOn的默认缓冲区会迅速被填满,同样会触发背压异常。

这里的onBackpressureBuffer相当于手动扩大了缓冲区,让生产者可以先把事件暂存起来,再逐步交给消费者处理,避免了缓冲区溢出的问题。


核心结论

背压的本质是解决生产者与消费者的速度不匹配问题,和Observable是Hot还是Cold没有直接关系。Cold Observable只是支持重新订阅、从头发射事件,但在单次订阅过程中,只要生产者速度远超消费者,就必须使用背压策略(比如缓冲区、节流、丢弃等)来平衡速度差,防止程序崩溃。

内容的提问来源于stack exchange,提问作者a.rah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 10:56:06