Pika:RabbitMQ经典镜像队列向仲裁队列迁移方案咨询
RabbitMQ经典镜像队列迁移至仲裁队列的无删除存量保留方案
你当前预设的先删除旧队列再声明同名称仲裁队列的方案会直接丢失所有存量消息,且会导致迁移窗口内生产者投递的消息丢失,不推荐使用。
目前存在无需提前删除原有队列、可完整保留存量消息、支持零停机的迁移方案,分别可以通过RabbitMQ内置工具或者pika代码实现:
方案一:无需修改业务代码,用Shovel插件实现迁移
该方案仅需要操作RabbitMQ命令即可完成消息同步,无需改动现有生产/消费业务代码:
- 首先开启Shovel插件:
rabbitmq-plugins enable rabbitmq_shovel rabbitmq_shovel_management - 声明新的仲裁队列,可先使用临时名称避免和原队列冲突,比如原队列名为
order_queue,新队列声明为order_queue_quorum,声明时指定队列类型参数即可,用pika实现的代码为:import pika conn = pika.BlockingConnection(pika.ConnectionParameters('你的RabbitMQ地址')) channel = conn.channel() channel.queue_declare( queue='order_queue_quorum', durable=True, arguments={'x-queue-type': 'quorum'} ) - 配置Shovel规则,自动将原经典队列的存量+增量消息同步到新仲裁队列,rabbitmqctl命令如下:
rabbitmqctl set_parameter shovel migrate_classic_to_quorum \ '{"src-uri": "amqp:///", "src-queue": "order_queue", "dest-uri": "amqp:///", "dest-queue": "order_queue_quorum", "src-prefetch-count": 1000, "ack-mode": "on-confirm"}' - 监控原队列消息数,待存量消息全部同步完成、原队列消息数为0后,灰度切换流量:
- 滚动发布生产者,将消息投递目标改为新仲裁队列,确认投递无异常
- 滚动发布消费者,将订阅目标改为新仲裁队列,确认消费无异常
- 业务稳定运行48小时以上后,再删除原经典队列和Shovel配置即可:
rabbitmqctl clear_parameter shovel migrate_classic_to_quorum rabbitmqctl delete_queue order_queue - 如果必须保留原队列的名称,可在原队列删除后,再次用Shovel将
order_queue_quorum的消息同步到新建的同名order_queue仲裁队列,全程不会丢失消息。
方案二:用pika代码实现自定义迁移
如果不想启用额外插件,可以自己写临时迁移脚本完成消息同步:
- 迁移脚本核心逻辑为作为原经典队列的消费者,拉取消息确认后投递到新仲裁队列,核心代码示例如下:
import pika # 原经典队列连接 src_conn = pika.BlockingConnection(pika.ConnectionParameters('你的RabbitMQ地址')) src_ch = src_conn.channel() src_ch.queue_declare(queue='order_queue', durable=True) # 新仲裁队列连接 dest_conn = pika.BlockingConnection(pika.ConnectionParameters('你的RabbitMQ地址')) dest_ch = dest_conn.channel() dest_ch.queue_declare(queue='order_queue_quorum', durable=True, arguments={'x-queue-type': 'quorum'}) def migration_callback(ch, method, properties, body): # 原样投递消息到新仲裁队列,保留原有消息属性 dest_ch.basic_publish( exchange='', routing_key='order_queue_quorum', body=body, properties=properties ) # 确认原队列消息消费完成,避免重复消费 ch.basic_ack(delivery_tag=method.delivery_tag) src_ch.basic_consume(queue='order_queue', on_message_callback=migration_callback) src_ch.start_consuming() - 后续流量切换步骤和Shovel方案完全一致。
关键注意事项
- RabbitMQ不支持直接修改已有队列的类型,不要尝试通过修改原队列的
x-queue-type参数完成切换,该操作会直接报错 - 仲裁队列默认不支持非持久化消息,迁移前请确认原队列的消息都是持久化的,非持久化消息同步到仲裁队列时会被自动丢弃
- 迁移全程原队列可以正常提供服务,仅在切换生产/消费目标时需要滚动发布,无业务停机时间
内容的提问来源于stack exchange,提问作者Benoy John
相关产品推荐
相关产品推荐

