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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 12:07:35