如何使用Pika向RabbitMQ发布持久化消息?解决容器重启消息丢失问题
如何通过Pika向RabbitMQ发布持久化消息?
你已经将队列和交换器声明为持久化,但重启容器后消息丢失的核心原因是发送的消息未标记为持久化。RabbitMQ中,只有同时满足以下三个条件,消息才能在服务重启后保留:
- 交换器是持久化的(你已设置
durable=True) - 队列是持久化的(你已设置
durable=True) - 消息本身标记为持久化
修改发送消息的代码
在调用basic_publish时,需要通过properties参数指定消息的持久化属性,具体是设置delivery_mode=2(RabbitMQ中2代表持久化消息)。
修改后的完整代码示例:
import json import pika from pika.spec import BasicProperties # 假设你已完成连接建立和通道初始化 channel.queue_declare(queue='Samt', durable=True) channel.exchange_declare(durable=True, exchange_type='direct', exchange='direct') channel.queue_bind(queue='Samt', exchange='direct', routing_key='samt_routing') # 发送持久化消息 for i in range(10): # 替换为你的实际循环逻辑 message_body = json.dumps({'name':'reza','lastname':'rastgoo','id':i}) # 定义持久化属性 persistent_props = BasicProperties(delivery_mode=2) channel.basic_publish( exchange='direct', routing_key='samt_routing', body=message_body, properties=persistent_props )
额外注意点
- 持久化消息会带来一定性能开销,因为RabbitMQ需要将消息写入磁盘,需根据业务场景权衡使用。
- 即使消息设为持久化,也无法保证100%不丢失(比如消息刚被接收但还未完成磁盘写入时节点崩溃),若需要更强可靠性,可搭配使用RabbitMQ的发布确认机制。
内容的提问来源于stack exchange,提问作者RezaRastgoo
相关产品推荐
相关产品推荐

