Rx.NET中Buffer操作符如何在取消时立即发出所有缓存项?
解决方案:Rx.NET Buffer取消时立即吐出缓存项
完全可行,你可以通过给Buffer方法添加一个由取消令牌触发的额外刷新信号来实现需求。
核心思路
Rx.NET的Buffer方法有一个重载:Buffer(TimeSpan timeSpan, int count, IObservable<TBufferClosing> bufferClosing),它会在三个条件任意满足时输出当前缓冲区:
- 达到指定时间间隔(5秒)
- 缓冲区项数达到上限(10项)
- 传入的
bufferClosing序列发出值
我们可以基于取消令牌创建这个bufferClosing序列,当取消令牌激活时,立即发出信号触发缓冲区刷新。
代码实现
var cts = new CancellationTokenSource(); // 创建取消触发的刷新信号 var flushOnCancel = Observable.Create<Unit>(observer => { // 注册取消回调,触发缓冲区刷新 return cts.Token.Register(() => { observer.OnNext(Unit.Default); observer.OnCompleted(); }); }); // 带取消刷新的Buffer逻辑 var bufferedObservable = observable.Buffer( TimeSpan.FromSeconds(5), 10, flushOnCancel); // 订阅处理 using var subscription = bufferedObservable.Subscribe(buffer => { Console.WriteLine($"收到缓冲区: {string.Join(", ", buffer)}"); }); // 模拟2秒后触发取消 Task.Delay(2000).ContinueWith(_ => cts.Cancel());
效果验证
按照你的示例:2秒内产生1、2、3三个项,此时触发取消,flushOnCancel会立即发出信号,Buffer会马上输出[1,2,3],无需等待剩余3秒。
补充说明
- 取消触发刷新后,
flushOnCancel序列会完成,若后续原序列还有项,Buffer会继续按时间/数量条件输出新的缓冲区(如果不需要后续输出,可以结合TakeUntil(cts.Token.AsObservable())终止整个序列)。 - 该方式完全兼容原
Buffer的时间、数量限制逻辑,只是额外增加了取消时的即时刷新触发。
内容的提问来源于stack exchange,提问作者Michal
相关产品推荐
相关产品推荐

