CombineLatest队列仅发最后项?Realm数据源背压策略疑问
问题:如何让RxJava的combineLatest仅发送最后一项?
我有两个Realm数据源,想用combineLatest操作符合并处理,伪代码如下:
Flowable.combineLatest( dataSource1.observeOn(Schedulers.io(), false, 1), dataSource2.observeOn(Schedulers.io(), false, 1), 1 ) .observeOn(Schedulers.computation(), false, 1) // 处理逻辑 .doOnNext(item -> Thread.sleep(1000)) .subscribe();
我已经通过限制observeOn的缓存项数来避免发送旧数据库缓存数据,现在想知道能不能让combineLatest的队列仅发送最后一项?
回答
当然可以!针对你的场景,这里有几个实用的方案,帮你实现只保留最新项的需求:
1. 用onBackpressureLatest()直接处理背压场景
既然你已经在使用observeOn设置缓存项数来处理背压,那onBackpressureLatest()就是最贴合你需求的选择。它会自动丢弃上游发送的旧项,只保留最新的那一个等待下游处理:
Flowable.combineLatest( dataSource1.observeOn(Schedulers.io(), false, 1), dataSource2.observeOn(Schedulers.io(), false, 1), 1 ) .onBackpressureLatest() // 丢弃积压的旧项,仅保留最新待处理项 .observeOn(Schedulers.computation(), false, 1) .doOnNext(item -> Thread.sleep(1000)) .subscribe();
这个操作符完美适配你的场景——当下游因为Thread.sleep(1000)处理速度跟不上上游更新时,它会自动只留最新的组合结果,不会让旧数据占用队列。
2. 用debounce过滤短时间内的连续更新
如果你的数据源会频繁触发连续更新,debounce可以帮你忽略掉这些中间更新,只保留最后一次更新后等待指定时长的项。你可以根据自己的处理速度调整时长:
Flowable.combineLatest( dataSource1.observeOn(Schedulers.io(), false, 1), dataSource2.observeOn(Schedulers.io(), false, 1), 1 ) .debounce(150, TimeUnit.MILLISECONDS) // 忽略150ms内的连续更新,仅保留最后一项 .observeOn(Schedulers.computation(), false, 1) .doOnNext(item -> Thread.sleep(1000)) .subscribe();
注意:这个方案会带来一点延迟,适合更新频率高但你不需要实时响应每一次更新的场景。
3. 用BehaviorSubject维护最新状态(精细控制)
如果你需要更精细的控制每个数据源的最新值,可以用BehaviorSubject来缓存每个数据源的最新状态,然后只发送基于最新状态组合的结果:
// 用BehaviorSubject缓存每个数据源的最新值 BehaviorSubject<Object> latestData1 = BehaviorSubject.create(); BehaviorSubject<Object> latestData2 = BehaviorSubject.create(); // 订阅数据源,更新Subject的状态 dataSource1.observeOn(Schedulers.io(), false, 1) .subscribe(latestData1); dataSource2.observeOn(Schedulers.io(), false, 1) .subscribe(latestData2); // 合并两个Subject的最新值,且仅当组合值变化时发送 Flowable.combineLatest(latestData1, latestData2, (v1, v2) -> combineYourData(v1, v2)) .distinctUntilChanged() .observeOn(Schedulers.computation(), false, 1) .doOnNext(item -> Thread.sleep(1000)) .subscribe();
这种方式可以确保你始终基于两个数据源的最新状态进行组合,还能通过distinctUntilChanged避免发送重复的组合结果。
小提示
- 优先推荐
onBackpressureLatest(),因为它直接针对你已经在处理的背压场景,代码改动最小,效果最直接。 - 如果你的更新不是连续爆发式的,
debounce可能会带来不必要的延迟,这时候onBackpressureLatest()更合适。
内容的提问来源于stack exchange,提问作者elnino
相关产品推荐
相关产品推荐

