Rx.Net:缓存设备连接数据直至IObserver订阅的实现方案
Rx.Net 热Observable缓存设备连接信息的解决方案
针对你遇到的热Observable订阅滞后导致错过设备连接通知的问题,结合不想用ReplaySubject无限缓存的需求,分两种常见场景给出解决方案:
场景1:维护当前已连接的设备集合(推荐)
如果你的视图需要展示当前处于连接状态的设备,而非所有历史连接记录,用BehaviorSubject维护设备集合是最优选择:
BehaviorSubject会保存最新的设备集合,新订阅者一订阅就能立即拿到当前所有已连接设备,后续还能实时接收设备的增删更新。- 搭配原始的设备连接/断开事件,动态更新集合:
// 假设原始的热Observable:设备连接、断开事件 IObservable<Device> deviceConnected = ...; IObservable<Device> deviceDisconnected = ...; // 初始化BehaviorSubject,初始为空设备集合 var currentConnectedDevices = new BehaviorSubject<HashSet<Device>>(new HashSet<Device>()); // 处理设备连接:添加到集合并推送更新 deviceConnected.Subscribe(device => { var updatedSet = new HashSet<Device>(currentConnectedDevices.Value); updatedSet.Add(device); currentConnectedDevices.OnNext(updatedSet); }); // 处理设备断开:从集合移除并推送更新(如果有断开事件的话) deviceDisconnected.Subscribe(device => { var updatedSet = new HashSet<Device>(currentConnectedDevices.Value); updatedSet.Remove(device); currentConnectedDevices.OnNext(updatedSet); }); // 视图订阅:立即获取当前设备列表,后续实时更新 currentConnectedDevices.Subscribe(devices => { // 刷新视图显示逻辑 });
场景2:保留历史连接事件但限制缓存
如果确实需要记录所有设备连接的历史事件,但要避免缓存无限增长,可以用以下方法:
方法A:带限制的ReplaySubject
创建ReplaySubject时指定缓存的最大数量或时间窗口,自动清理超出限制的旧事件:
// 方案1:最多缓存最近100条连接事件 var limitedReplay = ReplaySubject<Device>.Create(100); // 方案2:只缓存最近5分钟内的连接事件 // var limitedReplay = ReplaySubject<Device>.Create(TimeSpan.FromMinutes(5)); // 订阅原始热Observable deviceConnected.Subscribe(limitedReplay); // 视图订阅:先收到缓存的历史事件,再接收新事件 limitedReplay.Subscribe(device => { // 更新视图展示历史/新连接设备 });
方法B:仅给第一个订阅者补发历史事件
如果只有一个视图会订阅,且只需要给第一个订阅者补发之前的连接事件,后续不再缓存新事件,可以自定义逻辑实现:
var connectable = deviceConnected.Publish().Replay(); connectable.Connect(); // 第一个订阅者会收到所有缓存事件,之后直接转发新事件 var observable = connectable.TakeUntil(connectable.Subscribe(_ => { })) .Concat(deviceConnected); // 视图订阅 observable.Subscribe(device => { // 处理连接事件 });
方法C:自定义缓存Subject
如果上述方法都不满足需求,可以自己实现一个Subject:无订阅者时缓存事件,有订阅者时发送缓存并清空,之后直接转发新事件:
public class CacheUntilFirstSubscribeSubject<T> : ISubject<T> { private readonly List<T> _eventCache = new List<T>(); private readonly Subject<T> _innerSubject = new Subject<T>(); private bool _hasSubscribers = false; private readonly object _lockObj = new object(); public void OnCompleted() => _innerSubject.OnCompleted(); public void OnError(Exception error) => _innerSubject.OnError(error); public void OnNext(T value) { lock (_lockObj) { if (!_hasSubscribers) { _eventCache.Add(value); } else { _innerSubject.OnNext(value); } } } public IDisposable Subscribe(IObserver<T> observer) { lock (_lockObj) { if (!_hasSubscribers) { _hasSubscribers = true; // 补发缓存的事件 foreach (var cachedEvent in _eventCache) { observer.OnNext(cachedEvent); } _eventCache.Clear(); } return _innerSubject.Subscribe(observer); } } }
使用示例:
var cacheSubject = new CacheUntilFirstSubscribeSubject<Device>(); deviceConnected.Subscribe(cacheSubject); // 视图订阅时,先收到所有缓存的连接事件,之后接收新事件 cacheSubject.Subscribe(device => { // 更新视图 });
内容的提问来源于stack exchange,提问作者resp78
相关产品推荐
相关产品推荐

