RxJava的背压机制是用于解决异步场景下上游数据发送速度远快于下游处理速度,导致数据积压、内存溢出(OOM)等问题的策略,核心是通过控制上游发送速率实现流量控制。
触发条件
range、create)支持背压,Hot Observable(如interval、just)默认不支持。实现方式
request(n)主动告知上游可处理的数据量,上游按需发送(类似TCP流量控制)。
request(1),上游再发送下一个。Flowable的create方法指定背压策略(如BUFFER、DROP等),或使用onBackpressureXXX操作符转换不支持背压的Observable。常见策略
| 策略 | 行为 | 缓存大小 | 适用场景 |
|---|---|---|---|
| BUFFER | 缓存所有未处理数据,可能导致OOM | 无限制(默认128) | 数据可丢失场景 |
| DROP | 丢弃无法处理的新数据 | 0(直接丢弃) | 允许数据丢失的实时场景 |
| LATEST | 仅保留最新数据,丢弃旧数据 | 1(覆盖旧数据) | 需要最新数据的场景 |
| ERROR | 缓存满时抛出MissingBackpressureException异常 |
128 | 严格限制数据量的场景 |
使用场景
DROP或LATEST策略丢弃冗余数据,防止缓存溢出。// 使用Flowable指定背压策略(以BUFFER为例)
Flowable.create(emitter -> {
for (int i = 0; i < 1000; i++) {
emitter.onNext(i);
}
}, BackpressureStrategy.BUFFER)
.subscribe(new Subscriber<Integer>() {
@Override
public void onSubscribe(Subscription s) {
s.request(10); // 初始请求10条数据
}
@Override
public void onNext(Integer integer) {
System.out.println(integer);
if (integer % 10 == 0) {
s.request(10); // 每处理10条再请求10条
}
}
});
总结:背压机制通过异步协作+策略控制,平衡上下游处理速度,避免内存问题,是响应式编程中处理流速差异的核心手段。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。