【问题标题】:Is it possible to send a message not to all observers using .net Rx?是否可以不使用 .net Rx 向所有观察者发送消息?
【发布时间】:2015-06-10 14:18:30
【问题描述】:

我有一种情况,有一个 observable,假设有 10 个观察者附加到它上面。我只想将消息发送给每个新的观察者,直到观察者以某种方式对 observable 说它识别出消息并将处理它。此时我想停止向其他观察者发送此消息。

换句话说,每个观察者都知道如何处理特定类型的消息,并且每个观察者都会接受并处理它识别的消息。其他人不需要在识别它开始处理后接收它。如何使用响应式扩展来实现这种情况?我认为我们需要某种通知返回到 observable,但我不知道怎么做。

【问题讨论】:

  • 这听起来更像是chain of responsibility 的问题,而不是 Rx。我认为您不会与 Rx 进行任何双向通信 - 您能做的最好的事情就是在订阅观察者之前过滤序列。一些特定的代码可能在这里有用。
  • 确实这不是一个 Rx 问题。您可以订阅一个单独的“控制器”可观察对象,该可观察对象保留候选人列表,并且每当有项目到达时,控制器都会一个接一个地尝试候选者,直到有人接受该项目。或者每个观察者都可以使用.Where()将自己“附加”到链上,当它“认领”该项目时,它会返回false
  • 我同意这是在可观察对象之上的责任链。我只是希望 Rx 有一个快速的解决方案。
  • @alex.49.98 - 你没有以 Rx 的方式思考。您正在以面向对象的方式思考。 Rx 是功能性的。你可以通过在 observable 中添加 .Where(...) 子句来过滤你的消息来做这种事情。
  • 我知道我可以使用 .Where(..) 过滤消息。我正在考虑一种优化代码的方法。消息识别部分需要一定的时间。如果找到特定消息的观察者,我绝对不需要重复它。基本上看起来差不多。我必须将我的选择代码放入 .Where(...) 或 .Select(...) 中,然后使用 bool OnNext(T value) 之类的方法签名在循环中调用自定义观察者。每个观察者将检查消息是否是它的并返回真或假。一旦观察者返回 true,循环就会停止。

标签: c# .net system.reactive observer-pattern


【解决方案1】:

这就是我想出的。非常欢迎任何 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();

在示例中,第一个订阅者负责处理序列中的所有其他值,第二个订阅者负责其余部分。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2012-01-19
    • 2014-01-28
    • 1970-01-01
    • 1970-01-01
    • 2018-12-16
    • 2018-07-03
    • 1970-01-01
    相关资源
    最近更新 更多