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