【问题标题】:RabbitMQ Manual ACK on c# client在 C# 客户端上的 RabbitMQ 手动 ACK
【发布时间】:2016-05-18 18:25:23
【问题描述】:

我正在尝试在一个非常简单的控制台应用程序上使用手动 ACK,但我无法使其工作。

在发件人上,我有以下代码:

var factory = new ConnectionFactory() { HostName = "localhost" };
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
    channel.QueueDeclare(queue: "task_queue",
                         durable: true,
                         exclusive: false,
                         autoDelete: false,
                         arguments: null);

    var message = GetMessage(args);
    var body = Encoding.UTF8.GetBytes(message);

    channel.ConfirmSelect();
    channel.BasicAcks += (sender, e) =>
    {
        Console.Write("ACK received");
    };

    var properties = channel.CreateBasicProperties();

    channel.BasicPublish(exchange: "",
                         routingKey: "task_queue",
                         basicProperties: properties,
                         body: body);

    Console.WriteLine(" [x] Sent {0}", message);
}

Console.WriteLine(" Press [enter] to exit.");
Console.ReadLine();

在接收器上我有以下代码:

var factory = new ConnectionFactory() { HostName = "localhost" };
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
    channel.QueueDeclare(queue: "task_queue",
                         durable: true,
                         exclusive: false,
                         autoDelete: false,
                         arguments: null);

    channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);
    channel.ConfirmSelect();

    Console.WriteLine(" [*] Waiting for messages.");

    var consumer = new EventingBasicConsumer(channel);
    consumer.Received += (model, ea) =>
    {
        var body = ea.Body;
        var message = Encoding.UTF8.GetString(body);
        Console.WriteLine(" [x] Received {0}", message);

        int dots = message.Split('.').Length - 1;
        Thread.Sleep(dots * 1000);

        channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
        Console.WriteLine(" [x] Done");
    };
    channel.BasicConsume(queue: "task_queue",
                         noAck: false,
                         consumer: consumer);

    Console.WriteLine(" Press [enter] to exit.");
    Console.ReadLine();
}

我期望的是,当我在接收者上调用 channel.BasicAck() 时,发送者上的事件 BasicAcks 会被触发,但是当消息传递到客户端时,会在 consumer.Received 之前触发该事件。

我所期望的行为是正确的还是我遗漏了什么?

【问题讨论】:

    标签: c# rabbitmq


    【解决方案1】:

    您的期望不正确。 BasicAcks 是关于 publisher confirms,而不是来自接收者的 ack。因此,您向代理发布消息,broker(因此,RabbitMQ 本身)在处理此消息时(例如,当它将其写入磁盘以获取持久消息时)将确认或确认(否定确认)您, 或当 in 时将其放入队列中)。请注意,这里不涉及接收者 - 它完全在发布者和 RabbitMQ 之间。

    现在,当您在接收者处确认消息时 - 再次仅在接收者和 RabbitMQ 之间进行 - 您告诉 rabbit 消息已被处理并且可以安全地删除。这样做是为了处理接收器在处理过程中崩溃的情况 - 然后 rabbit 将能够将此消息传递给下一个接收器(如果有的话)。

    请注意,此类架构的全部目的是将发布者和接收者分开 - 它们不应相互依赖。

    如果您有一个接收器(可能有很多)并且您希望确保它处理您的消息 - 使用 RPC 模式:发送消息并等待来自该接收器的另一条消息。

    【讨论】:

      【解决方案2】:

      消费者:

      consumer.Received += async (model, ea) =>
      {
          var body = ea.Body.ToArray();
          var message = Encoding.UTF8.GetString(body);
          Console.WriteLine(" [x] Received {0}", message);
      
          int dots = message.Split('.').Length - 1;
          await Task.Delay(2000);
      
          channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
          Console.WriteLine(" [x] Done");
      };
      

      【讨论】:

        猜你喜欢
        • 2014-04-27
        • 2023-04-05
        • 1970-01-01
        • 1970-01-01
        • 2019-04-28
        • 1970-01-01
        • 1970-01-01
        • 2020-05-05
        • 2012-08-14
        相关资源
        最近更新 更多