【问题标题】:Reactive: Trying to understand how Subject<T> work反应式:试图理解 Subject<T> 是如何工作的
【发布时间】:2012-08-21 15:18:23
【问题描述】:

试图了解Subject&lt;T&gt;ReplaySubject&lt;T&gt; 和其他的工作原理。这是一个例子:

(主题是可观察的观察者)

public IObservable<int> CreateObservable()
{
     Subject<int> subj = new Subject<int>();                // case 1
     ReplaySubject<int> subj = new ReplaySubject<int>();    // case 2

     Random rnd = new Random();
     int maxValue = rnd.Next(20);
     Trace.TraceInformation("Max value is: " + maxValue.ToString());

     subj.OnNext(-1);           // specific value

     for(int iCounter = 0; iCounter < maxValue; iCounter++)
     {
          Trace.TraceInformation("Value: " + iCounter.ToString() + " is about to publish");
          subj.OnNext(iCounter);
      }

      Trace.TraceInformation("Publish complete");
      subj.OnComplete();

      return subj;
 }

 public void Main()
 {
     //
     // First subscription
     CreateObservable()
         .Subscribe(
               onNext: (x)=>{  
                  Trace.TraceInformation("X is: " + x.ToString()); 
      });

     //
     // Second subscribe
     CreateObservable()
         .Subscribe(
               onNext: (x2)=>{  
                  Trace.TraceInformation("X2 is: " + x.ToString());
     });

案例 1:奇怪的情况是 - 当我使用 Subject&lt;T&gt; 时没有订阅 (???) - 我从来没有看到“X 是:”文本 - 我只看到“值是:”和“Max value is"... 为什么Subject&lt;T&gt; 不推送值到订阅?

案例 2:如果我使用 ReplaySubject&lt;T&gt; - 我确实看到了订阅中的值,但我无法将 Defer 选项应用于任何内容。不是Subject,也不是 Observable....所以每个订阅都会收到不同的值,因为CreateObservable 函数是 cold 可观察的。 Defer 在哪里?

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    每当您需要凭空创建一个可观察对象时,Observable.Create 应该是首先想到的。主体进图分两种情况:

    • 您需要某种“可寻址端点”来提供数据,以便所有订阅者都能接收到它。将此与具有调用端(通过委托调用)和订阅端(通过委托结合 +- 和 -= 语法)的 .NET 事件进行比较。你会发现在很多情况下,你可以使用 Observable.Create 达到同样的效果。

    • 您需要在查询管道中多播消息,通过查询逻辑中的许多分支有效地共享可观察序列,而不会触发多个订阅。 (想想为你的宿舍订阅一次你最喜欢的杂志,然后在信箱后面放一台复印机。你仍然需要支付一次订阅费,尽管你所有的朋友都可以在信箱上阅读通过 OnNext 发送的杂志。)

    此外,在很多情况下,Rx 中已经有一个内置原语可以完全满足您的需求。例如,有 From* 工厂方法来桥接现有概念(例如事件、任务、异步方法、可枚举序列),其中一些使用了隐藏的主题。对于多播逻辑的第二种情况,有 Publish、Replay 等一系列运算符。

    【讨论】:

      【解决方案2】:

      您需要注意代码的执行时间。

      在“案例 1”中,当您使用 Subject&lt;T&gt; 时,您会注意到对 OnNextOnCompleted 的所有调用在 CreateObservable 方法返回 observable 之前完成。由于您使用的是Subject&lt;T&gt;,这意味着任何后续订阅都将丢失所有值,因此您应该期望得到您所得到的 - 什么都没有。

      您必须延迟对该主题的操作,直到您订阅了观察者。为此,请使用 Create 方法。方法如下:

      public IObservable<int> CreateObservable()
      {
          return Observable.Create<int>(o =>
          {
              var subj = new Subject<int>();
              var disposable = subj.Subscribe(o);
      
              var rnd = new Random();
              var maxValue = rnd.Next(20);
              subj.OnNext(-1);
              for(int iCounter = 0; iCounter < maxValue; iCounter++)
              {
                  subj.OnNext(iCounter);
              }
              subj.OnCompleted();
      
              return disposable;
          });
      }
      

      为简洁起见,我删除了所有跟踪代码。

      所以现在,对于每个订阅者,您都可以在 Create 方法中获得新的代码执行,并且您现在可以从内部 Subject&lt;T&gt; 中获取值。

      使用Create 方法通常是创建从方法返回的可观察对象的正确方法。

      或者,您可以使用ReplaySubject&lt;T&gt; 并避免使用Create 方法。然而,由于多种原因,这并不吸引人。它强制在创建时计算整个序列。这给了你一个冷的 observable,你可以在不使用回放主题的情况下更有效地制作它。

      现在,顺便说一句,您应该尽量避免使用主题。一般规则是,如果您使用主题,那么您做错了什么。 CreateObservable 方法最好这样写:

      public IObservable<int> CreateObservable()
      {
          return Observable.Create<int>(o =>
          {
              var rnd = new Random();
              var maxValue = rnd.Next(20);
              return Observable.Range(-1, maxValue + 1).Subscribe(o);
          });
      }
      

      根本不需要主题。

      如果这有帮助,请告诉我。

      【讨论】:

      • var 一次性 = subj.Subscribe(o);意味着:在o 中的每个OnNext() 上都会在subj 中提高OnNext() ???
      • 不,反过来。 o 正在订阅subj,因此subj 上的每个OnNext() 都会在o 上的订阅上调用OnNext 回调。
      • 我猜我对 Reactive 的理解有问题......我认为 Reactive 是关于“未来的收藏”。当您没有可用的整个集合时运行一些逻辑/命令。 IObservable 是集合的,当发生变化并引发OnNext 时,您可以在集合源中将SubscribeOnNext 事件(?)。我是对的还是我在描述 EAP(基于事件的模式)?
      • 我不太确定你在说什么。考虑反应式集合的方式是它们是跨越一段时间的集合,而不是一次瞬时集合(IEnumerable)。可观察对象可以“冷”(值在订阅者连接时生成)或“热”(值独立于订阅者生成)。如果您有一个在订阅者订阅之前产生其所有值的可观察对象,那么它们只会“错过”对方。考虑一个按钮被点击之前一个事件处理程序被附加 - 你不会得到那些以前的点击。
      • @Wilka - 我很高兴。不过从现在开始大约需要 10 个小时左右。
      猜你喜欢
      • 1970-01-01
      • 2020-05-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多