在 RxJava 中,背压(Backpressure) 指的是:
上游发射数据的速度 > 下游处理数据的速度 时,如何防止 OOM 或系统崩溃。
RxJava 主要通过 Flowable + 背压策略 来处理,而 Observable 默认不支持背压(会直接抛 MissingBackpressureException)。
下面按「核心概念 → 处理方式 → 示例」来说明。
| 类型 | 是否支持背压 | 使用场景 |
|---|---|---|
| Observable | ❌ 不支持 | 事件少、UI 事件 |
| Flowable | ✅ 支持 | 高频数据流(网络、数据库、传感器) |
Observable.create(...) // 无背压
Flowable.create(...) // 有背压
在创建 Flowable 时指定:
Flowable.create(emitter -> {
for (int i = 0; i < 1000; i++) {
emitter.onNext(i);
}
emitter.onComplete();
}, BackpressureStrategy.BUFFER)
| 策略 | 说明 |
|---|---|
| BUFFER | 缓存所有数据(可能 OOM) |
| DROP | 下游处理不过来就丢弃 |
| LATEST | 只保留最新的一条 |
| ERROR | 上游过快直接抛异常 |
| MISSING | 不处理背压(需手动 request) |
Flowable 是 响应式拉取模型:
Flowable.range(1, 100)
.onBackpressureDrop()
.subscribe(new Subscriber<Integer>() {
private Subscription subscription;
@Override
public void onSubscribe(Subscription s) {
subscription = s;
s.request(1); // 请求 1 个
}
@Override
public void onNext(Integer integer) {
System.out.println(integer);
subscription.request(1); // 处理完再要
}
@Override
public void onError(Throwable t) {}
@Override
public void onComplete() {}
});
✅ 这是最安全、推荐的方式
Flowable.interval(1, TimeUnit.MILLISECONDS)
.buffer(100)
.subscribe(list -> System.out.println(list.size()));
Flowable.interval(10, TimeUnit.MILLISECONDS)
.sample(1, TimeUnit.SECONDS)
.subscribe(System.out::println);
.onBackpressureDrop()
.onBackpressureLatest()
Observable.range(1, 100)
.toFlowable(BackpressureStrategy.BUFFER);
⚠️ 只是“强制兼容”,并不能真正解决背压
✅ 高频数据流 → Flowable
✅ 下游 request 控制消费速度
✅ 避免 BUFFER 无限增长
✅ UI 事件用 Observable,数据流用 Flowable
如果你愿意,我可以:
你更想看哪一种?
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。