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

多层软件架构下基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 11:42:14