温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

RxJava背压如何解决

发布时间:2026-09-01 14:36:28 来源:亿速云 阅读:89 作者:小樊 栏目:编程语言

RxJava 中,“背压(Backpressure)”指的是:上游发送数据的速度,快于下游处理数据的速度,从而导致内存堆积甚至 OOM。

RxJava 从 RxJava 2 开始,明确区分了:

  • Observable:不支持背压
  • Flowable:支持背压(基于 Reactive Streams 规范)

一、最根本的解决方式:使用 Flowable

Flowable.create((FlowableOnSubscribe<Integer>) emitter -> {
    for (int i = 0; i < 1000; i++) {
        emitter.onNext(i);
    }
}, BackpressureStrategy.BUFFER)
.subscribe(
    value -> System.out.println(value),
    Throwable::printStackTrace
);

二、BackpressureStrategy 的几种策略

Flowable.create() 时必须指定背压策略:

1. BUFFER(缓冲)

  • 把所有来不及处理的数据缓存起来
  • 可能 OOM
BackpressureStrategy.BUFFER

✅ 适合:数据不能丢
❌ 不适合:高频数据流


2. DROP(丢弃)

  • 下游处理不过来时,直接丢弃新数据
BackpressureStrategy.DROP

✅ 适合:实时性要求高,丢数据可接受(如传感器)


3. LATEST(只保留最新)

  • 只保留最新的一条,丢弃旧数据
BackpressureStrategy.LATEST

✅ 适合:只关心最新状态(如 UI 状态)


4. ERROR(直接报错)

  • 下游跟不上就抛 MissingBackpressureException
BackpressureStrategy.ERROR

✅ 适合:严格要求处理能力


5. MISSING

  • 不做任何背压处理(类似 Observable)
BackpressureStrategy.MISSING

三、常用操作符解决背压

1. onBackpressureBuffer

Flowable.range(1, 1000000)
    .onBackpressureBuffer(1000)
    .subscribe(...);

2. onBackpressureDrop

Flowable.interval(1, TimeUnit.MILLISECONDS)
    .onBackpressureDrop()
    .subscribe(...);

3. onBackpressureLatest

Flowable.interval(1, TimeUnit.MILLISECONDS)
    .onBackpressureLatest()
    .subscribe(...);

四、控制下游消费速度(推荐)

1. observeOn + 调度器

Flowable.range(1, 1000)
    .subscribeOn(Schedulers.io())
    .observeOn(Schedulers.computation(), false, 16) // 限制请求量
    .subscribe(...);

2. request() 手动拉取(高级)

Flowable.range(1, 100)
    .subscribe(new Subscriber<Integer>() {
        private Subscription s;

        @Override
        public void onSubscribe(Subscription s) {
            this.s = s;
            s.request(1);
        }

        @Override
        public void onNext(Integer integer) {
            System.out.println(integer);
            s.request(1);
        }

        ...
    });

五、Observable 有没有背压?

没有

  • 数据直接推给下游
  • 高频流(如 intervalflatMap)容易 OOM

✅ 解决方案:

  • 换成 Flowable
  • 或者在源头控制发射频率(如 samplethrottleFirst

六、实际开发建议

场景 建议
高频事件(点击、传感器) Flowable + DROP / LATEST
网络 + 数据库 Flowable + BUFFER
UI 更新 observeOn(main) + throttleFirst
大量数据 request 拉取

如果你愿意,可以告诉我:

  • RxJava 版本(1 / 2 / 3)
  • 具体场景(网络、UI、数据库、采集)

我可以给你 针对性的背压方案代码

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI