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

如何用RxJava实现双超时触发的元素Buffer功能?

RxJava双条件缓冲实现方案

需求说明

  • 当最后一个元素产生后间隔5秒无新元素时,发布当前缓冲的元素
  • 当当前缓冲窗口的第一个元素产生后间隔20秒时,发布当前缓冲的元素

元素时序示例:[A -> 2秒 -> B -> 3秒 -> C -> 6秒 -> D -> 4秒 -> E -> 9秒 -> F -> 1秒 -> G -> 15秒 -> H]
预期输出结果:[A,B,C]、[D,E,F]、[G,H]

现有代码(仅实现20秒超时逻辑)

fun <T> Observable<T>.buffered(): Observable<List<T>> = publish { shared ->

    val startEvent = shared.throttleFirst(20, TimeUnit.SECONDS, scheduler)

    shared.buffer(startEvent.mergeWith(startEvent.delay(20, TimeUnit.SECONDS, scheduler, false)))

}

修改后的完整代码

fun <T> Observable<T>.buffered(scheduler: Scheduler): Observable<List<T>> = publish { shared ->
    // 捕获每个缓冲窗口的第一个元素,作为窗口启动标记
    val windowStart = shared.throttleFirst(Long.MAX_VALUE, TimeUnit.NANOSECONDS, scheduler)
    
    // 触发条件1:窗口启动后20秒自动输出缓冲
    val timeout20sSignal = windowStart.flatMap {
        Observable.timer(20, TimeUnit.SECONDS, scheduler)
    }
    
    // 触发条件2:最后一个元素产生后5秒无新元素,输出缓冲
    val idle5sSignal = shared.debounce(5, TimeUnit.SECONDS, scheduler)
    
    // 合并两个触发信号,作为buffer的边界事件
    val bufferTriggers = Observable.merge(timeout20sSignal, idle5sSignal)
    
    // 基于触发信号控制buffer,每次触发后自动开启新窗口
    shared.buffer(bufferTriggers)
}

核心逻辑说明

  • 窗口启动标记:用throttleFirst配合Long.MAX_VALUE,确保只会捕获每个缓冲周期的第一个元素,以此标记窗口的启动时间点。
  • 20秒超时触发:每个窗口启动后,通过timer延迟20秒发送信号,强制输出当前缓冲内容并开启新窗口。
  • 5秒空闲触发:通过debounce监听元素流,当连续5秒没有新元素产生时,发送信号输出缓冲内容并开启新窗口。
  • 双信号合并:将两个触发信号合并后作为buffer的边界,任意一个信号触发都会执行缓冲输出,完美满足需求中的两个条件。

内容的提问来源于stack exchange,提问作者Serg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 08:12:35