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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 09:55:23