如何配置C# Reactive ReplaySubject,使缓存项被访问后清除
实现一次性重播的自定义Subject(满足你的需求)
嘿,这个需求原生的ReplaySubject<T>确实做不到——它的缓存是持久化的,只要没触发清理条件(比如大小/时间限制),后续订阅者还是能拿到之前缓存的项。不过我们可以通过自定义一个ISubject<T>来实现你要的行为:无订阅时缓存项,有订阅者时一次性推送所有缓存并清空,后续订阅者拿不到已推送过的项。
自定义Subject实现
下面是一个线程安全的实现,核心是用队列缓存无订阅时的项,新订阅者进来时先消费所有缓存再订阅内部Subject:
using System; using System.Collections.Generic; using System.Reactive.Subjects; using System.Reactive; public class OneTimeReplaySubject<T> : ISubject<T> { private readonly Subject<T> _innerSubject = new Subject<T>(); private readonly Queue<T> _messageCache = new Queue<T>(); private readonly object _lock = new object(); public void OnNext(T value) { lock (_lock) { // 如果有活跃订阅者,直接发送;否则缓存 if (_innerSubject.HasObservers) { _innerSubject.OnNext(value); } else { _messageCache.Enqueue(value); } } } public void OnError(Exception error) { lock (_lock) { _messageCache.Clear(); _innerSubject.OnError(error); } } public void OnCompleted() { lock (_lock) { // 完成前先推送所有缓存项 foreach (var item in _messageCache) { _innerSubject.OnNext(item); } _messageCache.Clear(); _innerSubject.OnCompleted(); } } public IDisposable Subscribe(IObserver<T> observer) { lock (_lock) { // 新订阅者先接收所有缓存项 foreach (var cachedItem in _messageCache) { observer.OnNext(cachedItem); } // 清空缓存,避免后续订阅者拿到这些项 _messageCache.Clear(); // 订阅内部Subject,接收后续的实时消息 return _innerSubject.Subscribe(observer); } } }
使用示例
我们来验证一下这个Subject的行为是否符合你的需求:
var subject = new OneTimeReplaySubject<int>(); // 阶段1:无订阅者,发送的项被缓存 subject.OnNext(1); subject.OnNext(2); subject.OnNext(3); // 阶段2:第一个订阅者进来,收到所有缓存项 var subscriber1 = subject.Subscribe( x => Console.WriteLine($"订阅者1收到:{x}"), ex => Console.WriteLine($"订阅者1出错:{ex.Message}"), () => Console.WriteLine("订阅者1完成") ); // 输出: // 订阅者1收到:1 // 订阅者1收到:2 // 订阅者1收到:3 // 阶段3:有活跃订阅者,新发送的项直接推送 subject.OnNext(4); // 输出:订阅者1收到:4 // 阶段4:取消订阅,再次无订阅者,发送的项被缓存 subscriber1.Dispose(); subject.OnNext(5); subject.OnNext(6); // 阶段5:第二个订阅者进来,收到最新的缓存项(1-3已经被清空了) var subscriber2 = subject.Subscribe( x => Console.WriteLine($"订阅者2收到:{x}"), () => Console.WriteLine("订阅者2完成") ); // 输出: // 订阅者2收到:5 // 订阅者2收到:6 // 阶段6:完成Subject,推送剩余缓存(这里没有)并通知完成 subject.OnCompleted(); // 输出:订阅者2完成 subscriber2.Dispose();
关键行为说明
- 无订阅时缓存:当没有任何订阅者时,所有
OnNext的项都会被存入内部队列。 - 订阅时一次性推送并清空:新订阅者会立即收到所有缓存项,之后缓存被清空,后续订阅者无法获取这些已推送的项。
- 线程安全:用
lock保证了多线程环境下的缓存操作和订阅操作的原子性,符合Rx的线程安全契约。
内容的提问来源于stack exchange,提问作者user2074945
相关产品推荐
相关产品推荐

