【问题标题】:Reactive Framework as Message queue using BlockingCollection使用 BlockingCollection 作为消息队列的响应式框架
【发布时间】:2013-04-02 19:45:13
【问题描述】:

我最近一直在使用响应式框架做一些工作,到目前为止我一直非常喜欢它。我正在考虑用一些过滤的 IObservable 替换传统的轮询消息队列,以清理我的服务器操作。以旧的方式,我处理进入服务器的消息是这样的:

// Start spinning the process message loop
   Task.Factory.StartNew(() =>
   {
       while (true)
       {
           Command command = m_CommandQueue.Take();
           ProcessMessage(command);
       }
   }, TaskCreationOptions.LongRunning);

这会导致持续轮询线程,将来自客户端的命令委托给 ProcessMessage 方法,在该方法中,我有一系列 if/else-if 语句来确定命令的类型并根据其类型委托工作

我正在用一个使用 Reactive 的事件驱动系统替换它,为此我编写了以下代码:

 private BlockingCollection<BesiegedMessage> m_MessageQueue = new BlockingCollection<BesiegedMessage>();
 private IObservable<BesiegedMessage> m_MessagePublisher;

 m_MessagePublisher = m_MessageQueue
       .GetConsumingEnumerable()
       .ToObservable(TaskPoolScheduler.Default);

        // All generic Server messages (containing no properties) will be processed here
 IDisposable genericServerMessageSubscriber = m_MessagePublisher
       .Where(message => message is GenericServerMessage)
       .Subscribe(message =>
       {
           // do something with the generic server message here
       }

我的问题是,虽然这可行,但使用阻塞集合作为这样的 IObservable 的支持是一种好习惯吗?我看不到 Take() 在哪里被这样调用,这让我认为消息会堆积在队列中而不会在处理后被删除?

将主题作为后备集合来驱动将接收这些消息的过滤后的 IObservable 会更有效吗?还有什么我在这里遗漏的可能有利于这个系统的架构的东西吗?

【问题讨论】:

  • +1 用于使用 Rx :) - 你的命令队列是什么 - 消息本质上是什么?它是临时字符 - 还是应用程序范围内可用的东西 - 所以一部分提交 - 另一部分订阅?或者,即如果两个订阅者“订阅”(还有什么),他们会在发生时看到相同的事情,还是每个订阅者都有自己的? - 与 Rx 打交道时,您需要回答那些简单的事情——性质、用途。如果有帮助就快点扔到这里
  • @NSGaga 消息是整个应用程序的基础,因为这里有两种方式进行通信,从客户端到服务器,反之亦然。这些消息只是一种共同的语言,双方都可以采取行动
  • 我脑海中闪现的第一件事是,您不再特别需要 BlockingCollection - 当您主动轮询新消息时,您需要这些,但是,如果您的消息已发送给您,则无需再“阻止直到我收到信号”
  • @JerKimball 我也想到了这个想法。在我的系统中,所有消息都通过客户端公开的称为 SendMessage(string message) 的方法发送到服务器。有没有办法从这里传递的所有消息中创建一个 IObservable 源,而不是将它们放入 BlockingCollection?
  • 为了记录,消息在处理时肯定会从 BlockingCollection 队列中删除。再次感谢出色的示例代码!

标签: c# system.reactive


【解决方案1】:

这是一个完整的工作示例,在 Visual Studio 2012 下测试。

  1. 创建一个新的 C# 控制台应用程序。
  2. 右键单击您的项目,选择“管理 NuGet 包”,然后添加“反应式扩展 - 主 图书馆”。

添加此 C# 代码:

using System;
using System.Collections.Concurrent;
using System.Reactive.Concurrency;
using System.Reactive.Linq;
namespace DemoRX
{
    class Program
    {
        static void Main(string[] args)
        {
            BlockingCollection<string> myQueue = new BlockingCollection<string>();
            {                
                IObservable<string> ob = myQueue.
                  GetConsumingEnumerable().
                  ToObservable(TaskPoolScheduler.Default);

                ob.Subscribe(p =>
                {
                    // This handler will get called whenever 
                    // anything appears on myQueue in the future.
                    Console.Write("Consuming: {0}\n",p);                    
                });
            }
            // Now, adding items to myQueue will trigger the item to be consumed
            // in the predefined handler.
            myQueue.Add("a");
            myQueue.Add("b");
            myQueue.Add("c");           
            Console.Write("[any key to exit]\n");
            Console.ReadKey();
        }
    }
}

