【问题标题】:How to create an IObservable<T> which reads from a MSMQ message queue?如何创建从 MSMQ 消息队列中读取的 IObservable<T>?
【发布时间】:2012-02-10 01:52:04
【问题描述】:

我正在从我们的 ASP.NET 站点中删除我们的电子邮件系统,该站点用于立即使用系统发送电子邮件,以在单独的服务中处理请求以减少网站上的工作量。我正在尝试围绕一组接口设计它,以便我可以根据需要交换实现,但最初它将基于消息队列(MSMQ)将请求发送到队列,让服务接收传入请求然后处理它们。我目前大致定义了以下接口:

// Sends one or more requests to be processed somehow
public interface IRequestSender
{
    void Send(IEnumerable<Request> requests);
}

// Listens for incoming requests and passes them to an observer to do the real work
public interface IRequestListener : IObservable<Request>
{
    void Start();
    void Stop();
}

// Processes a request given to it by a IRequestListener
public interface IRequestProcessor : IObserver<Request>
{
}

您会注意到 Listener 和 Processor 使用 observable 模式,因为我认为这似乎最合适。

我的问题是弄清楚如何编写从 MSMQ 接收的IRequestListener 的实现,基本上我如何创建合适的IObservable&lt;T&gt;?

我发现的第一个选择是根据MSDN documentation 给出的示例从头开始创建IObservable&lt;T&gt;,但这似乎需要做很多管道工作。

另一个选择是使用响应式扩展,因为它似乎旨在使创建可观察对象变得更容易。我发现最接近将 Rx 与 MSMQ 结合使用的是这些页面:

但我不确定如何将这些示例应用到我的IRequestListener 界面。

也欢迎任何其他想法,如果合适的话,甚至可以更改我的基本设计。

【问题讨论】:

    标签: .net msmq system.reactive


    【解决方案1】:

    我最初确实使用了 FromAsyncPattern,但后来为它编写了一个类,因为它可以更好地处理超时和中毒消息。一旦开始,队列无论如何都是热门的 Observables。您也可以使用 Observable.Defer 使其更接近 Rx 而不是 Start/Stop。

    这是 QueueObservable 的基本实现。您可以先致电ListenReceive。

    Subject<T> Subject = new Subject<T>();
    
    protected void ListenReceive()
    {
        Queue.BeginReceive(MessageQueue.InfiniteTimeout, null, OnReceive);
    }
    
    protected void OnReceive(IAsyncResult ar)
    {
        Message message = null;
    
        try
        {
            message = Queue.EndReceive(ar);
        }
        catch (TimeoutException ex)
        {
            //retry?
        }
    
        if (message != null)
            Subject.OnNext((T) message.Body);
    
        Thread.Yield();
    
        if (!IsDisposed)
            ListenReceive();
    }    
    
    public IObservable<T> AsObservable()
    {
            return Subject;
    }
    

    【讨论】:

    • 在看到这个答案之前,我正在按照与此相同的方式进行一些实验,在内部使用 Subject&lt;T&gt; 确实有助于跟踪订阅,我的实施与您的建议相距一百万英里,谢谢。
    猜你喜欢
    • 2014-05-14
    • 2013-01-27
    • 2011-06-16
    • 2013-12-31
    • 2017-06-16
    • 2010-12-10
    • 2010-09-11
    • 1970-01-01
    • 2018-02-20
    相关资源
    最近更新 更多