【问题标题】:Using Subject to decouple Observable subscription and initialisation使用 Subject 解耦 Observable 订阅和初始化
【发布时间】:2013-05-07 07:17:13
【问题描述】:

我有一个公开IObservable 状态的 API。但是这个状态取决于一个底层的 observable 源,它必须通过Init 初始化。

我想做的是保护用户不必按正确的顺序做事:按照目前的情况,如果他们在执行Init 之前尝试订阅Status,他们会得到一个例外因为它们的源没有初始化。

所以我有了使用Subject 将两者解耦的天才想法:订阅我的Status 的外部用户只是订阅了主题,然后当他们调用Init 时,我订阅了底层服务使用我的主题。

代码中的想法

private ISubject<bool> _StatusSubject = new Subject<bool>();
public IObservable<bool> Status { get { return _StatusSubject; } }

public void Init() 
{
    _Connection = new Connection();
    Underlying.GetDeferredObservable(_Connection).Subscribe(_StatusSubject);
}

但是,从对虚拟项目的测试来看,问题在于初始化 通过订阅 Subject 来“唤醒”我的底层 Observable,即使还没有人订阅该主题。如果可能的话,我想避免这种情况,但我不确定如何...

(我也注意到 the received wisdom 说“一般规则是,如果你使用一个主题,那么你做错了什么”)

【问题讨论】:

  • 为了它的价值,“有状态的”可观察流和Subject 在我的脑海中齐头并进;是的,它通常不受欢迎,但通常你会尝试建立一个无状态的流。如果状态的概念不能从流机制中抽象出来(比如需要特定的引导/初始化),我建议使用Subject;它方式更容易理解和使用。
  • 在我看来,避免Subject 的主要原因是人们通常会先达到它并最终得到一个复杂的解决方案(如下面的@Benjol 的解决方案),而没有意识到他们通常是- 实现RX 已经作为运营商提供的东西。 Subject 应该放在要尝试的事情列表的末尾附近。但是一旦你发现你不能以其他方式更简单地做到这一点,那么一定要使用Subject
  • @brandon 这是一个公平的观点;我通常会用主题制作原型,然后向后工作,将它们分解成直接的流。好的答案,顺便说一句

标签: c# system.reactive subject


【解决方案1】:

您似乎缺少的概念是如何知道某人何时开始收听并仅初始化您的基础来源。通常您使用Observable.Create 或其同级之一(DeferUsing、...)来执行此操作。

以下是没有Subject 的方法:

