【问题标题】:What is an elegant way of halting consumption of messages gracefully in the C# client for RabbitMQ?在 RabbitMQ 的 C# 客户端中优雅地停止使用消息的优雅方法是什么?
【发布时间】:2014-09-29 17:07:30
【问题描述】:

我正在用 C# 设置一个标准的独立线程来监听 RabbitMQ。假设线程中监听的方法是这样的:

public void Listen()
{
    using (var channel = connection.CreateModel())
    {
        var consumer = SetupQueues(channel);
        while (true)
        {
            var ea = consumer.Queue.Dequeue();    // blocking call
            handler.HandleMessage(channel, ea);
        }
    }
}

在 RabbitMQ 的 C# 客户端中优雅地停止消费消息的优雅方式是什么?请记住,我在 RabbitMQ 示例/文档或这些 SO 问题中没有发现任何用处:

这里的问题是consumer.Queue.Dequeue() 是一个阻塞调用。我已经尝试了这些选项:

  • 呼叫channel.BasicCancel(string tag)。这会导致阻塞调用中出现System.IO.EndOfStreamException。出于显而易见的原因,我不想将此异常用作控制流的一部分。

  • 调用consumer.Queue.Dequeue(int millisecondsTimeout, out T result) 并检查循环迭代之间的标志。这可以工作,但看起来很hacky。

我想让线程优雅地退出并清理我可能拥有的所有非托管资源,因此不会出现线程中止等情况。

感谢任何帮助。谢谢

【问题讨论】:

  • 据我所知,没有“干净”的方法可以做到这一点。我个人使用超时方法——我的消费者线程检查一个标志,并据此决定做什么。这实际上似乎工作正常。本质上,应用程序的状态决定了队列是否应该被消费——消费者循环只是在它尝试消费之前检查它是否应该消费。

标签: c# rabbitmq messaging amqp consumer


【解决方案1】:

带有超时和标志的 DeQueue 就是这样做的方法。这是一种非常常见的模式,这也是为什么许多阻塞调用都提供了启用超时的版本。

另外,抛出(已知的)异常对于控制流来说不一定是坏事。优雅地关闭可能意味着实际捕获异常,评论“这是在请求关闭通道时抛出的”,然后干净地返回。这就是 TPL 的一部分与 CancellationToken 一起工作的方式。

【讨论】:

  • 明白,感谢您的澄清。就使用异常而言,我很乐意使用它,但是,当通过其他方式关闭“SharedQueue”(网络故障、心跳丢失等)时,会引发相同的异常(System.IO.EndOfStreamException: SharedQueue closed)跨度>
  • 是的,这很不幸。 TPL 以同样的方式执行此操作,它们抛出一个 ThreadAbortException,这可能由于多种原因而发生。我一直不明白为什么——开发人员创建异常很便宜,为什么他们不使用更多的 OO 呢? :(
【解决方案2】:

阻塞方法不是属性事件驱动的。 我不明白他们为什么建议使用consumer.Queue.Dequeue();

反正我一般不会用consumer.Queue.Dequeue();

我以这种方式扩展了默认消费者:

class MyConsumer : DefaultBasicConsumer {

  public MyConsumer(IModel model):base(model)
  {

  }
  public override void HandleBasicDeliver(string consumerTag, ulong deliveryTag, bool redelivered, string exchange, string routingKey, IBasicProperties properties, byte[] body) {
    var message = Encoding.UTF8.GetString(body);
    Console.WriteLine(" [x] Received {0}", message);
  }
}


class Program
{
  static void Main(string[] args)
  {

    var factory = new ConnectionFactory() { Uri = "amqp://aa:bbb@lemur.cloudamqp.com/xxx" };
    using (var connection = factory.CreateConnection())
    {
      using (var channel = connection.CreateModel())
      {
      channel.QueueDeclare("hello", false, false, false, null);
      var consumer = new MyConsumer(channel);
      String tag = channel.BasicConsume("hello", true, consumer);
      Console.WriteLine(" [*] Waiting for messages." +
                               " any key to exit");
      Console.ReadLine();
      channel.BasicCancel(tag);


        /*while (true)
        {
          /////// DON'T USE THIS   
          var ea = (BasicDeliverEventArgs)consumer.Queue.Dequeue();
          var body = ea.Body;
          var message = Encoding.UTF8.GetString(body);
          Console.WriteLine(" [x] Received {0}", message);
        }*/
      }
    }


  }
}

这样你就没有阻塞的方法了,你可以正确的释放所有的资源。

编辑

我认为使用 ctrl+C 来破坏程序总是错误的。

【讨论】:

    猜你喜欢
    • 2019-09-23
    • 2010-11-05
    • 1970-01-01
    • 1970-01-01
    • 2010-12-15
    • 1970-01-01
    • 2018-09-27
    • 2019-11-10
    • 1970-01-01
    相关资源
    最近更新 更多