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

如何订阅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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 08:00:23