【发布时间】: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