如何用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
相关产品推荐
相关产品推荐

