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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 08:45:03