代码之家  ›  专栏  ›  技术社区  ›  bradgonesurfing

如何在C中为IObservable<IObserver<T>>的反应式扩展实现Switch#

  •  0
  • bradgonesurfing  · 技术社区  · 13 年前

    在反应式扩展中,我们有

    IObservable<T> Switch(this IObservable<IObservable<T>> This)
    

    我希望实现

    IObserver<T> Switch(this IObservable<IObserver<T>> This)
    

    这将把传出的事件切换到不同的观察者,但是 以单个观察者的身份呈现。

    2 回复  |  直到 13 年前
        1
  •  3
  •   Brandon    13 年前

    此版本处理了几个问题:

    • 有一种比赛条件可能会导致比赛失利。如果观察者在一个线程上观察到一个事件,而源观察者在另一个线程产生一个新的观察者,如果您不使用任何类型的同步,您最终可能会调用 OnCompleted 在一个线程上的当前观察器上,刚好在另一个线程调用之前 OnNext 在同一个观察者身上。这将导致事件丢失。

    • 与上述内容相关的是,默认情况下,观察者不是线程安全的。您不应该同时调用观察者,否则将违反主Rx协议。在没有任何锁定的情况下,用户可能会呼叫 已完成 上 currentObserver 而另一个线程正在调用 在下一个 在同一个观察者身上。开箱即用,这类事情可以通过使用同步主题来解决。但是,由于前面的问题也需要同步,所以我们可以只使用一个简单的互斥。

    • 我们需要一种方法来取消订阅源observable。我假设,当生成的观察者完成(或出错)时,这是取消订阅源代码的好时机,因为我们的观察者被告知不要再发生事件。

    这是代码:

    public static IObserver<T> Switch<T>(this IObservable<IObserver<T>> source)
    {
        var mutex = new object();
        var current = Observer.Create<T>(x => {});
        var subscription = source.Subscribe(o =>
        {
            lock (mutex)
            {
               current.OnCompleted();
               current = o;
            }
        });
    
        return Observer.Create<T>(
            onNext: v =>
            {
                lock(mutex)
                {                
                  current.OnNext(v);
                }
            },
            onCompleted: () =>
            {
                 subscription.Dispose();
                 lock (mutex)
                 {
                     current.OnCompleted();
                 }
            },
            onError: e =>
            {
                 subscription.Dispose();
                 lock (mutex)
                 {
                     current.OnError(e);
                 }
            });
    }
    
        2
  •  1
  •   bradgonesurfing    13 年前
    public static IObserver<T> Switch<T>(this IObservable<IObserver<T>> This)
    {
        IObserver<T> currentObserver = Observer.Create<T>(x => { });
    
        This.Subscribe(o => { currentObserver.OnCompleted(); currentObserver = o; });
    
    
        return Observer.Create<T>
            ( onNext: v => currentObserver.OnNext(v)
            , onCompleted: () => currentObserver.OnCompleted()
            , onError: v => currentObserver.OnError(v));
    }