【问题标题】:Why each observation delegate runs on a new thread为什么每个观察委托都在新线程上运行
【发布时间】:2016-01-24 03:15:25
【问题描述】:

在 Rx 中,当对 ObserveOn 方法使用 Scheduler.NewThread 时,让每个 Observation 委托 (OnNext) 在新线程上运行有什么好处,而 Rx 已经保证 OnNext 永远不会重叠。如果每个 OnNext 都将被一个接一个地调用,为什么需要为每个 OnNext 一个新的线程。

我明白为什么要在不同于订阅和应用程序线程的线程上运行观察委托,但在它们永远不会并行运行的情况下在新线程上运行每个观察委托?......没有意义我还是我在这里错过了什么?

举例

using System;
using System.Linq;
using System.Reactive.Concurrency;
using System.Reactive.Linq;
using System.Threading;

namespace RxTesting
{
    class Program
    {
        static void Main(string[] args)
        {
            Console.WriteLine("Application Thread : {0}", Thread.CurrentThread.ManagedThreadId);

            var numbers = from number in Enumerable.Range(1,10) select Process(number);

            var observableNumbers = numbers.ToObservable()
                .ObserveOn(Scheduler.NewThread)
                .SubscribeOn(Scheduler.NewThread);

            observableNumbers.Subscribe(
                n => Console.WriteLine("Consuming : {0} \t on Thread : {1}", n, Thread.CurrentThread.ManagedThreadId));

            Console.ReadKey();
        }

        private static int Process(int number)
        {
            Thread.Sleep(500);
            Console.WriteLine("Producing : {0} \t on Thread : {1}", number,
                              Thread.CurrentThread.ManagedThreadId);

            return number;
        }
    }
}

上面的代码产生以下结果。请注意,Consuming 每次都是在一个新线程上完成的。

Application Thread : 8
Producing : 1    on Thread : 9
Consuming : 1    on Thread : 10
Producing : 2    on Thread : 9
Consuming : 2    on Thread : 11
Producing : 3    on Thread : 9
Consuming : 3    on Thread : 12
Producing : 4    on Thread : 9
Consuming : 4    on Thread : 13
Producing : 5    on Thread : 9
Consuming : 5    on Thread : 14
Producing : 6    on Thread : 9
Consuming : 6    on Thread : 15
Producing : 7    on Thread : 9
Consuming : 7    on Thread : 16
Producing : 8    on Thread : 9
Consuming : 8    on Thread : 17
Producing : 9    on Thread : 9
Consuming : 9    on Thread : 18
Producing : 10   on Thread : 9
Consuming : 10   on Thread : 19

【问题讨论】:

  • 您在哪里看到让每个观察委托 (OnNext) 在自己的线程上运行的代码示例?
  • @Robert Harvey:我已经包含了重现问题的示例代码
  • 那你是说每次消费都要等上一次消费完成?这没有任何意义。如果您在单独的线程上消费每个观察,则该消费可以异步执行,而其他消费发生在它们自己的线程上。如果您在其自己的线程上使用每个 OnNext,则该线程立即释放 OnNext 队列以触发下一个线程。虽然消费是按顺序触发的,但它们仍然可以同时处理。
  • 是的,这正是我要说的,我最初认为 OnNexts 会并行触发,但事实并非如此。这是 Rx 的一个功能,它保证将保持集合的顺序,并且 OnNext for 2 不会在 OnNext for 1 完成之前触发。见akhildeshpande.com/2011/05/net-reactive-extensions-rx-2.html

标签: c# system.reactive


【解决方案1】:

NewThread 调度程序对于长时间运行的订阅者很有用。如果您不指定任何调度程序,则生产者将被阻止等待订阅者完成。通常,您可以使用 Scheduler.ThreadPool,但如果您希望有许多长时间运行的任务,您不会希望用它们阻塞线程池(因为它可能不仅仅被单个 observable 的订阅者使用)。

例如,考虑对您的示例进行以下修改。我将延迟移至订阅者,并添加了主线程何时准备好进行键盘输入的指示。请注意取消注释 NewThead 行时的区别。

