【问题标题】:Rx produce and consume on different threadsRx 在不同的线程上生产和消费
【发布时间】:2014-12-02 11:11:52
【问题描述】:

我试图通过此处的示例代码来简化我的问题。我有一个生产者线程不断地输入数据,我试图在批处理之间有时间延迟来批处理它,以便 UI 有时间渲染它。但是结果并不像预期的那样,生产者和消费者似乎在同一个线程上。

我不希望批处理缓冲区在正在生成的线程上休眠。试过SubscribeOn 并没有太大帮助。我在这里做错了什么,如何让它在生产者和消费者线程上打印不同的线程 ID。

static void Main(string[] args)
{
    var stream = new ReplaySubject<int>();

    Task.Factory.StartNew(() =>
    {
        int seed = 1;
        while (true)
        {
            Console.WriteLine("Thread {0} Producing {1}",
                Thread.CurrentThread.ManagedThreadId, seed);

            stream.OnNext(seed);
            seed++;

            Thread.Sleep(TimeSpan.FromMilliseconds(500));
         }
    });

    stream.Buffer(5).Do(x =>
    {
        Console.WriteLine("Thread {0} sleeping to create time gap between batches",
            Thread.CurrentThread.ManagedThreadId);

        Thread.Sleep(TimeSpan.FromSeconds(2));
    })
    .SubscribeOn(NewThreadScheduler.Default).Subscribe(items =>
    {
        foreach (var item in items)
        {
            Console.WriteLine("Thread {0} Consuming {1}",
                Thread.CurrentThread.ManagedThreadId, item);
        }
    });
    Console.Read();
}

【问题讨论】:

    标签: c# .net-4.0 system.reactive


    【解决方案1】:

    了解ObserveOnSubscribeOn 之间的区别是这里的关键。请参阅 - ObserveOn and SubscribeOn - where the work is being done 以获得对这些的深入解释。

    此外,您绝对不想在您的 Rx 中使用 Thread.Sleep。或任何地方。曾经。 Do 几乎一样邪恶,但Thead.Sleep 几乎总是完全邪恶。 Buffer 有几个你想使用的重载 - 这些包括基于时间的重载和接受计数限制的重载时间限制,当达到其中任何一个时返回一个缓冲区。基于时间的缓冲将在生产者和消费者之间引入必要的并发性——也就是说,在与生产者不同的线程上将缓冲区传递给它的订阅者。

    另请参阅这些问题和答案,它们对保持消费者响应式进行了很好的讨论(在 WPF 的上下文中,但这些要点通常适用)。

    上面的最后一个问题专门使用了基于时间的缓冲区过载。正如我所说,在调用链中使用BufferObserveOn 将允许您在生产者和消费者之间添加并发。您仍然需要注意缓冲区的处理速度仍然足够快,以免在缓冲区订阅者上建立队列。

    如果队列确实增加,您需要考虑应用背压、删除更新和/或合并更新的方法。这是一个太大的话题,无法在这里进行深入讨论 - 但基本上你可以:

    首先查看适当的缓冲是否有帮助,然后考虑在源头限制/合并事件(UI 只能显示这么多信息) - 然后考虑更智能的合并,因为这可能会变得非常复杂。 https://github.com/AdaptiveConsulting/ReactiveTrader 是使用一些高级合并技术的项目的一个很好的例子。

    【讨论】:

    • 我无法删除事件或获取最新值,我需要将所有消息传递到 UI,因为它们形成单独的行。添加 ObserveOn 似乎可以解决问题,为什么说 do with sleep 是邪恶的?有没有更好的方法在每批之间产生延迟?
    • @anivas - 你可以通过使用 ZipInterval observable 来造成延迟。
    • @Enigmativity 尝试了该方法,间隔中的计时器会造成内存泄漏。
    • @anivas - 当你处理掉 observable 时,它​​应该都消失了。
    • @anivas 他说的 ^ :) 睡眠通常是邪恶的,因为它浪费了线程资源。线程并不便宜,您可能会导致不必要的线程产生。 await Task.Delay(...) 通常更好,因为它允许线程做其他有用的工作——但即便如此,引入这样的等待几乎总是一种代码味道,更好的设计通常可以避免。在 Rx 中,它会阻塞在该线程上运行的整个链,其中可能包括生产者和其他订阅者,具体取决于所使用的调度程序。
    【解决方案2】:

    虽然其他答案是正确的,但我想确定您的实际问题可能是对 Rx 行为的误解。将生产者置于睡眠状态会阻止对OnNext 的后续调用,并且您似乎假设Rx 会同时自动调用OnNext,但实际上它并不是有很好的理由。实际上,Rx 有一个合约需要序列化通知。

    有关详细信息,请参阅Rx Design Guidelines 中的 §§4.2、6.7。

    最终,您似乎正在尝试从 Rxx 实现 BufferIntrospective 运算符。此运算符允许您传入一个引入并发的调度程序,类似于ObserveOn,以在生产者和消费者之间创建并发边界。 BufferIntrospective 是一种动态背压策略,它根据观察者不断变化的延迟推出不同大小的批次。当观察者处理当前批次时,操作者缓冲所有传入的并发通知。为了实现这一点,运算符利用OnNext 是阻塞调用(根据§4.2 合同)这一事实,因此该运算符应尽可能靠近查询边缘应用,通常在您之前拨打Subscribe

    正如 James 所描述的,您可以将其称为“智能缓冲”策略本身,或者将其视为实施此类策略的基准;例如,我还定义了一个 SampleIntrospective 运算符,它会删除每批中除最后一个通知之外的所有通知。

    【讨论】:

    • 有用的 cmets,这些设计指南是 Rx dev 恕我直言的必读。 +1
    【解决方案3】:

    ObserveOn 可能是你想要的。它需要一个 SynchronizationContext 作为参数,它应该是你的 UI 的 SynchronizationContext。如果不知道如何获取,请看Using SynchronizationContext for sending events back to the UI for WinForms or WPF

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-08-27
      • 1970-01-01
      • 2018-03-08
      • 1970-01-01
      • 1970-01-01
      • 2018-09-24
      相关资源
      最近更新 更多