如何用Monix实现随源元素触发的带重叠时间窗口缓冲?
实现随元素触发的滑动时间窗口收集
针对你需要的「源Observable每发射一个元素,就输出最近指定时长内所有元素」的需求,不需要自定义类似BufferTimedObservable的操作符,可以通过Monix现有操作符的组合实现,核心思路是维护一个带时间戳的元素队列,每次新元素进入时清理超时项并输出当前队列。
基础实现方案
import monix.reactive._ import monix.execution.Scheduler.Implicits.global import scala.concurrent.duration._ def slidingTimeWindow[T](windowDuration: FiniteDuration): Observable[T] => Observable[List[T]] = source => { source // 为每个元素添加发射时的时间戳 .map(t => (System.currentTimeMillis(), t)) // 用scan维护元素队列,每次更新时先清理超时元素 .scan(List.empty[(Long, T)]) { (acc, elem) => val currentTime = System.currentTimeMillis() val validElements = acc.filter { case (timestamp, _) => currentTime - timestamp <= windowDuration.toMillis } validElements :+ elem } // 提取队列中的元素值,去掉时间戳 .map(_.map(_._2)) } // 使用示例 Observable.range(0, 100) .delayExecution(1.second) .delayOnNext(1.second) // 模拟每秒发射一个元素 .compose(slidingTimeWindow(3.seconds)) .subscribe(list => println(list))
性能优化版本
如果窗口内元素数量较多,不可变列表的过滤操作会有性能开销,可以改用可变队列从头部移除超时元素:
import scala.collection.mutable def slidingTimeWindowOptimized[T](windowDuration: FiniteDuration): Observable[T] => Observable[List[T]] = source => { source .map(t => (System.currentTimeMillis(), t)) .scan(mutable.Queue.empty[(Long, T)]) { (queue, elem) => val currentTime = System.currentTimeMillis() // 从队列头部移除所有超时元素 while (queue.nonEmpty && currentTime - queue.head._1 > windowDuration.toMillis) { queue.dequeue() } queue.enqueue(elem) queue } .map(_.toList.map(_._2)) }
方案说明
- 核心逻辑:通过
scan维护状态队列,每次新元素进入时先过滤/移除超出时间窗口的旧元素,再将新元素加入队列,最后输出当前队列的元素列表。 - 为何不用
bufferTimed:bufferTimed是定时触发输出(每隔指定时长输出一次窗口内元素),而非源Observable每发射一个元素就触发输出,完全不符合需求。
内容的提问来源于stack exchange,提问作者Giuliano De Sabata
相关产品推荐
相关产品推荐

