这就是我想出的。非常欢迎任何 cmets、想法、建议和批评。最有趣的部分是 IObserver<TValue, TResult> 公共接口已经存在于 Rx 库中,但它仅在 Notification 对象中使用。我所做的是,我创建了IObservable<TValue, TResult> 对应物 SelectiveSubject 来处理调用观察者的逻辑,直到其中一个返回 true 和 ToSelective 方法扩展。我真的很惊讶它没有在图书馆完成,至少在IObservable<TValue, TResult> 部分。毕竟,IObserver<TValue, TResult> 存在。
public interface IObservable<out TValue, in TResult> {
IDisposable Subscribe(IObserver<TValue, TResult> objObserver);
}
internal class SelectiveSubject<T> : IObserver<T>, IObservable<T, bool> {
private readonly LinkedList<IObserver<T, bool>> _ObserverList;
public SelectiveSubject() {
_ObserverList = new LinkedList<IObserver<T, bool>>();
}
public void OnNext(T value) {
lock(_ObserverList) {
foreach(IObserver<T, bool> objObserver in _ObserverList) {
if(objObserver.OnNext(value)) {
break;
}
}
}
}
public void OnError(Exception exception) {
lock(_ObserverList) {
foreach(IObserver<T, bool> objObserver in _ObserverList) {
if(objObserver.OnError(exception)) {
break;
}
}
}
}
public void OnCompleted() {
lock(_ObserverList) {
foreach(IObserver<T, bool> objObserver in _ObserverList) {
if(objObserver.OnCompleted()) {
break;
}
}
}
}
public IDisposable Subscribe(IObserver<T, bool> objObserver) {
LinkedListNode<IObserver<T, bool>> objNode;
lock(_ObserverList) {
objNode = _ObserverList.AddLast(objObserver);
}
return Disposable.Create(() => {
lock(objNode.List) {
objNode.List.Remove(objNode);
}
});
}
}
public static IObservable<T, bool> ToSelective<T>(this IObservable<T> objThis) {
var objSelective = new SelectiveSubject<T>();
objThis.Subscribe(objSelective);
return objSelective;
}
现在用法就这么简单
IConnectableObservable<int> objGenerator = Observable.Generate(0, i => i < 100, i => i + 1, i => i).Publish();
IObservable<int, bool> objSelective = objGenerator.ToSelective();
var objDisposableList = new CompositeDisposable(2) {
objSelective.Subscribe(i => {
Console.Write("1");
if(i % 2 == 0) {
Console.Write("!");
return true;
}
else {
Console.Write(".");
return false;
}
}),
objSelective.Subscribe(i => {
Console.Write("2");
Console.Write("!");
return true;
})
};
objGenerator.Connect();
objDisposableList.Dispose();
在示例中,第一个订阅者负责处理序列中的所有其他值,第二个订阅者负责其余部分。