在 RxJava 中,“背压(Backpressure)”指的是:上游发送数据的速度,快于下游处理数据的速度,从而导致内存堆积甚至 OOM。
RxJava 从 RxJava 2 开始,明确区分了:
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
);
Flowable.create() 时必须指定背压策略:
BackpressureStrategy.BUFFER
✅ 适合:数据不能丢
❌ 不适合:高频数据流
BackpressureStrategy.DROP
✅ 适合:实时性要求高,丢数据可接受(如传感器)
BackpressureStrategy.LATEST
✅ 适合:只关心最新状态(如 UI 状态)
MissingBackpressureExceptionBackpressureStrategy.ERROR
✅ 适合:严格要求处理能力
BackpressureStrategy.MISSING
Flowable.range(1, 1000000)
.onBackpressureBuffer(1000)
.subscribe(...);
Flowable.interval(1, TimeUnit.MILLISECONDS)
.onBackpressureDrop()
.subscribe(...);
Flowable.interval(1, TimeUnit.MILLISECONDS)
.onBackpressureLatest()
.subscribe(...);
Flowable.range(1, 1000)
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.computation(), false, 16) // 限制请求量
.subscribe(...);
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);
}
...
});
❌ 没有
interval、flatMap)容易 OOM✅ 解决方案:
Flowablesample、throttleFirst)| 场景 | 建议 |
|---|---|
| 高频事件(点击、传感器) | Flowable + DROP / LATEST |
| 网络 + 数据库 | Flowable + BUFFER |
| UI 更新 | observeOn(main) + throttleFirst |
| 大量数据 | request 拉取 |
如果你愿意,可以告诉我:
我可以给你 针对性的背压方案代码。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。