如何订阅Observable最新n次发射并对最近n次发射应用运算函数?
解决Observable最新n次发射及滑动窗口运算问题
1. 实现订阅Observable的最新n次发射
要获取Observable发射的最新n次数据(形成滑动窗口,每次新元素加入时自动保留最近的n个),RxJava里的buffer()操作符是最直接的解决方案。它能将Observable的元素按指定的窗口大小和步长分组,返回一个发射这些分组列表的Observable。
如果需要每次新元素发射时,都拿到包含当前元素在内的最近n个元素,只需要设置:
windowSize = n:每个窗口包含n个元素skip = 1:每新增1个元素就生成一个新窗口(实现滑动效果)
比如n=2时,原Observable发射1、2、3、4,buffer(2,1)会依次输出[1,2]、[2,3]、[3,4]这三个列表,每个列表都是当前的最新2次发射结果。
2. 对最近n次发射进行运算(以n=2求和为例)
结合你给出的场景:原Observable依次发射1、2、3、4,需要计算最近2次发射的和并输出3、5、7。我们可以在buffer()之后,用map()操作符对每个窗口的元素做求和运算,再过滤掉元素数量不足n的窗口(避免初始阶段窗口未满的无效计算)。
完整代码示例
伪代码逻辑
// 定义原Observable:依次发射1、2、3、4 Observable<Integer> source = Observable.just(1, 2, 3, 4); // 应用滑动窗口获取最近2次发射,再计算求和 source.buffer(2, 1) // 过滤掉元素数量不足2的窗口(比如初始阶段仅发射1时的窗口) .filter(list -> list.size() == 2) // 计算窗口内两个元素的和 .map(list -> list.get(0) + list.get(1)) // 订阅输出结果 .subscribe(sum -> System.out.println(sum));
运行这段代码会依次输出3、5、7,完全符合需求。
扩展说明
- 如果n更大(比如n=3),只需要把
buffer(2,1)改成buffer(3,1),同时调整filter条件为list.size() ==3,求和逻辑改成累加窗口内所有元素即可。 - 对于持续发射元素的无限流Observable,这个方案同样适用——
buffer()是实时滑动的,只会保留当前窗口内的元素,不会占用过多内存。
内容的提问来源于stack exchange,提问作者fweigl
相关产品推荐
相关产品推荐

