【问题标题】:How can I observe values in a non blocking way using Rx?如何使用 Rx 以非阻塞方式观察值?
【发布时间】:2013-05-24 14:48:19
【问题描述】:

我试图观察一个计时器,它的处理程序比时间间隔长。 为此,我想安排对某种线程池、任务池或其他东西的观察。

我尝试了线程池、任务池和新线程,但都没有工作。 有谁知道怎么做? 示例:

var disposable = Observable.Timer(TimeSpan.Zero, TimeSpan.FromMilliseconds(100)).ObserveOn(Scheduler.NewThread).
    Subscribe(x =>
      {
      count++;
    Thread.Sleep(TimeSpan.FromMilliseconds(1000));
  });

  Thread.Sleep(TimeSpan.FromSeconds(5));
  disposable.Dispose();
  if (count > 10 )
  {
    //hurray...
  }

【问题讨论】:

    标签: .net system.reactive


    【解决方案1】:

    您的要求是一个坏主意,因为您最终会耗尽可用资源(因为创建线程的速率 > 线程完成速率)。相反,您为什么不在前一个项目完成后安排一个新项目?

    在您的具体示例中,您需要将 IScheduler 传递给 Observable.Timer 而不是尝试使用 ObserveOn。

    【讨论】:

    • 创建率与完成率相同。只是实际操作需要时间。我想将所有观察结果放在某种线程池/任务池上,以避免创建开销。我可以让订阅以异步方式完成它的工作,我只是希望在 rx 中内置这种操作。
    • 任务池并不神奇,在这种情况下,您最终会填满任务池队列,直到它抛出错误
    【解决方案2】:

    当保罗说这是一个坏主意时,他是对的。从逻辑上讲,您正在创建排队操作可能会耗尽系统资源的情况。您甚至可以发现它在您的计算机上工作,但在客户的计算机上失败。可用内存、32 位/64 位处理器等都会影响代码。

    但是,很容易修改您的代码以使其执行您想要的操作。

    首先,只要观察者在下一个计划事件之前完成,Timer 方法就会正确地安排计时器事件。如果观察者还没有完成,那么计时器将等待。请记住,可观察计时器是“冷”可观察对象,因此对于每个订阅的观察者,实际上都有一个新的可观察计时器。这是一对一的关系。

    此行为可防止计时器无意中耗尽您的资源。

    因此,按照您当前定义的代码,OnNext 每 1000 毫秒调用一次,而不是每 100 毫秒调用一次。

    现在,要让代码以 100 毫秒的时间表运行,请执行以下操作:

    Observable
        .Timer(TimeSpan.Zero, TimeSpan.FromMilliseconds(100))
        .Select(x => Scheduler.NewThread.Schedule(() =>
        {
            count++;
            Thread.Sleep(TimeSpan.FromMilliseconds(1000));
        }))
        .Subscribe(x => { });
    

    实际上,此代码是 IObservable<IDisposable>,其中每个一次性操作都是需要 1000 毫秒才能完成的预定操作。

    在我的测试中,这运行得很好并且正确地增加了计数。

    我确实尝试过耗尽我的资源,发现将计时器设置为每毫秒运行一次,我很快得到了System.OutOfMemoryException,但我发现如果我将设置更改为每两毫秒运行一次,代码就会运行。然而,在代码运行并创建了大约 500 个新线程时,这确实使用了超过 500 MB 的 RAM。一点都不好看。

    谨慎行事!

    【讨论】:

    • 我会将调度放在订阅处理程序中,而不是使用 Select。你有什么理由不那样做?
    • @NiallConnaughton - 差不多两年前我回答了这个问题。我不记得了。现在看来确实有点奇怪。
    【解决方案3】:

    如果你真的在不断地产生价值,而不是消耗它们,那么正如所指出的,你正在走向麻烦。如果你不能减慢生产速度,那么你需要看看如何更快地消耗它们。也许您希望对观察者进行多线程处理以使用多个内核?

    如果您对观察者进行多线程处理,您可能需要小心处理乱序事件。您将同时处理多个通知,并且所有关于哪个处理首先完成(或首先进入某个竞争条件临界状态)的赌注都没有了。

    如果您没有必须处理流中的每个事件,请查看浮动的 ObserveLatestOn 的几个不同实现。有线程在讨论它herehere

    ObserveLatestOn 将丢弃除观察者处理先前通知时发生的最新通知之外的所有通知。当观察者处理完之前的通知后,它会收到最新的通知,并错过所有之间发生的通知。

    这样做的好处是它可以防止来自比消费者更快的生产者的压力累积。如果消费者因为负载而变慢,那么处理更多通知只会变得更糟。删除不需要的通知可能会使负载降低到消费者可以跟上的程度。

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-10-29
    • 2015-03-31
    • 2020-10-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多