【问题标题】:Reactive Extensions: Concurrency within the subscriber反应式扩展:订阅者内的并发
【发布时间】:2013-05-15 13:34:34
【问题描述】:

我正试图围绕 Reactive Extensions 对并发的支持进行思考,并且很难获得我想要的结果。所以我可能还没有得到它

我有一个源,它向流中发送数据的速度比订阅者消耗它的速度要快。我更喜欢配置流,以便使用另一个线程为流中的每个新项目调用订阅者,以便订阅者有多个线程同时通过它运行。我能够确保订阅者的线程安全。

以下示例演示了该问题:

Observable.Interval( TimeSpan.FromSeconds(1))
    .Do( x => Console.WriteLine("{0} Thread: {1} Source value: {2}",
                                DateTime.Now, 
                                Thread.CurrentThread.ManagedThreadId, x))
    .ObserveOn(NewThreadScheduler.Default)
    .Subscribe(x =>
               {
                   Console.WriteLine("{0} Thread: {1} Observed value: {2}",
                                     DateTime.Now,
                                     Thread.CurrentThread.ManagedThreadId, x);
                   Thread.Sleep(5000); // Simulate long work time
               });

控制台输出如下所示(删除日期):

4:25:20 PM Thread: 6 Source value: 0
4:25:20 PM Thread: 11 Observed value: 0
4:25:21 PM Thread: 12 Source value: 1
4:25:22 PM Thread: 12 Source value: 2
4:25:23 PM Thread: 6 Source value: 3
4:25:24 PM Thread: 6 Source value: 4
4:25:25 PM Thread: 11 Observed value: 1
4:25:25 PM Thread: 12 Source value: 5
4:25:26 PM Thread: 6 Source value: 6

请注意“观察值”时间增量。即使源继续发出数据的速度快于订阅者处理它的速度,订阅者也不会被并行调用。虽然我可以想象当前行为会派上用场的大量场景,但我需要能够在消息可用时立即对其进行处理。

我已经使用 ObserveOn 方法尝试了几种调度程序的变体,但它们似乎都没有做我想要的。

除了在订阅操作中分离一个线程以执行长时间运行的工作之外,我是否还缺少任何可以将数据并发传递给订阅者的东西?

提前感谢所有答案和建议!

【问题讨论】:

  • 自从我上次使用RX已经有一段时间了,但是避免手动线程处理不是更好吗?也就是说,使用 TPL 在 Subscribe 方法中生成一个后台任务,然后可以立即返回。据我所知,这将解决您的并发问题并避免产生大量线程的风险(如果您的源比您的订阅者更快,则会发生这种情况)。

标签: c# concurrency system.reactive


【解决方案1】:

这里的根本问题是,您希望 Rx observable 以一种真正违反 observable 工作规则的方式调度事件。我认为在这里查看 Rx 设计指南会很有启发性:http://go.microsoft.com/fwlink/?LinkID=205219 - 最值得注意的是,“4.2 假设观察者实例以序列化方式调用”。即您不应该能够并行运行 OnNext 调用。事实上,Rx 的排序行为是其设计理念的核心。

如果您查看源代码,您会发现 Rx 在派生 ObserveOnObserver<T>ScheduledObserver<T> 类中强制执行此行为... OnNexts 从内部队列调度,并且每个必须在下一个之前完成被调度 - 在给定的执行上下文中。 Rx 不允许单个订阅者的 OnNext 调用同时执行。

这并不是说您不能让多个订阅者以不同的速率执行。实际上,如果您将代码更改如下,这很容易看出:

var source = Observable.Interval(TimeSpan.FromSeconds(1))
    .Do(x => Console.WriteLine("{0} Thread: {1} Source value: {2}",
                                DateTime.Now,
                                Thread.CurrentThread.ManagedThreadId, x))
    .ObserveOn(NewThreadScheduler.Default);

var subscription1 = source.Subscribe(x =>
    {
        Console.WriteLine("Subscriber 1: {0} Thread: {1} Observed value: {2}",
                            DateTime.Now,
                            Thread.CurrentThread.ManagedThreadId, x);
        Thread.Sleep(1000); // Simulate long work time
    });

var subscription2 = source.Subscribe(x =>
{
    Console.WriteLine("Subscriber 2: {0} Thread: {1} Observed value: {2}",
                        DateTime.Now,
                        Thread.CurrentThread.ManagedThreadId, x);
    Thread.Sleep(5000); // Simulate long work time
});

现在您会看到订阅者 1 领先于订阅者 2。

你不能轻易做到的是让 observable 做一些事情,比如向“准备好的”订阅者发送 OnNext 调用——这是你以迂回的方式要求的。我还假设您不会真的想在消费缓慢的情况下为每个 OnNext 创建一个新线程!

在这种情况下,听起来您可能最好使用一个订阅者,该订阅者除了尽快将工作推送到队列中之外什么都不做,而队列又由许多消费工作线程提供服务,然后您可以控制为有必要跟上步伐。

【讨论】:

  • 谢谢詹姆斯。你的解释很有道理。感谢您提供指向设计指南的指针……它涵盖了我在书籍和博客中没有看到的内容。
  • 欢迎您!顺便说一句,我写了关于删除事件的有效方法,以便让消费者在此处保持最新状态:zerobugbuild.com/?p=192 它可能与某些场景相关。
  • 您还可以执行source.Select(x => Observable.Defer(() => HandleAsync(x))).Merge(5).Subscribe(...) 之类的操作,其中HandleAsync(x) 是一种处理该项目并在完成时返回任务以发出信号的方法。当您的来源有时比您的订阅者可以处理的更快时,这种模式很好。以上将最多同时观察5个结果。仅供参考在订阅调用中,您实际上是在观察任务的结果。
  • 如果你想“扇出”你可以使用旧的 ..SelectMany(i=>Observable.Start(()=>DoSomething(i))..
  • 我偶然发现了这个问答,但我正在寻找 @JamesWorld 所写的内容。这是一个使用非物质化流的修改版本的要点 - gist.github.com/aniongithub/c650636189b68d0b1ece3e020bc51329
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-04
  • 2011-11-12
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多