RabbitMQ多服务器主备队列监听实现需求咨询
实现RabbitMQ主备服务器监听同名队列的方案
嘿,针对你要做的主备服务器监听同名队列、故障自动切换的需求,结合你现有用@RabbitListener和SimpleRoutingConnectionFactory的代码,我给你梳理一套可行的方案:
核心思路
要搞定主备切换,核心要解决两个关键问题:
- 主节点挂了之后队列不能丢——得用RabbitMQ的镜像队列,把队列同步到备节点
- 主节点正常时备节点别抢消息,主挂了备自动接棒——通过消费者优先级或者节点亲和性来控制消费权限
第一步:先搞定RabbitMQ的镜像队列配置
首先你得确保RabbitMQ已经搭好主备集群(哪怕是两节点的主备模式),然后给to_client队列设置镜像策略,让队列内容自动同步到备节点。
你可以用RabbitMQ管理控制台,或者直接敲命令行:
# 给to_client队列设置镜像策略,同步到所有集群节点,自动同步内容 rabbitmqctl set_policy ha-to-client "^to_client$" '{"ha-mode":"all", "ha-sync-mode":"automatic"}'
ha-mode":"all:意思是这个队列会在集群所有节点上创建镜像ha-sync-mode":"automatic":新节点(这里就是备节点)加入时自动同步队列里的消息,不用手动触发
第二步:调整Spring AMQP的消费者配置
结合你现有的代码,这里给你两种可行的配置方案,优先推荐第一种:
方案A:消费者优先级(最省心)
给主节点的消费者设置最高优先级,RabbitMQ会优先把消息发给主节点的消费者。只有主节点的消费者断开连接(比如主服务器挂了),备节点的低优先级消费者才会收到消息。
主节点的配置代码
@Configuration public class RabbitMainConfig { @Bean @Primary public SimpleMessageListenerContainer mainListenerContainer(ConnectionFactory connectionFactory) { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory); container.setQueueNames("to_client"); // 主消费者优先级设为10(数值越大优先级越高,范围0-255) container.setConsumerArguments(Map.of("x-priority", 10)); container.setMessageListener(new ClientRabbitService()); return container; } } // 主节点的监听类,和你原来的差不多 @RabbitListener(queues = "to_client") public class ClientRabbitService { @RabbitHandler public void handleMessage(String message) { System.out.println("主节点处理消息:" + message); // 你的业务逻辑写在这 } }
备节点的配置代码
@Configuration public class RabbitStandbyConfig { @Bean public SimpleMessageListenerContainer standbyListenerContainer(ConnectionFactory connectionFactory) { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory); container.setQueueNames("to_client"); // 备消费者优先级设为1(比主节点低就行) container.setConsumerArguments(Map.of("x-priority", 1)); container.setMessageListener(new StandbyClientRabbitService()); return container; } } // 备节点的监听类,逻辑和主节点完全一致 @RabbitListener(queues = "to_client") public class StandbyClientRabbitService { @RabbitHandler public void handleMessage(String message) { System.out.println("备节点处理消息:" + message); // 和主节点一样的业务逻辑 } }
方案B:节点亲和性(绑定主备节点)
如果不想用优先级,也可以让主节点的消费者只连主RabbitMQ节点,备节点的消费者只连备节点,同时靠镜像队列保证队列在两边都存在。当主节点故障时,备节点的消费者会自动接管镜像队列。
结合你现有的SimpleRoutingConnectionFactory,可以这么配:
@Configuration public class RabbitConnectionConfig { // 主节点连接工厂 @Bean public ConnectionFactory mainConnFactory() { CachingConnectionFactory factory = new CachingConnectionFactory(); factory.setHost("你的主节点IP"); factory.setPort(5672); factory.setUsername("rabbitmq用户名"); factory.setPassword("rabbitmq密码"); return factory; } // 备节点连接工厂 @Bean public ConnectionFactory standbyConnFactory() { CachingConnectionFactory factory = new CachingConnectionFactory(); factory.setHost("你的备节点IP"); factory.setPort(5672); factory.setUsername("rabbitmq用户名"); factory.setPassword("rabbitmq密码"); return factory; } @Bean @Primary public SimpleRoutingConnectionFactory routingConnFactory(Map<Object, ConnectionFactory> connFactories) { SimpleRoutingConnectionFactory rcf = new SimpleRoutingConnectionFactory(); rcf.setTargetConnectionFactories(connFactories); // 默认用主节点的连接工厂 rcf.setDefaultTargetConnectionFactory(mainConnFactory()); return rcf; } }
然后主节点的消费者用主连接工厂,备节点用备连接工厂,RabbitMQ默认会自动处理故障转移,主节点挂了之后备节点的消费者就会开始消费镜像队列里的消息。
第三步:验证故障切换
你可以这么测试:
- 启动主备两台服务器,发几条消息到
to_client,确认只有主节点在处理 - 关掉主节点的服务(或者重启主RabbitMQ),再发消息,看备节点是不是自动开始处理了
- 恢复主节点,主节点的消费者会重新连接,因为优先级更高(方案A),后续消息会回到主节点处理
一些注意点
- 镜像队列一定要配置对,不然主节点挂了队列可能直接丢了,备节点也接不到消息
- 消费者优先级方案要求RabbitMQ版本在3.2以上,这个版本之后才支持
x-priority参数 - 备节点的消费逻辑一定要和主节点完全一致,别搞两套逻辑,不然切换后出问题
内容的提问来源于stack exchange,提问作者Grafity08
相关产品推荐
相关产品推荐