你会在控制台上看到这个:

[any key to exit]
Consuming: a
Consuming: b
Consuming: c

使用 RX 的真正好处是您可以使用 LINQ 的全部功能来过滤掉任何不需要的消息。例如,添加一个.Where 子句来过滤“a”,然后观察会发生什么:

ob.Where(o => (o == "a")).Subscribe(p =>
{
    // This will get called whenever something appears on myQueue.
    Console.Write("Consuming: {0}\n",p);                    
});

哲学笔记

与启动专用线程来轮询队列相比,此方法的优势在于,您不必担心程序退出后会正确处理线程。这意味着您不必为 IDisposable 或 CancellationToken 操心(在处理 BlockingCollection 时总是需要这样做,否则您的程序可能会在退出时挂起,线程拒绝终止)。

相信我,编写完全健壮的代码来处理来自 BlockingCollection 的事件并不像您想象的那么容易。我更喜欢使用 RX 方法,如上所示,因为它更简洁、更健壮、代码更少,并且您可以使用 LINQ 进行过滤。

延迟

我对这种方法的速度感到惊讶。

在我的 Xeon X5650 @ 2.67Ghz 上,处理 1000 万个事件需要 5 秒,每个事件大约需要 0.5 微秒。将项目放入 BlockingCollection 需要 4.5 秒,因此 RX 将它们取出并处理它们的速度几乎与它们进入时一样快。

线程

在我所有的测试中,RX 只启动了一个线程来处理队列中的任务。

这意味着我们有一个非常好的模式:我们可以使用 RX 收集来自多个线程的传入数据,将它们放入共享队列中,然后在单个线程上处理队列内容(根据定义,这是线程安全的)。

这种模式消除了处理多线程代码时的大量麻烦,通过队列将数据的生产者和消费者解耦,其中生产者可以是多线程的,消费者是单线程的,因此是线程安全的。这就是使 Erlang 如此健壮的概念。有关此模式的更多信息,请参阅Multi-threading made ridiculously simple。

【讨论】:

  • 真的不知道这个答案的目的是什么。这个问题非常古老,更不用说您只是重复了我发布的代码 sn-p 。它没有增加任何建设性。
  • 更多供我自己参考,因为您最初给出的示例不完整。我的学习曲线与您一年前的学习曲线相同。
  • 同意杰西。你在这里所做的很好,但我认为你可以用一个主题和.Synchonize() 替换所有这些,以允许多个线程调用OnNext()。这就是我在回答中的观点,也许它没有很好地表达出来。基本上,队列需要一个线程专门用于读取它。当进程退出时,后台线程将关闭。 Linq 在 IEnumerable 上也很有效。
  • 考虑到这一点,阻塞队列的间接性确实将订阅者的 OnNext 处理程序的缓慢与阻塞生产者排队的能力分开。然而,这很容易通过ObserveOn(IScheduler) 解决,它还重新引入了与上面提供的答案相同的并发性。所以 queue.ToObservable 代码也可以是subject.Synchronize().ObserveOn(TaskPoolScheduler.Default),不是吗?可能会删除Synchronize(),因为ObserveOn 将使用其内部队列序列化数据。所以我认为 subject.ObserveOn(TaskPoolScheduler.Default)
【解决方案2】:

这是直接从我的后部提取的东西 - 任何真正的解决方案都非常很大程度上取决于您的实际使用情况,但这里是“有史以来最便宜的伪消息队列系统”:

想法/动机:

  • 故意曝光IObservable&lt;T&gt;,以便订阅者可以进行任何他们想要的过滤/交叉订阅
  • 整个 Queue 是无类型的,但 Register 和 Publish 是类型安全的(ish)
  • 带有Publish() 的YMMV - 尝试移动它
  • 通常Subject 是一个禁忌,尽管在这种情况下它确实会产生一些简单的代码。
  • 也可以将注册“内部化”以实际执行订阅,但随后队列将需要管理创建的 IDisposables - 嗯,让您的消费者处理它!

代码:

public class TheCheapestPubSubEver
{    
    private Subject<object> _inner = new Subject<object>();

