ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RabbitMQ实战指南:SpringBoot整合、死信队列与消息可靠性设计

RabbitMQ实战指南:SpringBoot整合、死信队列与消息可靠性设计 1. 先说结论RabbitMQ到底解决什么问题直接回答被问了几百次的困惑——什么场景会用到RabbitMQ一句话讲清楚当你的业务里存在“生产者只管把消息发出去但不想等消费者马上处理完”的异步需求或者多个服务需要对同一件事做不同的响应时RabbitMQ就是那个中间缓冲区。举个最典型的电商下单场景。用户点击“提交订单”传统同步写法是扣库存、生成订单、发短信、发邮件、加积分全部在一个请求线程里做完。高峰期一个下单接口可能因为发短信服务超时整个请求被拖到5秒才返回用户早跑光了。用RabbitMQ之后主流程只做扣库存、生成订单然后把“订单已创建”这个事件丢进交换机发短信、发邮件、加积分三个消费者各自监听队列谁慢谁快互不影响接口响应直接降到200毫秒以内。我在项目里还遇到过另一类场景——多个系统需要同一份业务数据的场景比如订单创建后数据仓库要同步、搜索引擎要索引、风控系统要检测。如果每个系统都用HTTP去调订单服务接口订单服务会被打成筛子。改成RabbitMQ的广播或者路由模式订单服务只负责把消息投递到交换机三个队列分别绑定自己的路由键各取所需。用不用RabbitMQ有个很粗糙的判断标准这个操作是不是用户无感的用户点了按钮立刻要看到结果的就别上消息队列用户根本不关心结果什么时候到但系统之间必须完成的联动适合丢消息队列。还有一类是流量突刺场景比如秒杀数据库扛不住瞬间写入RabbitMQ天然带削峰填谷的能力。这套技术栈为什么现在这么流行因为SpringBoot已经成了Java后端的事实标准而Spring官方提供的SpringAMQP封装把RabbitMQ原生客户端那一堆繁琐的ConnectionFactory、Channel创建逻辑全部收敛了。你不用再手写一堆样板代码一个RabbitListener注解就能收消息一个RabbitTemplate就能发消息配好Exchange、Queue、Binding之后业务代码只要关心消息本身。这套组合上手成本低、社区资料多、生产环境验证充分从入门到能干活一个周末就够了。接下来我从环境搭建开始走一条完整的路径安装、核心概念、SpringAMQP实战、手动确认与重试、死信队列、广播模式、集群部署、高频面试题把这套东西彻底盘清楚。2. 先把环境跑起来Windows和Docker两条路径2.1 Windows安装最容易被官网下载坑到Windows装RabbitMQ记住一个顺序先装Erlang再装RabbitMQ版本必须匹配。不懂的人直接去RabbitMQ官网下个最新版exe装完启动服务报错一查日志发现Erlang版本不对又要卸了重装。RabbitMQ每个大版本对Erlang版本都有严格要求官网的“Installing on Windows”页面底部有Erlang Version Compatibility表格安装前先看一眼。以RabbitMQ 3.12/3.13系列为例官方推荐Erlang 25.2以上、26.x不要拿一个特别新的Erlang 27去配老版本RabbitMQ兼容问题非常隐蔽。安装完之后RabbitMQ默认是开机自启的Windows服务在服务管理工具里能看到RabbitMQ服务。接着执行两个命令打开管理界面cd C:\Program Files\RabbitMQ Server\rabbitmq_server-3.12.x\sbin rabbitmq-plugins enable rabbitmq_management然后浏览器访问 http://localhost:15672默认账号guest/guest。这里有个巨坑guest只能在localhost登录如果要从别的机器访问管理界面必须新建用户并授权。rabbitmqctl add_user admin admin123 rabbitmqctl set_user_tags admin administrator rabbitmqctl set_permissions -p / admin .* .* .*这个坑我实际踩过开发机上开了管理端口远程连不上一查才知道guest的localhost限制。顺便建议把默认guest用户直接禁用安全很多rabbitmqctl disable_user guest2.2 Docker安装开发环境最省心的方式Docker跑RabbitMQ是我最推荐的方式一条命令的事换版本也方便不会污染宿主机。带管理界面的镜像要用rabbitmq:3-management不带management标签的官方镜像没有Web界面很多人在这一步栽过。version: 3.8 services: rabbitmq: image: rabbitmq:3.13-management container_name: rabbitmq restart: always ports: - 5672:5672 # AMQP协议端口Java客户端连这个 - 15672:15672 # Web管理界面端口 - 15692:15692 # Prometheus监控端口可选 environment: RABBITMQ_DEFAULT_USER: admin RABBITMQ_DEFAULT_PASSWORD: admin123 RABBITMQ_DEFAULT_VHOST: / volumes: - rabbitmq_data:/var/lib/rabbitmq - rabbitmq_log:/var/log/rabbitmq volumes: rabbitmq_data: rabbitmq_log:执行docker-compose up -d然后访问管理界面。这里解释下5672和15672的区别5672是AMQP协议端口你的Java程序通过RabbitTemplate连的就是它15672是HTTP管理端口给浏览器打开管理页面用的。很多新手拿15672去配置Spring Boot怎么都连不上就是这个原因。Docker部署还有一层好处后面搞集群可以直接用docker-compose编排多个容器比在实体机器上拉多个实例省太多事。这个我放到后面的集群章节具体讲。2.3 Spring Boot项目引入SpringAMQP环境起来之后创建一个Spring Boot项目引入依赖这段没什么可说的dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependencyapplication.yml配置spring: rabbitmq: host: localhost port: 5672 username: admin password: admin123 virtual-host: / publisher-confirm-type: correlated # 开启发送方确认 publisher-returns: true # 开启消息路由失败回调 template: mandatory: true # 路由不到队列时让消息返回给发送方 listener: simple: acknowledge-mode: manual # 手动ACK生产环境强烈建议 prefetch: 10 # 每次预取10条根据业务调整 retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2这些参数里最重要的两个acknowledge-mode和prefetch。默认是AUTO模式Spring帮你确认但业务抛异常时消息会被重新投递搞不好就无限循环manual模式完全由你决定什么时候算处理成功可控性最高。prefetch是每个消费者同时持有的未确认消息数默认250条如果消费者处理慢、消息量大内存直接爆。我一般业务处理在100ms级别就配10~30在秒级就配1~3。3. 核心概念搞不懂后面全是糊涂账3.1 用快递站的思维理解四大核心模型RabbitMQ的核心概念说多不多说少不少生产者、交换机、队列、消费者、绑定、路由键、虚拟主机。我用快递站来打个比方保证一遍记住。生产者就是你你要寄一个包裹消息。交换机是快递站的分拣台它不管包裹最终怎么送只看包裹上的标签路由键Routing Key然后决定递给哪条传送带队列。传送带后面站着的快递员消费者再把包裹取走派送。关键点在于队列本身是没有路由能力的它只能被动地接收交换机甩过来的消息。消息最终能不能进到某个队列取决于交换机的类型和路由键的匹配规则。绑定Binding就是把交换机和队列关联起来的那张小卡片上面写着什么样的路由键会被匹配进来。虚拟主机Virtual Host相当于快递站的不同楼层每层楼的设备和人员互不干扰。RabbitMQ默认有个“/”虚拟主机实际项目中建议按环境或业务线拆多个vhost实现逻辑隔离。我在公司就按order、product、user三个vhost分开部署互不影响也方便权限控制。3.2 交换机四种类型的选择逻辑交换机是RabbitMQ最核心的组件类型决定了消息的投递规则。很多人背了四种类型的定义但不知道什么时候该用哪个我来梳理清楚。Direct直连交换机完全精确匹配路由键。订单服务发一条路由键为order.create的消息只有绑定了order.create这个路由键的队列能收到。这是最常用的点对点通知用它就对了。Topic主题交换机支持通配符匹配。order.*能匹配order.create、order.payorder.#能匹配order.create.success等多级路径。系统间集成、按业务类型分发时用得多比如订单模块发消息库存、积分、通知各取一瓢。Fanout扇形交换机不关心路由键把消息复制到所有绑定的队列。广播场景专属比如用户登录后要同步到多个子系统用Fanout一把梭最省事。Ruoyi框架里集成RabbitMQ广播模板用的就是这种模式一次登录事件多个队列各自消费。Headers头交换机不看路由键根据消息头的键值对匹配。现在实际项目里用得很少面试可能会问业务上基本被Topic取代了了解概念即可。选型逻辑很简单一对一选Direct一对多按规则分发选Topic一对多全部通知选Fanout。3.3 消息从发出到消费的完整旅程一条消息的生命周期我建议口述一遍这是面试必考理解透了排错也有方向生产者通过Channel发送消息到交换机connection创建、channel复用这些细节由SpringAMQP处理。交换机根据类型和路由键找到匹配的队列把消息放进去。如果交换机找不到匹配的队列消息要么丢弃要么通过mandatory参数回传给生产者由ReturnCallback处理。队列里的消息等消费者来取RabbitMQ默认采用公平分发轮询但prefetch参数可以改变分发策略。消费者收到消息后执行业务完成后发送ACK确认。如果消费者宕机没有ACK的消息会被重新投递给其他消费者。如果消息一直被拒绝且requeue为false或者消息在队列中超过TTL就会进入死信队列。这个流程里有三个容易出问题的环节第3步路由失败第5步消费者没发送ACK导致消息积压第6步死信配置不到位导致消息堆积在正常队列死循环。后面我分别展开讲。4. SpringAMQP实战把生产者和消费者跑通4.1 声明交换机、队列、绑定的两种姿势SpringAMQP里声明队列和交换机有两种方式一种是用Java配置类显式声明Bean一种是用RabbitListener注解自动声明。前者适合架构清晰、预先规划好的场景后者适合快速开发、随用随建。配置类方式Configuration public class RabbitMQConfig { public static final String EXCHANGE_NAME order.exchange; public static final String QUEUE_NAME order.create.queue; public static final String ROUTING_KEY order.create; Bean public DirectExchange orderExchange() { // durabletrue表示持久化重启不丢失 return new DirectExchange(EXCHANGE_NAME, true, false); } Bean public Queue orderCreateQueue() { // durabletrue队列持久化exclusivefalse多消费者共享autoDeletefalse不自动删除 return new Queue(QUEUE_NAME, true); } Bean public Binding orderBinding() { return BindingBuilder.bind(orderCreateQueue()) .to(orderExchange()) .with(ROUTING_KEY); } }用注解方式更简洁直接在消费端声明Component public class OrderConsumer { RabbitListener(queuesToDeclare Queue( name order.create.queue, durable true )) public void handleOrder(String message) { System.out.println(收到订单消息 message); } }但queuesToDeclare这种方式有个坑它声明队列时不会自动绑定交换机如果你需要绑定交换机还是要写配置类。我在项目中基本都是配置类声明全套只有测试代码才用注解方式保持结构清晰避免消息队列的路由关系散落在各个消费者里面。4.2 RabbitTemplate发送消息的完整玩法发送消息用RabbitTemplate注入之后直接调用即可。最简单的发法Service public class OrderService { Autowired private RabbitTemplate rabbitTemplate; public void createOrder(OrderDTO orderDTO) { // 业务逻辑 rabbitTemplate.convertAndSend( RabbitMQConfig.EXCHANGE_NAME, RabbitMQConfig.ROUTING_KEY, orderDTO // 会被Jackson自动序列化为JSON ); } }SpringAMQP默认用SimpleMessageConverter会把对象转成字节流。生产环境强烈建议自定义Jackson消息转换器消息体是JSON跨系统消费方便得多Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); }在application.yml里不需要额外配置声明了这个Bean之后SpringBoot会自动装配到RabbitTemplate和监听器容器里。发送方确认有两个层级ConfirmCallback确认消息是否到达交换机ReturnCallback确认消息是否从交换机路由到了队列。两个都配上才能完整感知消息投递结果PostConstruct public void init() { rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { log.error(消息发送到交换机失败cause{}, cause); } }); rabbitTemplate.setReturnsCallback(returned - { log.error(消息路由到队列失败exchange{}, routingKey{}, body{}, returned.getExchange(), returned.getRoutingKey(), new String(returned.getMessage().getBody())); }); }这两个回调配合publisher-confirm-type: correlated和template.mandatory: true配置才能正常工作。回调属于异步回调如果发送方在消息发出后立即抛异常回调可能在异常处理后才触发排查乱序问题时要留意这点。4.3 广播模式下项目里的实际用法之前热词里提到Ruoyi集成SpringBootRabbitMQ广播模板我实际也做过类似的。Ruoyi框架里的业务往往包含多端同步需求比如用户下线要通知所有已登录设备或者配置变更要广播给集群里的所有节点。用Fanout交换机实现广播代码极简Configuration public class BroadcastConfig { public static final String FANOUT_EXCHANGE system.broadcast.exchange; Bean public FanoutExchange broadcastExchange() { return new FanoutExchange(FANOUT_EXCHANGE, true, false); } // 多个队列绑定同一个Fanout交换机 Bean public Queue broadcastQueue1() { return new Queue(system.broadcast.queue.1, true); } Bean public Queue broadcastQueue2() { return new Queue(system.broadcast.queue.2, true); } Bean public Binding binding1() { return BindingBuilder.bind(broadcastQueue1()).to(broadcastExchange()); } Bean public Binding binding2() { return BindingBuilder.bind(broadcastQueue2()).to(broadcastExchange()); } }发送方rabbitTemplate.convertAndSend(BroadcastConfig.FANOUT_EXCHANGE, , message);注意Fanout交换机连路由键都不需要传空字符串就行写什么都一样。接收方各写各的RabbitListener各干各的事。这种模式的好处是新增一个广播消费者只需要新增队列绑定不用改任何已有代码。5. 手动确认与重试机制生产环境能不能稳住的胜负手5.1 为什么默认的AUTO确认靠不住SpringAMQP默认的确认模式是AUTO它的行为逻辑是监听器方法正常返回Spring自动发送ACK抛出异常Spring自动发送NACK并决定是否重新入队。看起来很智能但坑就藏在“自动”两个字里。举个例子消费者处理消息时数据库连接池满了抛了个数据库连接异常。AUTO模式下Spring会把消息NACK并重新放回队列然后立刻再次投递又抛异常又放回队列无限循环。如果消息本身有问题比如JSON解析失败AUTO模式下也会无限重投消息永远卡在队头后面正常消息全部堵死。这就是网上常说的“毒消息问题”。手动确认模式把控制权完全交给你处理成功才ACK处理失败可以把消息拒绝掉或者让它进死信队列甚至可以记录日志后直接丢弃。代价是代码必须自己兜底漏了ACK会让消息一直处于Unacked状态堆积在消费者本地造成队列阻塞。5.2 手动确认的标准姿势消费端开启手动确认后监听器方法签名要加上Channel和DeliveryTag两个参数Component public class OrderConsumer { RabbitListener(queues order.create.queue) public void handleOrder(String message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { // 1. 解析消息 OrderDTO orderDTO JSON.parseObject(message, OrderDTO.class); // 2. 执行业务逻辑比如写入数据库 orderService.create(orderDTO); // 3. 成功确认只确认当前这条消息 channel.basicAck(deliveryTag, false); } catch (Exception e) { // 4. 处理失败拒绝消息且不重新入队 channel.basicReject(deliveryTag, false); // 如果需要记录到死信或日志系统 log.error(订单消息处理失败: {}, message, e); } } }这里有几个细节要注意deliveryTag是消息在当前Channel上的自增序号确认时必须携带。它不代表全局消息ID。basicAck的第二个参数multiplefalse表示确认当前这条true表示确认当前及之前所有未确认的消息。批量确认可以提高吞吐但如果中途一条失败会导致前面的成功消息也算未确认特殊情况别乱用。basicReject的第二个参数requeuefalse表示不重新入队true表示重新入队。不要随便设为true否则毒消息无限循环的坑就又回来了。还有一种情况是手动确认时使用basicNack它比basicReject多一个multiple参数支持批量拒绝其他行为类似。5.3 重试机制的两层设计RabbitMQ的重试有两个维度很多人搞混SpringAMQP层面的重试和RabbitMQ消息投递层面的重试。SpringAMQP的listener.simple.retry配置是消费者本地重试消息已经出队了在本地处理失败后会根据配置重试N次重试期间不发送ACK。这个机制下要注意重试期间消息不返回队列消费者线程会长时间占用如果重试全部失败消息会被丢弃或者进入死信队列而不是回到RabbitMQ。而RabbitMQ层面的重新投递是消费者进程崩溃、连接断开或者channel.basicReject(deliveryTag, true)时消息回到队列头部等待重新投递给其他消费者。这两个层面互不干扰但容易叠加。我的推荐设计是本地重试3次重试期间做间隔递增重试耗尽后手动拒绝进死信队列。SpringAMQP的本地重试配置只管自己重试它不负责把重试失败的消息转投到死信这个转投逻辑要自己写或者靠死信交换机配置实现。spring: rabbitmq: listener: simple: acknowledge-mode: manual retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2这个配置的意思是首次失败后等1秒重试第二次失败后等2秒重试最多尝试3次。我在一个订单回调项目里试过数据库连接中断恢复后重试2次内基本能恢复配合死信基本没有消息丢失。5.4 如何取当前重试次数热词里有一条很具体的问题怎么在消费者里取当前是第几次重试这里涉及一个命名的坑。SpringAMQP重试时的消息头里带了一个x-retried-count头我们可以通过注解把它拿出来RabbitListener(queues order.create.queue) public void handleOrder(String message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Header(name x-retried-count, defaultValue 0) int retryCount) throws IOException { try { // 业务处理 orderService.create(JSON.parseObject(message, OrderDTO.class)); channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(订单处理失败当前第{}次重试, retryCount 1, e); if (retryCount 2) { // 已重试3次拒绝并进死信 channel.basicReject(deliveryTag, false); } else { // 抛出异常触发Spring重试机制 throw e; } } }注意一个矛盾点手动确认模式下如果SpringAMQP的retry配置了本地重试其实是容器在管理先于我们自己的catch代码。上面的代码里判断retryCount其实是在容器重试之后才到达的实际上SpringAMQP的本地重试是一次方法调用的内部重试对于方法内部而言每次重试都会带着递增的x-retried-count头重新进入方法。但如果你想走手动确认自己控制重试可以关闭SpringAMQP的retry配置在catch里自己判断重试次数并决定是重新入队还是进死信RabbitListener(queues order.create.queue) public void handleOrder(String message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { // 业务处理 orderService.create(JSON.parseObject(message, OrderDTO.class)); channel.basicAck(deliveryTag, false); } catch (Exception e) { // 自己统计重试次数这里假设通过Redis或本地计数实现 int retryCount getRetryCount(message); // 从Redis读取 if (retryCount 3) { channel.basicReject(deliveryTag, false); return; } incrementRetryCount(message); channel.basicNack(deliveryTag, false, true); // 重新入队 } }但这种自己控制的方式有个隐藏问题重新入队的消息会放到队列头部如果不断重试可能导致后面的消息持续等待。生产环境我更推荐把重试逻辑交给SpringAMQP的retry机制然后配合死信队列兜底代码最简单行为最可控。6. 死信队列最后一道防线怎么设计6.1 死信到底是什么为什么非配不可死信队列Dead Letter QueueDLQ是消息队列里最后一道兜底防线理解它的触发条件比死信本身重要。消息变成死信有三种情况消费者调用了basicReject或basicNack且requeuefalse消息被“判死刑”。消息超过了队列设置的TTL过期时间没人消费直接进死信。队列达到了最大长度新消息无法入队Mapping时旧消息可以进死信。没有死信队列的后果我实际经历过消费者一直拒绝消息消息不断被丢弃业务数据悄无声息地消失。你从日志里发现消费失败率暴涨但消息已经没了想要重新处理根本没机会。配了死信队列所有处理不掉的消息会汇聚到DLQ里你可以写一个独立程序定期扫描DLQ做补偿也可以盯监控报警出了问题还留了“后悔药”。6.2 死信交换机配置的完整代码死信不是RabbitMQ默认开启的功能需要你在声明正常队列时指定死信交换机。配置方式如下Configuration public class DLQConfig { public static final String ORDER_EXCHANGE order.exchange; public static final String ORDER_QUEUE order.create.queue; public static final String ORDER_ROUTING_KEY order.create; public static final String ORDER_DLX order.dlx; public static final String ORDER_DLQ order.create.dlq; public static final String ORDER_DLQ_ROUTING_KEY order.create.dlq; // 正常业务交换机 Bean public DirectExchange orderExchange() { return new DirectExchange(ORDER_EXCHANGE, true, false); } // 死信交换机实际生产可以多个业务共用一个 Bean public DirectExchange orderDlx() { return new DirectExchange(ORDER_DLX, true, false); } // 正常队列绑定死信交换机 Bean public Queue orderQueue() { MapString, Object args new HashMap(); args.put(x-dead-letter-exchange, ORDER_DLX); args.put(x-dead-letter-routing-key, ORDER_DLQ_ROUTING_KEY); // 可选队列TTL比如订单消息30分钟没消费就进死信 // args.put(x-message-ttl, 30 * 60 * 1000); return new Queue(ORDER_QUEUE, true, false, false, args); } // 死信队列 Bean public Queue orderDlq() { return new Queue(ORDER_DLQ, true); } Bean public Binding orderBinding() { return BindingBuilder.bind(orderQueue()).to(orderExchange()).with(ORDER_ROUTING_KEY); } Bean public Binding orderDlqBinding() { return BindingBuilder.bind(orderDlq()).to(orderDlx()).with(ORDER_DLQ_ROUTING_KEY); } }这个配置有两个参数特别关键x-dead-letter-exchange指定死信交换机x-dead-letter-routing-key指定死信消息的路由键。如果后者不指定死信消息会沿用原消息的路由键导致死信消息进不了死信队列。我在一个项目里就因为这个坑死信队列里一条消息都没有业务数据丢了都不知道。正常消费者处理失败时配合手动拒绝RabbitListener(queues ORDER_QUEUE) public void handleOrder(String message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { orderService.create(JSON.parseObject(message, OrderDTO.class)); channel.basicAck(deliveryTag, false); } catch (Exception e) { // requeuefalse消息触发死信转发 channel.basicReject(deliveryTag, false); log.error(订单处理失败已转发死信队列, e); } }死信队列的消费者要单独写一个类职责就是补偿、告警、人工介入。我在实际项目中死信消费者做的事很简单解析消息记录到一张compensation表然后钉钉报警。等开发排查完问题再从补偿表里捞数据重新投递。6.3 死信30分钟会压多少一个实际推算题热词里有个具体问题“死信30分钟会压多少”。这个实际是个容量评估问题。假设你配了队列TTL为30分钟那么30分钟内的消息都应该被正常消费只有消费失败或超时的消息会进死信。死信30分钟的积压量 单位时间失败率 × 30分钟内的消息总量再乘以重试失败的比例。举个例子业务高峰期每秒进队列100条消息正常处理失败率5%SpringAMQP本地重试3次后成功救回60%最终进死信的约为每秒100 × 5% × 40% 2条30分钟积压3600条。每条消息按1KB算死信队列30分钟数据量约3.6MB对RabbitMQ来说毫无压力。但如果失败率不是5%而是50%同样的参数下30分钟会压36000条、约36MB这就会对死信消费者的处理能力提出要求。所以设计死信队列时一定要先估算业务失败率再决定要不要在死信队列上加TTL做二次失效。我在高并发项目里的习惯是死信队列再配一个二级死信防止补偿不及时导致消息越积越多。7. 集群部署与高可用用Docker Compose三步搭建7.1 为什么要做集群而不是单机单机RabbitMQ的问题很明显节点挂了所有队列、交换机、绑定全部不可用整个消息链路瘫痪。磁盘损坏可能连持久化消息一起丢。生产环境至少要保证高可用而RabbitMQ集群的核心是镜像队列Quorum Queue或者仲裁队列让同一个队列的数据在多个节点上有副本某个节点挂了其他节点还能继续服务。这里说个面试高频题RabbitMQ集群分为普通集群和镜像队列两种模式。普通集群下每个节点只存队列元数据和一个完整的队列副本其实普通集群的队列数据只存在一个节点上其他节点只存储队列的元信息位置。消费者连到其他节点时会转发到真正持有队列的节点。这种方式能提高吞吐但不解决单点问题。镜像队列才是把队列内容复制到多个节点才叫真正的高可用。现在官方更推荐Quorum Queue仲裁队列基于Raft协议替代了老的镜像队列行为更稳定但配置方式略有不同。7.2 Docker Compose搭建三节点集群用Docker Compose在开发机上验证集群最省事。我给出一个三节点集群的部署配置使用rabbitmq:3.13-management镜像version: 3.8 services: rabbit1: image: rabbitmq:3.13-management hostname: rabbit1 container_name: rabbit1 ports: - 5672:5672 - 15672:15672 environment: RABBITMQ_ERLANG_COOKIE: secret_cookie_value RABBITMQ_DEFAULT_USER: admin RABBITMQ_DEFAULT_PASS: admin123 volumes: - rabbit1_data:/var/lib/rabbitmq rabbit2: image: rabbitmq:3.13-management hostname: rabbit2 container_name: rabbit2 ports: - 5673:5672 - 15673:15672 environment: RABBITMQ_ERLANG_COOKIE: secret_cookie_value RABBITMQ_DEFAULT_USER: admin RABBITMQ_DEFAULT_PASS: admin123 volumes: - rabbit2_data:/var/lib/rabbitmq depends_on: - rabbit1 rabbit3: image: rabbitmq:3.13-management hostname: rabbit3 container_name: rabbit3 ports: - 5674:5672 - 15674:15672 environment: RABBITMQ_ERLANG_COOKIE: secret_cookie_value RABBITMQ_DEFAULT_USER: admin RABBITMQ_DEFAULT_PASS: admin123 volumes: - rabbit3_data:/var/lib/rabbitmq depends_on: - rabbit1 volumes: rabbit1_data: rabbit2_data: rabbit3_data:启动三个容器之后进入rabbit1执行集群命令docker exec -it rabbit1 bash rabbitmqctl stop_app rabbitmqctl reset rabbitmqctl start_app然后让rabbit2和rabbit3加入集群docker exec -it rabbit2 bash rabbitmqctl stop_app rabbitmqctl reset rabbitmqctl join_cluster rabbitrabbit1 rabbitmqctl start_app docker exec -it rabbit3 bash rabbitmqctl stop_app rabbitmqctl reset rabbitmqctl join_cluster rabbitrabbit1 rabbitmqctl start_app关键点是所有节点的Erlang Cookie必须相同否则节点间无法互相认证。上面的配置里我已经统一成secret_cookie_value了。集群起来之后如果你的队列还是默认的单副本节点挂了队列照样不可用。需要把队列设置为镜像队列或Quorum Queue。老版本镜像队列策略如下docker exec -it rabbit1 bash rabbitmqctl set_policy ha-all ^ {ha-mode:all}这条命令把所有以任意名字开头的队列都设置为所有节点同步镜像。对于新项目建议直接使用Quorum QueueBean public Queue orderQueue() { MapString, Object args new HashMap(); args.put(x-queue-type, quorum); // 死信配置同样适用 args.put(x-dead-letter-exchange, ORDER_DLX); args.put(x-dead-letter-routing-key, ORDER_DLQ_ROUTING_KEY); return new Queue(ORDER_QUEUE, true, false, false, args); }x-queue-typequorum会让队列在集群里以仲裁副本方式存在多数节点确认写入才返回成功比镜像队列更安全。但要注意Quorum Queue不支持部分参数如TTL per-message使用前确认官方文档。7.3 客户端连接集群的方式Spring Boot项目连集群地址可以用逗号分隔spring: rabbitmq: addresses: 192.168.1.10:5672,192.168.1.11:5672,192.168.1.12:5672 username: admin password: admin123 virtual-host: /SpringAMQP底层用的是RabbitMQ Java Client它支持自动重连和地址列表故障转移某个节点挂了会自动切换到其他可用节点。这里要注意连接均衡和队列副本是两回事连接随便连哪个节点都行但队列数据副本取决于你队列的类型和策略配置。8. RabbitMQ、RocketMQ和Kafka怎么选热词里也出现了RabbitMQ、RocketMQ和Kafka的对比这个在面试里几乎必问。我用一张表讲清楚再给个选型策略。维度RabbitMQRocketMQKafka语言ErlangJavaScala/Java吞吐量万级/秒十万级/秒百万级/秒消息可靠性高支持事务和Confirm高支持事务较高需配置消息顺序单队列有序队列有序分区内有序延迟微秒级毫秒级毫秒级管理界面完善一般一般社区生态非常成熟阿里系国内用得多大数据生态标配适用场景业务系统解耦、异步交易链路、金融级可靠日志采集、数据管道、流处理选型策略我总结成一句话业务系统内部的消息传递、异步解耦、RPC替代选RabbitMQ处理大数据量的日志流、埋点数据、需要长时间保存和回放的数据管道选Kafka业务和日志之间需要兼顾同时又追求极致的可靠性选RocketMQ。RabbitMQ的优势不在吞叶量而在灵活的路由模型、完善的管理界面、丰富的插件生态。绝大多数互联网应用的业务消息量RabbitMQ完全够用。我做过日活百万级别的电商系统峰值每秒消息量两三千RabbitMQ毫无压力。真正上到每秒十万级才需要考虑换RocketMQ或者Kafka。9. 常见排查链路消息积压和丢消息9.1 消息积压了怎么查消息积压是最常见的线上事故特征是管理界面能看到queue里Ready的消息数持续上涨消费者处理不过来。排查链路我按经验排序第一步看消费者的并发和prefetch。如果单机只有一个消费者实例prefetch又很大消息会累积在消费者本地表面看队列Ready在下降实际消费者内存要爆了。增加消费者实例数比如部署多个应用实例远比调大prefetch有效。第二步看消费者是否在疯狂重试。如果业务抛异常、SpringAMQP的retry又开着消费者会反复处理同一条消息什么都不消费物理上积压。管理界面能看到Unacked消息数抬升。这种时候去看日志的错误信息通常都是下游依赖挂了。第三步看死信有没有生效。如果你没有配死信重试耗尽的消息会被丢弃队列积压不明显但你丢了数据。如果配了死信但消息没进DLQ多半是x-dead-letter-routing-key没配置对。第四步用命令行捞数据。消息积压要快速清理除了修好消费者可以临时加一个消费者做批量转存把队列里积压的消息转移到另一个临时队列慢慢消化或者直接丢到文件/数据库里做离线处理# 查看队列状态 rabbitmqctl list_queues name messages messages_ready messages_unacknowledged # 查看消费者状态 rabbitmqctl list_consumers9.2 消息丢失的三个环节消息丢失的高发环节有三个每个都要防生产者丢消息消息发出后交换机或者队列不可用消息直接丢失。解决办法是开启ConfirmCallback和ReturnCallback发送失败就重新发送或记录报警。它的保障前提是开启了publisher-confirm-type: correlated。RabbitMQ自身丢消息队列和消息都非持久化重启之后全没了。解决办法是持久化三件套队列durabletrue交换机durabletrue消息发送时设置MessageDeliveryMode.PERSISTENT。SpringAMQP里可以用convertAndSend发送对象时消息转换器默认写入持久化属性但如果用Message对象自定义构造的要手动设置Message message MessageBuilder.withBody(body) .setContentType(application/json) .setDeliveryMode(MessageDeliveryMode.PERSISTENT) .build(); rabbitTemplate.send(EXCHANGE_NAME, ROUTING_KEY, message);消费者丢消息消费者收到消息后还没处理完就宕机如果没等ACK消息会重新投递。但如果用的是AUTO确认且方法内出异常前已经发了ACK消息就丢了。所以手动确认是消费者不丢消息的底线。综合一句话生产端开ConfirmReturn中间件开持久化消费端手动ACK三重保障叠加消息才能做到基本不丢。10. 面试实战这六个问题讲清楚就过关10.1 为什么使用消息队列直接回答三个字解耦、异步、削峰这个问题的标准回答框架“我们系统有A和B两个服务A产生核心业务事件B是旁路系统如果用HTTP同步调用B挂掉会影响A响应耗时也会增加。引入RabbitMQ后A只负责发消息B什么时候消费、消费成不成功都不影响A的主链路。另外在秒杀场景瞬时流量远大于数据库能承受的量通过队列做缓冲消费端按自己能力慢慢处理。”关键是要结合自己的项目讲不要只背概念。讲清楚你的业务里哪个环节是异步解耦哪个环节是削峰填谷面试官马上知道你是真做过还是背的。10.2 如何保证消息不丢失答三端参照上一节的方案生产端Confirm、中间件持久化、消费端手动ACK。补充一个细节持久化不是写磁盘的绝对保证RabbitMQ是先把消息写到内存再异步刷盘极端情况宕机可能丢几毫秒的数据。如果对可靠性有极致要求考虑镜像队列发布确认消费手动确认三层组合。10.3 如何保证消息不被重复消费关键在幂等RabbitMQ不保证消息只被消费一次它在网络抖动、消费者重试等场景下天然会投递重复消息。解决办法不是消除重复投递而是让消费逻辑幂等——即同一条消息处理两次和一次最终结果一致。常用的幂等方案有三种数据库唯一约束消息里带一个业务唯一ID比如订单号在订单表建唯一索引重复插入直接报DuplicateKey异常捕获后视为成功。Redis分布式锁用业务ID做Key消费前先setnx处理完再删除。状态机校验比如订单状态由“待支付”到“已支付”只能流转一次第二次到达时发现状态已经是“已支付”就直接跳过。我在项目里最常用第一个方案简单可靠不需要额外的中间件依赖。每次业务逻辑操作事务表时把消息ID塞进一张消费记录表做唯一约束重复消息直接被数据库挡下。10.4 消息积压怎么处理参考第9.1节的排查链路面试时按步骤回答重点提扩容消费者和临时转存两个运维操作。补充一点如果是消费速度追不上生产速度可以考虑多队列分摊把一个大队列拆成多个分区队列每个消费者只处理一个分区的消息。10.5 如何保证消息顺序消费这个题目比较坑因为很多场景根本不需要全局顺序。RabbitMQ的消息顺序保障是“单队列内有序”消费者并发处理时多线程消费无法保证顺序。方案一把需要保证顺序的关联消息发到同一个队列用同一个路由键消费者单线程处理。比如同一个订单的操作事件永远发到同一个queues消费者用单线程监听顺序就保证了。方案二如果消费者需要并发可以用基于业务ID的取模路由到不同队列同一个业务ID的消息只会进入一个队列进而被同一个消费者线程处理。这个方法用Topic交换机把order.${orderId}作为路由键消费端每个队列只消费一个orderId下的消息。10.6 Exchange和Queue的区别是什么绑定关系是什么这个问题考察基础概念我用一句话回答Exchange负责接收消息并判断路由规则Queue负责存储消息并等待消费者拉取Binding定义了Exchange和Queue的关联关系包括交换机类型和路由键匹配规则。之后可以展开四种交换机类型的匹配逻辑基本就过关了。11. 一套可以直接抄的完整项目实战最后给一个可以直接复制到项目里的完整示例以一个订单创建发送通知的场景为例包含配置、生产者、消费者、死信全流程跑通。生产者端完整代码Service public class OrderMessageProducer { Autowired private RabbitTemplate rabbitTemplate; /** * 发送订单创建消息 */ public void sendOrderCreated(OrderDTO orderDTO) { CorrelationData correlationData new CorrelationData(orderDTO.getOrderId()); rabbitTemplate.convertAndSend( DlxConfig.ORDER_EXCHANGE, DlxConfig.ORDER_ROUTING_KEY, orderDTO, correlationData ); } }消费端完整代码含手动确认、重试记录、失败转死信Component Slf4j public class OrderConsumer { Autowired private OrderService orderService; RabbitListener(queues DlxConfig.ORDER_QUEUE) public void handleOrder(String message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { OrderDTO orderDTO JSON.parseObject(message, OrderDTO.class); // 业务幂等检查订单是否已处理 if (orderService.isProcessed(orderDTO.getOrderId())) { channel.basicAck(deliveryTag, false); return; } // 执行真正的业务逻辑 orderService.createOrder(orderDTO); channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(订单消息处理异常msgId{}, message, e); channel.basicReject(deliveryTag, false); // 可在这里额外报警 } } }死信消费者负责补偿和告警Component Slf4j public class OrderDlqConsumer { RabbitListener(queues DlxConfig.ORDER_DLQ) public void handleDeadMessage(String message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { // 1. 记录补偿表 // 2. 发送钉钉/邮件告警 // 3. 如果具备补偿条件重新投递到业务队列 log.error(死信消息人工介入处理内容{}, message); channel.basicAck(deliveryTag, false); } catch (Exception e) { channel.basicReject(deliveryTag, true); } } }这套结构可以直接放到任何基于Spring Boot的业务系统里改改包名和业务方法就能用。12. 避坑经验汇总写了这么多最后把实战中踩过、见过的坑集中列出来给后来人提个醒第一条队列、交换机的名字不要用魔法值散落在代码里。我见过一个项目路由键在五个类里各写了一遍后来改了一次路由键漏掉一个线上两条链路静默断掉排查了整整一天。常量类或者配置类统一管理每处引用都走同一个静态字段。第二条消费者监听器方法里不要做太重的耗时操作。默认情况下SpringAMQP的监听器线程数量是固定的默认最大并发256如果每条消息处理耗时超过秒级有限的线程会被占满新消息全部积压。要么把耗时操作改成异步要么调大并发数listener.simple.concurrency和max-concurrency。但这只是临时方案根本解法是把大任务拆小。第三条RabbitMQ不要当作数据库用消息要设置TTL。队列里的消息如果没有TTL也不被消费会一直在内存/磁盘里占空间。偶发积压不可怕可怕的是积压消息一直不清堆积几百万条之后整个RabbitMQ性能崩塌。重要数据消费完要存档到业务数据库队列只是过河用的桥。第四条spring-boot-starter-amqp的版本要和Spring Boot版本匹配。这不算RabbitMQ的锅Spring生态老问题。版本不一致通常表现为启动时一堆NoSuchMethodError、BeanCreationException去Maven仓库查一下Spring Boot和amqp的依赖关系就能解决。第五条RabbitMQ的默认心跳是60秒网络不稳定时连接容易断。如果部署在跨机房或者WAN环境把心跳时间调大客户端连接属性设置requestedHeartbeat或者用SpringAMQP的connectionTimeout和handshakeTimeout避免服务间偶发断开又自动重连造成消息重复投递。第六条生产环境一定要开监控。RabbitMQ管理界面的Charts只能看最近一段时间生产上建议接Prometheus监控采集队列积压、连接数、Channel数、消费者数量配合AlertManager设置积压阈值报警。我干活的项目里消息积压超过1000就触发告警运维能把事故扼杀在摇篮里。我在实际项目里反复体会最深的是RabbitMQ的入门和一知半解之间隔着一个“故障处理”的距离而一知半解和真正熟练之间隔着一个“迁移/扩展”的距离。把这篇文章里的每个配置亲手敲一遍再把场景跑通尤其是手动确认、重试、死信这条链路走完你对这套技术栈的理解就算真正到位了。后面做高可用、性能调优、跨系统集成都会顺手很多。
返回列表