using System;
using System.Linq;
using System.Reactive.Concurrency;
using System.Reactive.Linq;
using System.Threading;

namespace RxTesting
{
    class Program
    {
        static void Main(string[] args)
        {
            Console.WriteLine("Application Thread : {0}", Thread.CurrentThread.ManagedThreadId);

            var numbers = from number in Enumerable.Range(1, 10) select Process(number);

            var observableNumbers = numbers.ToObservable()
//              .ObserveOn(Scheduler.NewThread)
//              .SubscribeOn(Scheduler.NewThread)
            ;

            observableNumbers.Subscribe(
                n => {
                    Thread.Sleep(500);
                    Console.WriteLine("Consuming : {0} \t on Thread : {1}", n, Thread.CurrentThread.ManagedThreadId);
                });

            Console.WriteLine("Waiting for keyboard");
            Console.ReadKey();
        }

        private static int Process(int number)
        {
            Console.WriteLine("Producing : {0} \t on Thread : {1}", number,
                              Thread.CurrentThread.ManagedThreadId);

            return number;
        }
    }
}

那么为什么 Rx 不优化为每个订阅者使用相同的线程呢?如果订阅者运行时间太长以至于您需要一个新线程,那么线程创建开销无论如何都是微不足道的。一个例外是,如果大多数订阅者都很短,但少数订阅者长时间运行,那么重用同一线程的优化确实很有用。

【讨论】:

  • +1 我要补充一点,在线程池被其他操作大量加载的情况下,它允许比 ThreadPool 更可预测的执行。
  • @Edward:我不明白你的最后一点“如果订阅者运行时间太长以至于你需要一个新线程,那么线程创建开销将是微不足道的”?如果 Rx 保证按顺序运行每个订阅者,那么为每个订阅者创建一个新线程有什么意义?
  • @nabeelfarid:没有意义。但是,它可能需要 Rx 库中的额外代码才能在订阅者之间共享线程,所以我的猜测是,没有人会费心编写优化代码,假设它永远不会被注意到。
【解决方案2】:

我不确定您是否注意到,但如果消费者比生产者慢(例如,如果您在订阅操作中添加更长的睡眠),他们将共享同一个线程,因此它可能是一种确保订阅者在内容发布后立即使用。

namespace RxTesting
{
    class Program
    {
        static void Main(string[] args)
        {
            Console.WriteLine("Application Thread : {0}", Thread.CurrentThread.ManagedThreadId);

            var numbers = from number in Enumerable.Range(1,10) select Process(number);

            var observableNumbers = numbers.ToObservable()
                .ObserveOn(Scheduler.NewThread)
                .SubscribeOn(Scheduler.NewThread);

            observableNumbers.Subscribe(
                n => 
                    {
                        Console.WriteLine("Consuming : {0} \t on Thread : {1}", n, Thread.CurrentThread.ManagedThreadId);
                        Thread.Sleep(600);
                    }
                        );

            Console.ReadKey();
        }

        private static int Process(int number)
        {
            Thread.Sleep(500);
            Console.WriteLine("Producing : {0} \t on Thread : {1}", number,
                              Thread.CurrentThread.ManagedThreadId);

            return number;
        }
    }
}

输出:

Application Thread : 1
Producing : 1    on Thread : 3
Consuming : 1    on Thread : 4
Producing : 2    on Thread : 3
Consuming : 2    on Thread : 4
Producing : 3    on Thread : 3
Consuming : 3    on Thread : 4
Producing : 4    on Thread : 3
Consuming : 4    on Thread : 4
Producing : 5    on Thread : 3
Consuming : 5    on Thread : 4
Producing : 6    on Thread : 3
Consuming : 6    on Thread : 4
Producing : 7    on Thread : 3
Producing : 8    on Thread : 3
Consuming : 7    on Thread : 4
Producing : 9    on Thread : 3
Consuming : 8    on Thread : 4
Producing : 10   on Thread : 3
Consuming : 9    on Thread : 4
Consuming : 10   on Thread : 4

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-02
    • 2012-04-26
    相关资源
    最近更新 更多