    public IObservable<T> Register<T>()
    {
        return _inner.OfType<T>().Publish().RefCount();
    }
    public void Publish<T>(T message)
    {
        _inner.OnNext(message);
    }
}

用法:

void Main()
{
    var queue = new TheCheapestPubSubEver();

    var ofString = queue.Register<string>();
    var ofInt = queue.Register<int>();

    using(ofInt.Subscribe(i => Console.WriteLine("An int! {0}", i)))
    using(ofString.Subscribe(s => Console.WriteLine("A string! {0}", s)))
    {
        queue.Publish("Foo");
        queue.Publish(1);
        Console.ReadLine();
    }
}

输出:

A string! Foo
An int! 1

但是,这并没有严格执行“消费消费者” - 特定类型的多个寄存器会导致多个观察者调用 - 即:

var queue = new TheCheapestPubSubEver();

var ofString = queue.Register<string>();
var anotherOfString = queue.Register<string>();
var ofInt = queue.Register<int>();

using(ofInt.Subscribe(i => Console.WriteLine("An int! {0}", i)))
using(ofString.Subscribe(s => Console.WriteLine("A string! {0}", s)))
using(anotherOfString.Subscribe(s => Console.WriteLine("Another string! {0}", s)))

{
    queue.Publish("Foo");
    queue.Publish(1);
    Console.ReadLine();
}

结果:

A string! Foo
Another string! Foo
An int! 1

【讨论】:

  • 是 - 多个Publish-es - 表示多个提要。您可以将所有这些(以及字符串/int 组合)合并为一个 Publish。只是为了 OP-s 的清晰度。
【解决方案3】:

我没有在这种情况下使用BlockingCollection - 所以我在“推测” - 你应该运行它来批准,反驳。

BlockingCollection 可能只会使这里的事情变得更加复杂(或提供很少的帮助)。看看this post from Jon - 只是为了确认。 GetConsumingEnumerable 将为您提供可枚举的“每个订阅者”。最终使他们筋疲力尽 - Rx 需要考虑的事情。

IEnumerable&lt;&gt;.ToObservable 也进一步扁平化了“来源”。当它起作用时(你可以查找源——我更推荐使用 Rx)——每个订阅者都会创建一个自己的“枚举器”——所以所有人都会得到他们自己版本的提要。我真的不确定,在像这样的 Observable 场景中是如何实现的。

无论如何 - 如果您想提供应用程序范围的消息 - IMO 您需要引入 Subject 或以其他形式(例如发布等)声明。从这个意义上说,我认为 BlockingCollection 不会有任何帮助 - 但同样,您最好自己尝试一下。

注意(一个哲学的)

如果您想组合消息类型,或组合不同的来源 - 例如在更“真实世界”的场景中 - 它变得更加复杂。我必须说它变得非常有趣。

注意将它们“扎根”到单一共享流中(并避免 Jer 正确建议的内容)。

我建议您不要尝试使用Subject 来逃避。对于你所需要的,那是你的朋友——不管所有与无状态相关的讨论以及主题有多糟糕——你实际上都有一个状态(并且你需要一个“状态”)——Rx 在“事后”开始,所以你无论如何享受它的好处。

我鼓励你这样做,因为我喜欢它的结果。

【讨论】:

    【解决方案4】:

    我的问题是,我们已经将队列(我通常将它与一个消费者的破坏性读取联系起来,特别是如果您使用 BlockingCollection 时)变成了广播(发送给现在正在收听的任何人和每个人)。

    这似乎是两个相互矛盾的想法。

    我已经看到这样做了,但后来被丢弃了,因为它是“对错误问题的正确解决方案”。

    【讨论】:

    • 在上面的示例代码中,您根本无法将多个订阅添加到单个 BlockingCollection,因为每个订阅者都会对 BlockingCollection 进行破坏性读取。上面的示例代码只是干净地处理将来任何时候添加到队列中的项目的好方法。
    • 这与不使用 Rx 而只拥有一个专用于从队列中消费的线程有什么不同?虽然我喜欢 Rx,但我认为在这种情况下它不会为自己付出代价。
    猜你喜欢
    • 2013-12-16
    • 2014-04-09
    • 1970-01-01
    • 1970-01-01
    • 2014-02-13
    • 2018-10-12
    • 1970-01-01
    • 1970-01-01
    • 2020-11-27
    相关资源
    最近更新 更多