You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 06:06:23