【发布时间】:2017-04-14 21:09:42
【问题描述】:
我正在使用此处发布的示例消费者:
http://pika.readthedocs.org/en/latest/examples/asynchronous_consumer_example.html
我使用 ExampleConsumer 的原因是当工作任务开始花费更长的时间时,我与 rabbitmq 的连接失败,其中时间超过 10 分钟。在长时间运行的任务完成并且进程失败后,连接被关闭。之前,它需要一分钟左右的时间来处理 1000 条消息。
ExampleConsumer 似乎可以重新连接,但是,在确认消息中,该消息实际上并未得到确认,因为连接已失效。它似乎从下面的确认消息方法正常返回。然后它会尝试重新连接,然后重新发送刚刚完成的消息。
def acknowledge_message(self, delivery_tag):
"""Acknowledge the message delivery from RabbitMQ by sending a
Basic.Ack RPC method for the delivery tag.
:param int delivery_tag: The delivery tag from the Basic.Deliver frame
"""
LOGGER.info('Acknowledging message %s', delivery_tag)
self._channel.basic_ack(delivery_tag)
【问题讨论】: