【问题标题】:C# RabbitMQ wait for one message for specified timeout?C# RabbitMQ 等待一条消息指定超时?
【发布时间】:2017-05-22 10:11:03
【问题描述】:

RabbitMQ Wait for a message with a timeout 和 Wait for a single RabbitMQ message with a timeout 中的解决方案似乎不起作用,因为官方 C# 库中没有下一个交付方法,并且 QueueingBasicConsumer 已被弃用,所以它只是到处抛出 NotSupportedException。

如何等待队列中的单个消息达到指定的超时时间?

附言

可以通过 Basic.Get() 来完成,是的,但是,在指定的时间间隔(流量过多,CPU 过多)拉取消息是不好的解决方案。

更新

EventingBasicConsumer 通过实现不支持立即取消。即使您在某个时候调用 BasicCancel,即使您通过 BasicQos 指定预取 - 它仍然会在 Frames 中获取,并且这些帧可以包含多个消息。因此,它不适合单任务执行。不要打扰 - 它只是不适用于单个消息。

【问题讨论】:

    标签: c# rabbitmq


    【解决方案1】:

    有很多方法可以做到这一点。例如,您可以将EventingBasicConsumer 与ManualResetEvent 一起使用,如下所示(这仅用于演示目的 - 最好使用以下方法之一):

    var factory = new ConnectionFactory();
    using (var connection = factory.CreateConnection()) {
        using (var channel = connection.CreateModel()) {
            // setup signal
            using (var signal = new ManualResetEvent(false)) {
                var consumer = new EventingBasicConsumer(channel);
                byte[] messageBody = null;                        
                consumer.Received += (sender, args) => {
                    messageBody = args.Body;
                    // process your message or store for later
                    // set signal
                    signal.Set();
                };               
                // start consuming
                channel.BasicConsume("your.queue", false, consumer);
                // wait until message is received or timeout reached
                bool timeout = !signal.WaitOne(TimeSpan.FromSeconds(10));
                // cancel subscription
                channel.BasicCancel(consumer.ConsumerTag);
                if (timeout) {
                    // timeout reached - do what you need in this case
                    throw new Exception("timeout");
                }
    
                // at this point messageBody is received
            }
        }
    }
    

    正如您在 cmets 中所述 - 如果您希望在同一个队列中有多个消息,这不是最好的方法。好吧,无论如何这都不是最好的方法,我包含它只是为了演示ManualResetEvent 的使用,以防库本身不提供超时支持。

    如果您正在执行 RPC(远程过程调用,请求-回复) - 您可以在服务器端使用 SimpleRpcClient 和 SimpleRpcServer。客户端将如下所示:

    var client = new SimpleRpcClient(channel, "your.queue");
    client.TimeoutMilliseconds = 10 * 1000;
    client.TimedOut += (sender, args) => {
        // do something on timeout
    };                    
    var reply = client.Call(myMessage); // will return reply or null if timeout reached
    

    更简单的方法:使用基本的Subscription 类(它在内部使用相同的EventingBasicConsumer,但支持超时,因此您无需自己实现),如下所示:

    var sub = new Subscription(channel, "your.queue");
    BasicDeliverEventArgs reply;
    if (!sub.Next(10 * 1000, out reply)) {
         // timeout
    }
    

    【讨论】:

    • 第一个解决方案无效。 BasicConsume 不能保证在 BasicCancel 上停止消费,由于 rabbit 的动态实现,它可以稍后执行此操作(只需尝试每个请求使用一条消息,您会看到在某些情况下您分配 messageBody 几次)。您仍然需要重新排队多余的消息。第二个和第三个,我现在试试=)
    • 虽然你的上半部分无关紧要,但订阅类正是我想要的!谢谢,效果很好!您可以编辑答案,以便其他人知道最后一个解决了吗?
    • 不过,它通过实现在内部存储了一堆消息,而我只需要一个 =/
    • 为什么需要超时获取?已发送请求并等待回复,还是出于其他原因?
    • 您可以使用basicQos 将预取计数限制为1,然后禁用自动确认并手动确认您的消息吗?然后您将收到一条一条消息,并且兔子不会将更多消息推送到您的频道,除非您确认当前一条(我认为)。
    猜你喜欢
    • 2011-02-17
    • 1970-01-01
    • 2015-04-15
    • 1970-01-01
    • 2022-11-03
    • 1970-01-01
    • 1970-01-01
    • 2013-12-07
    • 2021-01-31
    相关资源
    最近更新 更多