消息响应-避免消息丢失

当处理一个比较耗时的任务的时候,也许想知道消费者(consumers)是否运行到一半就挂掉。在当前的代码中,当RabbitMQ将消息发送给消费者(consumers)之后,马上就会将该消息从队列中移除。此时,如果把处理这个消息的工作者(worker)停掉,正在处理的这条消息就会丢失。同时,所有发送到这个工作者的还没有处理的消息都会丢失。

我们不想丢失任何任务消息。如果一个工作者(worker)挂掉了,我们希望该消息会重新发送给其他的工作者(worker)。

为了防止消息丢失,RabbitMQ提供了消息响应(acknowledgments)机制。消费者会通过一个ack(响应),告诉RabbitMQ已经收到并处理了某条消息,然后RabbitMQ才会释放并删除这条消息。

如果消费者(consumer)挂掉了,没有发送响应,RabbitMQ就会认为消息没有被完全处理,然后重新发送给其他消费者(consumer)。这样,即使工作者(workers)偶尔挂掉,也不会丢失消息。

消息是没有超时这个概念的;当工作者与RabbitMQ断开连接的时候,RabbitMQ会重新发送消息。这样在处理一个耗时非常长的消息任务的时候就不会出问题了。

消息响应默认是开启的。在之前的例子中使用了no_ack=True标识把它关闭。是时候移除这个标识了,当工作者(worker)完成了任务,就发送一个响应。

channel.BasicConsume("hello", false, consumer);Code language: JavaScript (javascript)

如果不关闭,会报下面异常:

RabbitMQ.Client.Exceptions.AlreadyClosedException:"Already closed: The AMQP operation was interrupted: AMQP close-reason, 
initiated by Peer, code=406, text='PRECONDITION_FAILED - unknown delivery tag 1', classId=60, methodId=80"Code language: PHP (php)

现在,可以保证,即使正在处理消息的工作者被停掉,这些消息也不会丢失,所有没有被应答的消息会被重新发送给其他工作者。

一个很常见的错误就是忘掉了BasicAck这个方法,这个错误很常见,但是后果很严重。当客户端退出时,待处理的消息就会被重新分发,但是RabbitMQ会消耗越来越多的内存,因为这些没有被应答的消息不能够被释放。调试这种case,可以使用rabbitmqctl打印messages_unacknowledged字段。

while (true)
{
    var ea = (BasicDeliverEventArgs)consumer.Queue.Dequeue();

    var body = ea.Body;
    var message = Encoding.UTF8.GetString(body);

    var time = Convert.ToInt32(message.Split(':')[1]);
    Thread.Sleep(time * 1000);

    Console.WriteLine("Received {0}", message);
    Console.WriteLine("Done");

    channel.BasicAck(ea.DeliveryTag, false);//添加响应
}Code language: JavaScript (javascript)

现在我们先启动消费者,然后启动生产者。当消费者处理一部分的时候,关闭消费者,然后再重新打开。

我们可以看到,消费者会继续把后面的消息消费完,不会丢失。

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注