RxJava如何生成持续监听另一Observable最近4条数据的Observable
RxJava实现持续滑动窗口获取最近4条数据方案
方案1:使用buffer(int count, int skip)实现滑动窗口
buffer(4,1)可以实现每次滑动1个元素的4元素窗口,正好匹配你需要的(1,2,3,4)→(2,3,4,5)→(3,4,5,6)的效果。结合过滤和映射逻辑即可满足需求:
Observable<A> generateObservable() { return rep.asObservable() // 生成大小为4、步长为1的滑动窗口 .buffer(4, 1) // 只保留包含完整4条数据的窗口(过滤初始阶段不足4条的情况) .filter(window -> window.size() == 4) .map(window -> { // 检查窗口中每条数据的b属性,返回符合条件的A for (A item : window) { if (checkB(item.b)) { // 替换为你的b属性检查逻辑 return item; } } // 若无符合条件的数据,默认返回窗口最后一条(可根据需求调整) return window.get(window.size() - 1); }); }
方案2:利用scan维护滑动队列
通过scan操作符维护一个最多存储4条数据的队列,每次新数据加入时自动移除最旧的元素,实现滑动窗口效果:
Observable<A> generateObservable() { return rep.asObservable() // 用scan维护一个最多4条数据的队列 .scan(new ArrayDeque<>(4), (queue, newItem) -> { queue.add(newItem); if (queue.size() > 4) { queue.poll(); // 移除最旧的元素 } return queue; }) // 只处理队列满4条的情况 .filter(queue -> queue.size() == 4) .map(queue -> { // 检查b属性并返回目标A for (A item : queue) { if (checkB(item.b)) { return item; } } return queue.getLast(); }); }
方案3:借助ReplayRelay的缓存特性
由于你已经使用了大小为4的ReplayRelay,它本身会缓存最近4条数据。可以在每次新数据到来时直接获取其缓存列表:
Observable<A> generateObservable() { return rep.asObservable() // 等待Relay累积到4条数据后再开始处理 .skipUntil(rep.take(4)) .map(ignored -> { // 获取Relay当前缓存的最近4条数据 List<A> recentFour = Arrays.asList(rep.getValues()); // 检查b属性并返回目标A for (A item : recentFour) { if (checkB(item.b)) { return item; } } return recentFour.get(recentFour.size() - 1); }); }
说明
- 你之前尝试的
buffer(4)是跳跃窗口(每次取4条后跳过4条),而buffer(4,1)才是滑动窗口,这是关键区别。 take/takeLast会终止流,不符合持续监听的需求;scan默认只处理前后两个元素,但可以通过自定义累积器维护队列来实现多元素窗口。
内容的提问来源于stack exchange,提问作者problem_asker_2022
相关产品推荐
相关产品推荐

