【问题标题】:Combine two Observables, but only when the first obs is immediately preceded by the second合并两个 Observable,但仅当第一个 obs 紧接在第二个之前
【发布时间】:2020-02-23 22:31:08
【问题描述】:

我的目标可能最容易用大理石图来解释。 我有两个 Observable,xs 和 ys。我想返回一个名为 rs 的 Observable。

xs --x---x---x---x-------x---x---x-
      \           \       \
ys ----y-----------y---y---y-------
       |           |       |
rs ----x-----------x-------x-------
       y           y       y

所以,我需要类似于 CombineLatest 的东西,除了它应该只在 xs 后跟 ys 时触发。该模式之外的其他 xs 或 ys 不应触发输出,应丢弃。

CombineLatest、Zip 或 And/Then/When 不做我需要的事情,我找不到任何方法来指定更复杂的连接结构。

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    我最终使用了 Join。

    var rs = xs.Join(ys, 
                     _ => xs.Merge(ys),
                     _ => Observable.Empty<Unit>(),
                     Tuple.Create);
    

    这是一篇很好的文章,它解释了 Join 的工作原理: http://blogs.microsoft.co.il/blogs/bnaya/archive/2012/04/04/rx-join.aspx

    【讨论】:

      【解决方案2】:

      对于此类问题,提供测试场景非常有用。你已经给出了你想要的一个很棒的小大理石图,然后将它转换为测试用例真的很容易。然后其他论坛读者只需获取测试代码,实施他们的想法并了解他们是否满足您的要求。

      [Test]
      public void Should_only_get_latest_value_from_Y_but_not_while_x_produes()
      {
          //            11111111112222222222333
          //   12345678901234567890123456789012
          //xs --x---x---x---x-------x---x---x-
          //      \           \       \
          //ys ----y-----------y---y---y-------
          //       |           |       |
          //rs ----x-----------x-------x-------
          //       y           y       y
      
      
          var testScheduler = new TestScheduler();
          //            11111111112222222222333
          //   12345678901234567890123456789012
          //xs --x---x---x---x-------x---x---x-
          var xs = testScheduler.CreateColdObservable(
              new Recorded<Notification<char>>(3, Notification.CreateOnNext('1')),
              new Recorded<Notification<char>>(7, Notification.CreateOnNext('2')),
              new Recorded<Notification<char>>(10, Notification.CreateOnNext('3')),
              new Recorded<Notification<char>>(15, Notification.CreateOnNext('4')),
              new Recorded<Notification<char>>(23, Notification.CreateOnNext('5')),
              new Recorded<Notification<char>>(27, Notification.CreateOnNext('6')),
              new Recorded<Notification<char>>(31, Notification.CreateOnNext('7')));
      
          //            11111111112222222222333
          //   12345678901234567890123456789012
          //ys ----y-----------y---y---y-------
          var ys = testScheduler.CreateColdObservable(
              new Recorded<Notification<char>>(5, Notification.CreateOnNext('A')),
              new Recorded<Notification<char>>(17, Notification.CreateOnNext('B')),
              new Recorded<Notification<char>>(21, Notification.CreateOnNext('C')),
              new Recorded<Notification<char>>(25, Notification.CreateOnNext('D')));
      
      
          //Expected :
          //Tick  x   y
          //5     1   A
          //17    4   B
          //25    5   D
          var expected = new[]
          {
              new Recorded<Notification<Tuple<char, char>>>(5, Notification.CreateOnNext(Tuple.Create('1', 'A'))),
              new Recorded<Notification<Tuple<char, char>>>(17, Notification.CreateOnNext(Tuple.Create('4', 'B'))),
              new Recorded<Notification<Tuple<char, char>>>(25, Notification.CreateOnNext(Tuple.Create('5', 'D')))
          };
      
          var observer = testScheduler.CreateObserver<Tuple<char, char>>();
      
          //Passes HOT, fails Cold. Doesn't meet the requirements due to Timeout anyway.
          //xs.Select(x => ys.Take(1)
          //                    .Timeout(TimeSpan.FromSeconds(0.5),
          //                            Observable.Empty<char>())
          //                    .Select(y => Tuple.Create(x, y))
          //    )
          //    .Switch()
          //    .Subscribe(observer);
      
          //Passes HOT. Passes Cold
          xs.Join(ys,
                  _ => xs.Merge(ys),
                  _ => Observable.Empty<Unit>(),
                  Tuple.Create)
              .Subscribe(observer);
      
          testScheduler.Start();
      
          //You may want to Console.WriteLine out the two collections to validate for yourself.
      
          CollectionAssert.AreEqual(expected, observer.Messages);
      }
      

      附:顺便说一句,Merge 很好地使用了答案。

      【讨论】:

      • 使用ReactiveTest中的OnNext方法避免重复new Recorded&lt;Notification&lt;char&gt;&gt;
      【解决方案3】:

      换一种说法,每个x需要取第一个 y 在y之前没有额外的x。第一个要求建议Take(1)。届时,您将拥有一个IObservable&lt;IObservable&lt;...&gt;&gt;。第二个需求和前两个生成的类型指向Switch。综上所述,您应该能够将图表与:

      var rs = xs.Select(x => ys.Take(1)
                                .Select(y => Tuple.Create(x, y))
                        )
                 .Switch()
      

      【讨论】:

      • 这些事件发生的时间不规律,所以不能使用Timeout。但是我可以合并 ys 和 xs,Take(1),然后过滤掉我有两个来自 xs 的值的情况。它不漂亮,我需要向 xs 和 ys 添加类型信息,但它应该可以工作。
      • 它的工作方式类似于窗口。我想知道我是否可以用它做一些更优雅的东西。
      • @oillio:我建议您将其添加为您的问题的答案:)
      • 嘿,由于超时,我投了反对票。这不是 IMO 的正确答案。
      • @LeeCampbell 你是对的。出于某种原因,我最初读到的问题是说 y 必须接近 x。进一步审查,没有这样的时间要求,所以超时是不正确的。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-04-04
      相关资源
      最近更新 更多