Rx使用Buffer批量处理操作时如何在新请求到达时重置时间窗
你需要实现的是带事件收集的防抖逻辑,原生Buffer操作符使用固定时间窗,不支持收到新事件时重置计时,可通过GroupByUntil组合操作符实现需求,示例实现代码如下:
_Event // 将所有事件归入同一个虚拟分组,方便统一控制生命周期 .GroupBy(_ => 0) .SelectMany(group => group // 分组终止条件:连续2000ms没有新事件流入 .TakeUntil(group.Throttle(TimeSpan.FromMilliseconds(2000))) // 收集分组内的所有事件 .ToList() ) .Subscribe(this._Observer);
上述代码的执行逻辑完全匹配你的需求:
- 第一个事件到达时创建分组,开始计时
- 每收到新事件,
Throttle的计时会自动重置,相当于时间窗被刷新 - 连续2000ms无新事件时,
Throttle触发信号终止当前分组,打包输出期间收集的所有事件,触发OnNext回调 - 不会输出空事件列表,原有代码中的
Where(obj => obj.Count > 0)可以直接省略
如果要复用该逻辑,可封装为通用的Rx扩展方法:
public static IObservable<IList<TSource>> BufferWithReset<TSource>(this IObservable<TSource> source, TimeSpan idleThreshold) { return source .GroupBy(_ => 0) .SelectMany(g => g.TakeUntil(g.Throttle(idleThreshold)).ToList()); }
封装后你的业务代码仅需修改一行即可使用:
_Event.BufferWithReset(TimeSpan.FromMilliseconds(2000)) .Subscribe(this._Observer);
内容的提问来源于stack exchange,提问作者Yogesh
相关产品推荐
相关产品推荐

