多层软件架构下基于System.Reactive的IObservable<T>转发方案探讨
多层架构中转发IObservable序列的通用模式
下面是几个能解决这类转发场景的实用模式,避免重复写Subject相关的冗余代码:
1. ConnectableObservable + AutoConnect/RefCount:延迟激活+共享订阅
当你需要在调用ConnectAsync这类方法后才启动底层序列,同时让多个订阅者共享同一个序列时,用这个模式最适合。它能自动管理订阅的生命周期,不用手动写Subscribe到Subject的逻辑。
示例代码:
public class Device { private readonly IObservable<Packet> _incomingPackets; private IDisposable? _connectionSubscription; private readonly IIpDevice _ipDevice; public IObservable<Packet> IncomingPackets => _incomingPackets; public Device(IIpDevice ipDevice) { _ipDevice = ipDevice; // 将底层序列转为可连接的,只有调用Connect才会开始发射元素 _incomingPackets = ipDevice.IncomingPackets .ObserveOn(Scheduler.Default) .Publish() // 转为ConnectableObservable .RefCount(); // 自动管理连接:有订阅时启动,无订阅时停止 } public async Task<bool> ConnectAsync() { var connected = await _ipDevice.ConnectAsync(); return connected; } public void Disconnect() { _connectionSubscription?.Dispose(); } }
适用场景:底层序列需要延迟启动、多个订阅者共享同一数据流,且希望自动管理订阅启停的情况。
2. Observable.Defer + 异步初始化:处理动态创建的底层序列
如果底层的IncomingPackets不是在构造时就存在,而是在ConnectAsync之后才可用,用Defer延迟序列的创建,结合异步逻辑确保订阅时能拿到正确的序列。
示例代码:
public class Device { private IIpDevice? _ipDevice; private readonly TaskCompletionSource<bool> _connectedTcs = new(); public IObservable<Packet> IncomingPackets => Observable.Defer(() => { // 等待连接完成后,返回底层的序列 return _connectedTcs.Task .Where(isConnected => isConnected) .SelectMany(_ => _ipDevice!.IncomingPackets.ObserveOn(Scheduler.Default)); }); public async Task<bool> ConnectAsync() { // 模拟创建或初始化_ipDevice的逻辑 _ipDevice = await IpDeviceFactory.CreateAsync(); var connected = await _ipDevice.ConnectAsync(); _connectedTcs.SetResult(connected); return connected; } }
适用场景:底层序列依赖异步初始化,订阅者可能在连接完成前就订阅的情况,Defer会在每个订阅时检查状态,确保返回有效的序列。
3. 复用转发工具类:减少重复代码
如果多个类都需要做类似的转发逻辑,可以封装一个通用的转发器,把Subject的管理逻辑抽离出来,避免重复写相同的代码。
示例代码:
// 通用转发器类 public class ObservableForwarder<T> : IDisposable { private readonly Subject<T> _subject = new(); public IObservable<T> Observable => _subject.AsObservable(); public void Forward(T item) => _subject.OnNext(item); public void Complete() => _subject.OnCompleted(); public void Error(Exception ex) => _subject.OnError(ex); public void Dispose() { _subject.Dispose(); } } // 使用转发器的Device类 public class Device : IDisposable { private readonly ObservableForwarder<Packet> _packetForwarder = new(); private readonly IIpDevice _ipDevice; private IDisposable? _subscription; public IObservable<Packet> IncomingPackets => _packetForwarder.Observable; public Device(IIpDevice ipDevice) { _ipDevice = ipDevice; } public async Task<bool> ConnectAsync() { var connected = await _ipDevice.ConnectAsync(); if (connected) { _subscription = _ipDevice.IncomingPackets .ObserveOn(Scheduler.Default) .Subscribe( packet => _packetForwarder.Forward(packet), ex => _packetForwarder.Error(ex), () => _packetForwarder.Complete() ); } return connected; } public void Dispose() { _subscription?.Dispose(); _packetForwarder.Dispose(); } }
适用场景:多个类需要相同的转发逻辑,通过工具类统一封装Subject的操作,减少代码重复,同时也能统一处理Complete和Error的传播(你之前的代码没处理这两个,实际场景中很重要)。
4. ReplaySubject:缓存历史元素给晚订阅者
如果需要让后续订阅的消费者能拿到连接后已经产生过的元素(比如Logger在连接后才订阅,想拿到之前的数据包),可以用ReplaySubject代替普通Subject,它会缓存指定数量或时间范围内的元素。
示例代码:
public class Device : IDisposable { private readonly ReplaySubject<Packet> _packetPublisher = new(10); // 缓存最近10个数据包 private readonly IIpDevice _ipDevice; private IDisposable? _subscription; public IObservable<Packet> IncomingPackets => _packetPublisher.AsObservable(); public Device(IIpDevice ipDevice) { _ipDevice = ipDevice; } public async Task<bool> ConnectAsync() { var connected = await _ipDevice.ConnectAsync(); if (connected) { _subscription = _ipDevice.IncomingPackets .ObserveOn(Scheduler.Default) .Subscribe( packet => _packetPublisher.OnNext(packet), _packetPublisher.OnError, _packetPublisher.OnCompleted ); } return connected; } public void Dispose() { _subscription?.Dispose(); _packetPublisher.Dispose(); } }
适用场景:需要新订阅者能获取历史事件的场景,比如日志、监控类组件晚于业务组件启动的情况。
内容的提问来源于stack exchange,提问作者this.myself
相关产品推荐
相关产品推荐