private IObservable<bool> _status = Observable.Defer(() =>
{
    _Connection = new Connection();
    return Underlying.GetDeferredObservable(_Connection);
};

public IObservable<bool> Status { get { return _status; } }

Defer 在有人真正订阅之前不会调用初始化代码。

但这有几个潜在的问题:

  1. 每个观察者都会建立一个新的连接
  2. 当观察者取消订阅时,连接不会被清理。

第二个问题很容易解决,所以让我们先解决这个问题。假设您的 Connection 是一次性的,在这种情况下您可以这样做:

private IObservable<bool> _status = Observable
    .Using(() => new Connection(),
           connection => Underlying.GetDeferredObservable(connection));

public IObservable<bool> Status { get { return _status; } }

通过这个迭代,每当有人订阅时,都会创建一个新的Connection 并传递给第二个 Lamba 方法来构造 observable。每当观察者取消订阅时,Connection 就是Disposed。如果Connection 不是IDisposable,那么您可以使用Disposable.Create(Action) 创建一个IDisposable,它将运行您需要运行的任何操作来清理连接。

你仍然有每个观察者创建一个新连接的问题。我们可以使用PublishRefCount 来解决这个问题:

private IObservable<bool> _status = Observable
    .Using(() => new Connection(),
           connection => Underlying.GetDeferredObservable(connection))
    .Publish()
    .RefCount();

public IObservable<bool> Status { get { return _status; } }

现在,当 first 观察者订阅时,将创建连接并订阅底层的 observable。后续观察者将共享连接并获取当前状态。当 last 观察者取消订阅时,连接将被释放并关闭一切。如果之后有其他观察者订阅,一切都会重新开始。

在底层,Publish 实际上是使用Subject 来共享单个可观察源。而RefCount 正在跟踪目前有多少观察者正在观察。

【讨论】:

  • 感谢您花时间回答。我需要做的另一件事是 not 订阅底层,即使 其他人订阅了 Status:这应该只发生在“Init”...
  • 也许我误解了你,但这就是最后一个例子的工作方式。 Publish 返回一个 ConnectableObservable,它只在您调用 Connect 时订阅底层源。所以其他人可以订阅它,但它会等到你调用Connect 之后才会订阅源。 RefCount 为您管理。第一个订阅者将导致RefCount 调用Connect。第二个以上的订阅者将导致RefCount 共享相同的订阅。当订阅者全部退订时,RefCount 将来自源的Disconnect
  • IMO 这是一个比使用主题更简洁的解决方案。它更惯用,更具声明性和更少状态(至少在您自己的代码中)。正如@Brandon 所说,使用 RefCount 将确保底层连接仅在第一次订阅。
【解决方案2】:

我可能在这里过于简单化了,但让我按要求使用Subject

你的Thingy

public class Thingy
{
    private BehaviorSubject<bool> _statusSubject = new BehaviorSubject<bool>(false);    
    public IObservable<bool> Status
    {
        get
        {
            return _statusSubject;
        }
    }

    public void Init()
    {
        var c = new object();
        new Underlying().GetDeferredObservable(c).Subscribe(_statusSubject);
    }
}

伪造的Underlying

public class Underlying
{
    public IObservable<bool> GetDeferredObservable(object connection)
    {
        return Observable.DeferAsync<bool>(token => {
            return Task.Factory.StartNew(() => {
                Console.WriteLine("UNDERLYING ENGAGED");
                Thread.Sleep(1000);
                // Let's pretend there's some static on the line...
                return Observable.Return(true)
                    .Concat(Observable.Return(false))
                    .Concat(Observable.Return(true));
            }, token);
        });
    }
}

安全带:

void Main()
{
    var thingy = new Thingy();
    using(thingy.Status.Subscribe(stat => Console.WriteLine("Status:{0}", stat)))
    {
        Console.WriteLine("Waiting three seconds to Init...");
        Thread.Sleep(3000);
        thingy.Init();
        Console.ReadLine();
    }
}

输出:

Status:False
Waiting three seconds to Init...
UNDERLYING ENGAGED
Status:True
Status:False
Status:True

【讨论】:

  • 谢谢。与我正在尝试做的唯一功能区别(但我承认我没有在问题中说明)是当所有订阅者都离开时,主题需要从底层取消订阅......顺便说一句,在哪里DeferAsync 从何而来?
  • 啊,是的,我会做类似 Brandon 的事情,使用 RefCounted 可连接。我认为 DeferAsync 在 Rx 平台服务程序集中,但我是凭记忆进行的。
  • 它似乎在 System.Reactive.Linq 程序集中。如果您不需要CancellationToken(在这个伪造的示例中似乎就是这种情况),那么还有一个Defer 重载,它采用异步工厂方法。我真的很惊讶DeferAsync 不仅仅是Defer 的另一个重载。
  • @brandon 是的,为了权宜之计,我只是把它放在 linqpad 中 - 同意过载,我认为三个“异步变体”属于同一类别。我猜是不同的团队?
【解决方案3】:

嗯,玩过这个,我不认为我可以只用一个主题。

尚未完成测试/尝试,但这是我目前提出的似乎可行的方法,但它并不能保护我免受主题问题的影响,因为我仍在内部使用它。

public class ObservableRouter<T> : IObservable<T>
{
    ISubject<T> _Subject = new Subject<T>();
    Dictionary<IObserver<T>, IDisposable> _ObserverSubscriptions 
                               = new Dictionary<IObserver<T>, IDisposable>();
    IObservable<T> _ObservableSource;
    IDisposable _SourceSubscription;

    //Note that this can happen before or after SetSource
    public IDisposable Subscribe(IObserver<T> observer)
    {
        _ObserverSubscriptions.Add(observer, _Subject.Subscribe(observer));
        IfReadySubscribeToSource();
        return Disposable.Create(() => UnsubscribeObserver(observer));
    }

    //Note that this can happen before or after Subscribe
    public void SetSource(IObservable<T> observable)
    {
        if(_ObserverSubscriptions.Count > 0 && _ObservableSource != null) 
                  throw new InvalidOperationException("Already routed!");
        _ObservableSource = observable;
        IfReadySubscribeToSource();
    }

    private void IfReadySubscribeToSource()
    {
        if(_SourceSubscription == null &&
           _ObservableSource != null && 
           _ObserverSubscriptions.Count > 0)
        {
            _SourceSubscription = _ObservableSource.Subscribe(_Subject);
        }
    }

    private void UnsubscribeObserver(IObserver<T> observer)
    {
        _ObserverSubscriptions[observer].Dispose();
        _ObserverSubscriptions.Remove(observer);
        if(_ObserverSubscriptions.Count == 0)
        {
            _SourceSubscription.Dispose();
            _SourceSubscription = null;
        }
    }
}

【讨论】:

    猜你喜欢
    • 2019-12-29
    • 1970-01-01
    • 1970-01-01
    • 2017-08-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多