此版本处理了几个问题:
-
有一种比赛条件可能会导致比赛失利。如果观察者在一个线程上观察到一个事件,而源观察者在另一个线程产生一个新的观察者,如果您不使用任何类型的同步,您最终可能会调用
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);
}
});
}