You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 20:15:55