【问题标题】:Perform action after all Observers complete on Observable sequence在所有 Observer 完成 Observable 序列后执行操作
【发布时间】:2017-02-18 10:55:50
【问题描述】:

我有一些可观察的序列,例如:

var period = TimeSpan.FromSeconds(0.5);
var observable = Observable
    .Interval(period)
    .Publish()
    .RefCount();

我想在后台线程上对该序列的元素执行一些硬计算,并在所有计算完成后执行一些最终操作。所以我想要这样的东西:

observable.ObserveOn(Scheduler.Default).Subscribe(i => ComplexComputation1(i));
observable.ObserveOn(Scheduler.Default).Subscribe(i => ComplexComputation2(i));
// next observer must be called only after ComplexComputation1/2 complete on input i
observable.Subscribe(i => FinalAction(i));

我可以在 Rx 中执行此操作吗?或者这可能违反了反应式编程的一些原则,我应该在这种情况下使用另一种方法?

【问题讨论】:

    标签: c# system.reactive observable


    【解决方案1】:

    在反应模式中具有计算有序序列是非常危险的。

    您可以做的一件事是让每个复杂计算在完成后发出一个事件。然后你可以有一个消费观察者,一旦他收到前面步骤完成的消息,他就会执行他的计算。


    另一种可能的解决方案是创建一个定期触发的具体序列块。这降低了解决方案的并行性。

    observable.ObserveOn(Scheduler.Default).Subscribe(i => 
    {     
        ComplexComputation1(i));
        ComplexComputation2(i));
        FinalAction(i);
    }
    

    【讨论】:

    • 我想过这种方法。但实际上,我不知道有多少观察者会处理每个元素。这种解决方案可以应用于这种情况吗?
    • 我添加了一个序列块方法作为可能的解决方案。
    【解决方案2】:

    为了测试这一点,我创建了以下方法来帮助说明事件的顺序:

    public void ComplexComputation1(long i)
    {
        Console.WriteLine("Begin ComplexComputation1");
        Thread.Sleep(100);
        Console.WriteLine("End ComplexComputation1");
    }
    
    public void ComplexComputation2(long i)
    {
        Console.WriteLine("Begin ComplexComputation2");
        Thread.Sleep(100);
        Console.WriteLine("End ComplexComputation2");
    }
    
    public void FinalAction(long i)
    {
        Console.WriteLine("Begin FinalAction");
        Thread.Sleep(100);
        Console.WriteLine("End FinalAction");
    }
    

    你的原始代码是这样运行的:

    开始最终行动 开始复杂计算1 开始复杂计算2 结束复杂计算2 结束最终行动 结束复杂计算1 开始最终行动 开始复杂计算1 开始复杂计算2 结束最终行动 结束复杂计算2 结束复杂计算1 开始最终行动 开始复杂计算1 开始复杂计算2 结束复杂计算2 结束复杂计算1 结束最终行动 ...

    强制代码在单个后台线程上按顺序运行很容易。只需使用EventLoopScheduler

    var els = new EventLoopScheduler();
    
    observable.ObserveOn(els).Subscribe(i => ComplexComputation1(i));
    observable.ObserveOn(els).Subscribe(i => ComplexComputation2(i));
    // next observer must be called only after ComplexComputation1/2 complete on input i
    observable.ObserveOn(els).Subscribe(i => FinalAction(i));
    

    这给了:

    开始复杂计算1 结束复杂计算1 开始复杂计算2 结束复杂计算2 开始最终行动 结束最终行动 开始复杂计算1 结束复杂计算1 开始复杂计算2 结束复杂计算2 开始最终行动 结束最终行动 开始复杂计算1 结束复杂计算1 开始复杂计算2 结束复杂计算2 开始最终行动 结束最终行动

    但是你一介绍Scheduler.Default就不行了。

    或多或少简单的选择是这样做:

    var cc1s = observable.ObserveOn(Scheduler.Default).Select(i => { ComplexComputation1(i); return Unit.Default; });
    var cc2s = observable.ObserveOn(Scheduler.Default).Select(i => { ComplexComputation2(i); return Unit.Default; });
    
    observable.Zip(cc1s.Zip(cc2s, (cc1, cc2) => Unit.Default), (i, cc) => i).Subscribe(i => FinalAction(i));
    

    按预期工作。

    你会得到这样一个很好的序列:

    开始复杂计算1 开始复杂计算2 结束复杂计算1 结束复杂计算2 开始最终行动 结束最终行动 开始复杂计算2 开始复杂计算1 结束复杂计算2 结束复杂计算1 开始最终行动 结束最终行动 开始复杂计算1 开始复杂计算2 结束复杂计算2 结束复杂计算1 开始最终行动 结束最终行动

    【讨论】:

      【解决方案3】:

      这似乎是一个嵌套 observable 被展平(SelectMany/Merge/Concat)和 Zip 组合的简单案例

      在这里,我冒昧地假设 Long Running 方法返回 Task。 但是,如果他们不这样做,那么慢阻塞同步方法可以改为使用 Observable.Start(()=>ComplexComputation1(x)) 包装。

      void Main()
      {
          var period = TimeSpan.FromSeconds(0.5);
          var observable = Observable
              .Interval(period)
              .Publish()
              .RefCount();
      
          var a = observable.Select(i => ComplexComputation1(i).ToObservable())
                      .Concat();
          var b = observable.Select(i => ComplexComputation2(i).ToObservable())
                      .Concat();
      
          a.Zip(b, Tuple.Create)
              .Subscribe(pair => FinalAction(pair.Item1, pair.Item2));
      }
      
      // Define other methods and classes here
      Random rnd = new Random();
      private async Task<long> ComplexComputation1(long i)
      {
          await Task.Delay(rnd.Next(50, 1000));
          return i;
      }
      private async Task<long> ComplexComputation2(long i)
      {
          await Task.Delay(rnd.Next(50, 1000));
          return i;
      }
      
      private void FinalAction(long a, long b)
      {
      
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-05-17
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多