尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

RabbitMQ消息可靠性实战:从生产者确认到死信队列

RabbitMQ消息可靠性实战:从生产者确认到死信队列 之前在做订单系统的时候遇到一个很典型的问题支付回调消息发送到 RabbitMQ消费者那边偶尔就是收不到日志里也没有报错。排查了大半天最后发现是生产者发消息时没有开启确认机制路由失败的消息直接丢了。这种问题在面试里几乎成了固定的“送命题”——“RabbitMQ 消息为什么会丢失你会怎么处理”本文将围绕这个高频场景题从消息投递的全链路拆解“消息失联”的各种可能原因并给出完整的 Spring Boot RabbitMQ 实战代码包含生产者确认、消费者手动 ACK、持久化、死信队列、幂等消费等内容。无论你是准备 Java 面试还是在项目中真正需要保障消息可靠性这篇文章都建议看完。1. 背景与核心概念1.1 什么是“消息投递失联”“消息投递失联”并不是一个严格的技术术语而是面试官对“消息丢失、消息没有到达消费者、消费者没有成功处理消息”这一类问题的统称。在 RabbitMQ 中一条消息从生产者发出到消费者最终处理完成要经过多条链路生产者 - Exchange交换机 - Queue队列 - 消费者这条链路上任何一个环节出现闪失消息都可能“失联”。常见的失联表现有几种生产端 send 方法执行成功但消息没有进入队列。消息进入了队列但消费者没有收到。消费者收到了消息但处理失败消息被误删。消息本身没问题但 RabbitMQ 节点重启后消息丢了。在 Java 项目里这些问题往往不会直接抛出异常而是“静默丢失”。这也是它坑人的地方系统看起来一切正常但数据就是对不上。1.2 为什么这是面试重灾区这道题在 Java 面试中出现频率极高主要原因有三个。第一消息可靠性是 MQ 中间件的核心设计点能考察候选人是否真正理解消息中间件底层机制而不只是会调用 API。第二生产环境消息丢失问题很难复现需要候选人具备完整的链路排查意识。第三这类场景题通常不会只停留在“丢消息怎么办”面试官会继续追问如何保证消息不重复消费如何处理消费失败如何保证消息顺序这些问题层层递进。所以把这道题吃透不只是在背一道“八股文”而是建立一套可靠消息投递的工程思维。2. 环境准备与版本说明本文采用 Spring Boot 2.x Spring AMQP 的方式演示这是目前 Java 后端最主流的 RabbitMQ 集成方式。组件说明JDK1.8 或 11本文示例默认 JDK 8Spring Boot2.5.x 或 2.7.xSpring AMQP由 Spring Boot 依赖管理统一控制版本RabbitMQ3.8示例使用 Docker 安装构建工具Maven 3.6IDEIntelliJ IDEA 或 Eclipse如果你本地没有安装 RabbitMQ推荐用 Docker 快速启动一个实例docker run -d --name rabbitmq \ -p 5672:5672 \ -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3.8-management其中 5672 是 AMQP 协议端口15672 是管理控制台端口。启动成功后通过浏览器访问http://localhost:15672使用admin / admin123登录即可。如果你用的不是 Docker也可以在 Windows 或 Linux 上直接安装 RabbitMQ核心配置思路完全一致只是部署方式不同。不同版本的 RabbitMQ 在原生队列和仲裁队列上有所差异本文以普通持久化队列为例仲裁队列会在最佳实践小节单独说明。3. 消息投递全链路拆解丢失发生在哪一环要回答“消息为什么会丢”首先要把链路分成三段来分析。每一段都有对应的可靠性机制。3.1 第一段生产者 - 交换机3.1.1 丢失原因生产者的消息是发送到 Exchange交换机而不是直接发送到队列。在这一段消息丢失最常见的原因有两个网络抖动消息根本没有到达 RabbitMQ 服务端。消息到达 Exchange但 Exchange 根据 RoutingKey 找不到匹配的队列。第一个原因属于网络层面的不可控因素TCP 连接断了消息自然发不出去。第二个原因常见于 RoutingKey 拼写错误或者队列没有正确绑定到交换机。3.1.2 解决方案Publisher ConfirmRabbitMQ 提供了 Publisher Confirm 机制也就是生产者确认。当生产者发送消息后RabbitMQ 服务端会异步返回一个确认结果ack消息已被服务端接收。nack消息被服务端拒绝。在 Spring Boot 中开启发布者确认的方式如下spring: rabbitmq: publisher-confirm-type: correlatedcorrelated表示支持异步回调回调结果中可以包含你发送消息时携带的 correlationData。3.1.3 解决方案Mandatory Return针对“消息到达交换机但路由不到队列”的情况需要开启 mandatory 参数。开启后如果消息路由失败RabbitMQ 会将消息回传给生产者触发 ReturnCallback。3.2 第二段交换机 - 队列这一段相对简单。只要交换机、队列都存在绑定关系正确消息就可以进入队列。但这里有一个隐藏风险那就是持久化问题。如果队列没有设置持久化或者交换机没有设置持久化RabbitMQ 节点重启后相关的元数据会丢失队列里的消息自然也没了。RabbitMQ 的持久化分为三个层面交换机持久化durabletrue交换机元数据不丢失。队列持久化durabletrue队列本身不丢失。消息持久化投递模式设置为MessageDeliveryMode.PERSISTENT消息会被写入磁盘。三者缺一不可。很多人只记得给队列加 durable却忘了消息本身的投递模式。3.3 第三段队列 - 消费者这是“消息失联”的高发区也是最容易在面试中展开聊的部分。3.3.1 丢失原因自动 ACK默认情况下消费者使用的是自动确认AUTO ACK。RabbitMQ 一旦把消息发给消费者就立即从队列中移除该消息无论消费者是否处理成功。如果消费者在接收消息后、业务逻辑还没执行完就宕机了这条消息就彻底丢了。具体流程是这样的RabbitMQ 投递消息 - 消费者收到消息 - RabbitMQ 删除消息 - 消费者宕机 - 消息丢失3.3.2 解决方案手动 ACK改成手动确认模式后消费者处理完成业务逻辑再调用channel.basicAck告诉 RabbitMQ 可以删除消息了。如果处理失败可以调用basicNack拒绝消息或调用basicReject丢弃消息。手动 ACK 的关键代码在消费者监听器中下文实战部分会给出完整示例。3.4 面试高频追问消息确认机制总结面试官通常还会追问一个对比问题事务机制和 Confirm 机制有什么区别RabbitMQ 也提供了事务机制通过txSelect、txCommit来确保消息发送成功。但事务机制是同步阻塞的会严重拉低吞吐量所以在生产环境中基本不会使用。Confirm 机制是异步的性能好得多这也是官方推荐的方式。机制性能可靠性使用场景AMQP 事务低同步阻塞高不推荐吞吐量要求低的场景Publisher Confirm高异步回调高生产环境主流方案4. 完整实战Spring Boot 实现可靠消息投递下面我们用代码把整条可靠链路串起来。项目结构如下rabbitmq-reliable-demo/ ├── pom.xml └── src/main/java/com/example/rabbitmq/ ├── RabbitmqReliableDemoApplication.java ├── config/ │ └── RabbitConfig.java ├── producer/ │ ├── OrderMessage.java │ └── OrderMessageProducer.java ├── consumer/ │ └── OrderMessageConsumer.java └── common/ └── Result.java4.1 创建项目并添加依赖创建一个 Maven 项目在pom.xml中加入 Spring Boot 和 AMQP 相关依赖parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version relativePath/ /parent dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency /dependencies4.2 配置文件在application.yml中加入如下配置spring: application: name: rabbitmq-reliable-demo rabbitmq: host: localhost port: 5672 username: admin password: admin123 # 开启发布者确认 publisher-confirm-type: correlated # 开启路由失败回调 publisher-returns: true template: # 模板级别的 mandatory确保消息路由失败时能走 ReturnCallback mandatory: true listener: simple: # 手动 ACK acknowledge-mode: manual # 单条线程消费便于观察 concurrency: 1 max-concurrency: 1这里有两个关键配置publisher-confirm-type: correlated开启 Confirm 异步回调。acknowledge-mode: manual消费者手动确认而不是自动确认。4.3 配置交换机、队列和绑定关系创建一个RabbitConfig配置类声明订单交换机、订单队列以及一个死信队列。package com.example.rabbitmq.config; import org.springframework.amqp.core.*; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; Configuration public class RabbitConfig { public static final String ORDER_EXCHANGE order.exchange; public static final String ORDER_QUEUE order.queue; public static final String ORDER_ROUTING_KEY order.create; public static final String DEAD_LETTER_EXCHANGE order.dead.exchange; public static final String DEAD_LETTER_QUEUE order.dead.queue; public static final String DEAD_LETTER_ROUTING_KEY order.dead; Bean public MessageConverter messageConverter() { // 使用 JSON 序列化消息体便于调试和跨语言消费 return new Jackson2JsonMessageConverter(); } Bean public DirectExchange orderExchange() { // durabletrue交换机持久化 return new DirectExchange(ORDER_EXCHANGE, true, false); } Bean public DirectExchange orderDeadExchange() { return new DirectExchange(DEAD_LETTER_EXCHANGE, true, false); } Bean public Queue orderQueue() { MapString, Object args new HashMap(); // 死信交换机 args.put(x-dead-letter-exchange, DEAD_LETTER_EXCHANGE); // 死信路由键 args.put(x-dead-letter-routing-key, DEAD_LETTER_ROUTING_KEY); // durabletrue队列持久化排他false不自动删除 return new Queue(ORDER_QUEUE, true, false, false, args); } Bean public Queue orderDeadQueue() { return new Queue(DEAD_LETTER_QUEUE, true, false, false); } Bean public Binding orderBinding() { return BindingBuilder.bind(orderQueue()) .to(orderExchange()) .with(ORDER_ROUTING_KEY); } Bean public Binding deadLetterBinding() { return BindingBuilder.bind(orderDeadQueue()) .to(orderDeadExchange()) .with(DEAD_LETTER_ROUTING_KEY); } }这段代码里比较重要的是队列构造函数中的参数new Queue(ORDER_QUEUE, true, false, false, args)五个参数依次是队列名称、durable持久化、exclusive排他、autoDelete自动删除、arguments扩展参数。x-dead-letter-exchange和x-dead-letter-routing-key的作用是当消费者处理失败且拒绝消息时消息会转入死信队列由专门的程序处理或人工排查。4.4 定义消息体package com.example.rabbitmq.producer; import java.io.Serializable; import java.time.LocalDateTime; public class OrderMessage implements Serializable { private String orderId; private String userId; private BigDecimal amount; private LocalDateTime createTime; public OrderMessage() { } public OrderMessage(String orderId, String userId, BigDecimal amount) { this.orderId orderId; this.userId userId; this.amount amount; this.createTime LocalDateTime.now(); } // getter / setter 略请用 IDE 自动生成 public String getOrderId() { return orderId; } public void setOrderId(String orderId) { this.orderId orderId; } public String getUserId() { return userId; } public void setUserId(String userId) { this.userId userId; } public BigDecimal getAmount() { return amount; } public void setAmount(BigDecimal amount) { this.amount amount; } public LocalDateTime getCreateTime() { return createTime; } public void setCreateTime(LocalDateTime createTime) { this.createTime createTime; } }4.5 生产者Confirm Return 双回调生产者是这道场景题的重头戏。这里演示完整的 Confirm 和 Return 回调写法。package com.example.rabbitmq.producer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.rabbit.connection.CorrelationData; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import java.util.UUID; Component public class OrderMessageProducer { private static final Logger log LoggerFactory.getLogger(OrderMessageProducer.class); Autowired private RabbitTemplate rabbitTemplate; /** * 初始化回调逻辑 * 这里的回调可以单独封装成类这里为了演示简洁直接在初始化时设置。 */ PostConstruct public void init() { // Confirm 回调消息是否成功到达 RabbitMQ 服务端 rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (correlationData null) { return; } String messageId correlationData.getId(); if (ack) { log.info(消息确认成功messageId: {}, messageId); } else { log.error(消息确认失败messageId: {}, cause: {}, messageId, cause); // 实际生产环境这里需要记录失败日志或者将消息写入本地补偿表 } }); // Return 回调消息到达交换机但路由不到队列时触发 rabbitTemplate.setReturnsCallback(returned - { log.error(消息路由失败exchange: {}, routingKey: {}, replyText: {}, message: {}, returned.getExchange(), returned.getRoutingKey(), returned.getReplyText(), returned.getMessage()); }); } /** * 发送订单消息 */ public void sendOrderMessage(OrderMessage message) { CorrelationData correlationData new CorrelationData(UUID.randomUUID().toString()); // convertAndSend 方法会自动使用 RabbitConfig 中的 JSON 消息转换器 rabbitTemplate.convertAndSend( RabbitConfig.ORDER_EXCHANGE, RabbitConfig.ORDER_ROUTING_KEY, message, correlationData ); log.info(消息已发送消息ID: {}, correlationData.getId()); } }这段代码把消息发送的可靠性体现得很清楚CorrelationData中携带消息唯一 ID确认回调返回时可以通过这个 ID 判断是哪条消息。ConfirmCallback负责确认消息是否到达服务端。ReturnsCallback负责捕获路由失败的消息。注意Spring Boot 2.x 中setReturnCallback已经过时推荐使用setReturnsCallback。4.6 消费者手动 ACK 失败重试消费者使用手动 ACK 模式处理成功调用basicAck处理失败调用basicNack并让消息进入死信队列。package com.example.rabbitmq.consumer; import com.example.rabbitmq.config.RabbitConfig; import com.example.rabbitmq.producer.OrderMessage; import com.rabbitmq.client.Channel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; Component public class OrderMessageConsumer { private static final Logger log LoggerFactory.getLogger(OrderMessageConsumer.class); RabbitListener(queues RabbitConfig.ORDER_QUEUE) public void onOrderMessage(OrderMessage message, Channel channel, Message amqpMessage) throws IOException { long deliveryTag amqpMessage.getMessageProperties().getDeliveryTag(); try { log.info(收到订单消息orderId: {}, userId: {}, amount: {}, message.getOrderId(), message.getUserId(), message.getAmount()); // 模拟业务处理例如创建订单、扣库存、发通知等 boolean success processOrder(message); if (success) { // 手动确认只确认当前消息 channel.basicAck(deliveryTag, false); log.info(订单消息处理成功orderId: {}, message.getOrderId()); } else { // requeuefalse消息不重新入队转入死信队列 channel.basicNack(deliveryTag, false, false); log.warn(订单消息处理失败转入死信队列orderId: {}, message.getOrderId()); } } catch (Exception e) { log.error(消费订单消息异常orderId: {}, message.getOrderId(), e); // 捕获到未知异常拒绝消息并转入死信队列 channel.basicNack(deliveryTag, false, false); } } private boolean processOrder(OrderMessage message) { // 模拟业务失败场景金额大于 10000 视为处理失败 if (message.getAmount() ! null message.getAmount().compareTo(new BigDecimal(10000)) 0) { return false; } return true; } }这里有几个细节需要理解deliveryTag是消息投递的唯一标记确认时必须带上。basicAck(deliveryTag, false)的第二个参数表示是否批量确认这里只确认当前消息。basicNack(deliveryTag, false, false)的第三个参数表示是否重新入队这里选择不重新入队消息会进入死信队列。这样设计的好处是消费失败不会被无限重试也不会因为消费者抛异常就把消息直接丢回队列。4.7 死信队列消费者死信队列用于处理那些反复消费失败的消息。它可以把异常消息记录下来或者做人工补偿。package com.example.rabbitmq.consumer; import com.example.rabbitmq.config.RabbitConfig; import com.rabbitmq.client.Channel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; Component public class DeadLetterConsumer { private static final Logger log LoggerFactory.getLogger(DeadLetterConsumer.class); RabbitListener(queues RabbitConfig.DEAD_LETTER_QUEUE) public void onDeadMessage(Message message, Channel channel) throws IOException { String content new String(message.getBody()); try { log.error(收到死信消息原消息体: {}, 死信原因: {}, content, message.getMessageProperties().getHeader(x-first-death-reason)); // 这里可以发送告警、写入补偿表或者人工处理 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); } catch (Exception e) { log.error(处理死信消息异常消息体: {}, content, e); channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false); } } }x-first-death-reason是 RabbitMQ 自动添加的消息头用来标记消息第一次进入死信队列的原因可选值有rejected消费者拒绝、expired消息过期、maxlen队列长度溢出等。4.8 启动类与接口测试在启动类中注入生产者写一个简单的 HTTP 接口触发消息发送。package com.example.rabbitmq; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; SpringBootApplication public class RabbitmqReliableDemoApplication { public static void main(String[] args) { SpringApplication.run(RabbitmqReliableDemoApplication.class, args); } }创建一个发送消息的 Controllerpackage com.example.rabbitmq.controller; import com.example.rabbitmq.producer.OrderMessage; import com.example.rabbitmq.producer.OrderMessageProducer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import java.math.BigDecimal; RestController public class SendController { Autowired private OrderMessageProducer producer; GetMapping(/send) public String send(RequestParam String orderId, RequestParam String userId, RequestParam BigDecimal amount) { producer.sendOrderMessage(new OrderMessage(orderId, userId, amount)); return send success; } }启动项目后访问curl http://localhost:8080/send?orderId1001userIdu001amount99.5观察日志输出消息已发送消息ID: 5f7a1f2e-xxxx-xxxx-xxxx 消息确认成功messageId: 5f7a1f2e-xxxx-xxxx-xxxx 收到订单消息orderId: 1001, userId: u001, amount: 99.5 订单消息处理成功orderId: 1001如果发送amount20000的消息消费者会处理失败并拒绝消息你会在管理控制台的order.dead.queue中看到这条消息这就是死信队列转储的效果。5. 面试深度考点消息幂等与重复消费5.1 为什么会产生重复消息消息可靠性和消息不重复其实是两个相互矛盾的目标。当消费者处理完消息后还没来得及确认网络出现抖动RabbitMQ 没有收到 ACK会重新投递这条消息。这时候消息就被消费了两次。尤其是订单、支付、积分这类核心业务重复消费可能造成严重后果。5.2 消费端幂等方案面试中通常会继续追问你怎么保证消息不被重复消费常见的方案有三种按推荐程度排序方案一业务唯一 ID 数据库唯一约束在消息体中携带业务唯一 ID例如订单号。数据库表中对订单号建立唯一索引插入时如果重复会抛出 DuplicateKeyException此时直接返回成功。这是最稳妥的方案性能、可靠性和可解释性都很好。try { orderMapper.insert(order); } catch (DuplicateKeyException e) { log.warn(订单已存在跳过重复消费orderId: {}, order.getOrderId()); // 幂等处理成功正常 ACK }方案二Redis 去重利用 Redis 的SETNX命令记录消息 ID如果设置成功说明第一次消费如果设置失败说明已经消费过。Boolean first redisTemplate.opsForValue() .setIfAbsent(msg:order:1001, 1, Duration.ofHours(24)); if (Boolean.FALSE.equals(first)) { log.warn(重复消息直接返回成功消息ID: {}, messageId); return; }这种方式需要额外引入 Redis 依赖并且要考虑锁过期时间是否覆盖业务处理时长。方案三本地消息表在生产者本地建一张消息表发送消息前先写库消费者处理完成后回调标记状态。这个方案实现复杂度较高通常用于对可靠性要求极高的金融级场景。5.3 消费失败重试次数怎么控制很多同学会在消费者里写for循环重试这是不推荐的。更好的做法是用 Spring Retry 组件spring: rabbitmq: listener: simple: retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2设置最大重试 3 次重试间隔指数增加。超过次数后消息进入死信队列。这样既避免无意义的一直重试又能保留数据现场。6. 常见问题与排查思路下面把消息投递失联相关的高频问题整理成一张排查表面试时可以按这个逻辑回答。问题现象常见原因排查与解决思路生产者发送成功但队列中没有消息1. RoutingKey 拼写错误 2. 交换机绑定关系错误 3. 未开启 mandatory开启 Publisher Confirm 和 ReturnsCallback管理控制台查看绑定关系消费者收不到消息1. 队列没有绑定到对应交换机 2. 消费者监听队列名不同 3. 消息被其他消费者消费且调度不均控制台查看 Queue 的 Consumers 数量检查注解中的队列名考虑使用 Direct或Topic 精确路由消费者收到消息但一处理就报错1. 消息体反序列化失败 2. 业务代码异常导致自动 ACK 删除消息开启手动 ACK用 try-catch 包裹业务逻辑异常消息转入死信队列RabbitMQ 重启后消息消失1. 队列未持久化 2. 消息未设置 PERSISTENT创建队列时 durabletrue发送消息时设置持久化投递模式消息重复消费网络抖动导致 ACK 丢失RabbitMQ 重新投递在消费端做幂等唯一约束、Redis SETNX、本地消息表消费者处理速度很慢消息堆积消费者并发度不够或出现死循环调整 concurrency 参数检查是否有阻塞调用配合死信队列分开处理慢消息如果你在项目中遇到“消息凭空消失”的问题推荐按下面的顺序排查先打开管理控制台切换到 Queues 页面观察消息的 Ready 和 Unacked 数量。看生产者日志有没有 Confirm 回调失败、Return 回调触发。看消费者日志有没有异常有没有手动 ACK。检查交换机、队列、绑定关系是否一一对应。检查是否开启了持久化和手动确认机制。这五步走完大部分“失联”问题都能定位到具体环节。7. 最佳实践与工程建议7.1 消息可靠性的“三道保险”在生产环境我建议对核心消息链路做三道保险第一生产者开启 Confirm 模式并记录发送日志失败时写入补偿表。第二队列、交换机、消息全部持久化核心队列配置死信队列。第三消费者手动 ACK业务失败走重试重试失败进死信。这里要特别注意Confirm 确认成功并不代表业务处理成功它只表示 RabbbitMQ 接收了消息。业务成功与否完全取决于消费端是否正常 ACK。很多线上事故就是误解了“消息已发送成功”的含义。7.2 消息补偿表对于金融、订单这类强一致场景建议生产者在发送消息前先把消息写入本地业务表状态为“待发送”。Confirm 回调成功后更新为“已发送”失败则定时任务扫描重发。CREATE TABLE t_message_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, message_id VARCHAR(64) NOT NULL COMMENT 消息唯一ID, exchange_name VARCHAR(128) NOT NULL, routing_key VARCHAR(128) NOT NULL, message_body TEXT NOT NULL, status TINYINT NOT NULL DEFAULT 0 COMMENT 0-待发送 1-已发送 2-消费成功 3-发送失败, retry_count INT NOT NULL DEFAULT 0, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, UNIQUE KEY uk_message_id (message_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT消息发送记录表;定时任务每分钟扫描一次 status0 且 retry_count 5 的记录重新发送。这样即使 RabbitMQ 完全不可用消息也不会丢只是延迟送达。7.3 消费端监控很多消息丢失问题不是“一开始就丢”而是“丢了很久才发现”。所以监控很重要。监控队列积压数量Ready 消息数长期大于 0 说明消费速度跟不上。监控死信队列数量死信突然增加说明消费者频繁处理失败。监控 Confirm 失败率失败率突然升高说明 RabbitMQ 集群可能有问题。在 Spring Boot 中可以通过 Actuator 暴露 RabbitMQ 指标也可以对接 Prometheus Grafana。不管用什么工具原则是一样的可靠的消息链路必须有可观测性。7.4 队列选型普通队列还是仲裁队列RabbitMQ 3.8 之后引入了 Quorum Queue仲裁队列它基于 Raft 协议实现消息在多个节点间复制可靠性比普通持久化队列更高在节点故障时能自动选主避免数据丢失。如果你使用的是 RabbitMQ 3.9 以上版本且集群规模在 3 节点及以上可以优先考虑仲裁队列。但仲裁队列也有自己的特性比如不支持某些懒加载策略使用前需要做充分测试。单机开发环境仍然推荐使用普通队列。spring: rabbitmq: listener: simple: acknowledge-mode: manual单机环境手动 ACK 持久化 死信队列已经能覆盖绝大多数场景。7.5 配置最终检查项上线前可以照着这份清单检查一遍[ ] 交换机 durabletrue[ ] 队列 durabletrue[ ] 消息发送时设置持久化投递模式[ ] 生产者开启 publisher-confirm-type[ ] 生产者开启 mandatory[ ] 消费者关闭自动 ACK使用手动 ACK[ ] 消费失败有重试机制或死信队列[ ] 消费端有幂等处理[ ] 核心消息写入了补偿表8. 总结与面试话术再回到面试场景。当面试官问“RabbitMQ 消息投递失联你会怎么处理”时你不需要背一个固定答案但要展示完整的分析链路。推荐的回答结构是先说明消息链路分为三段生产者到交换机、交换机到队列、队列到消费者。再分阶段说明防护方案生产者 Confirm Return交换机队列消息持久化消费者手动 ACK。补充异常处理消费失败走死信队列。最后提业务兜底消息补偿表和消费幂等。这样回答说明你既有整体架构视角又知道具体落地的代码怎么写还能考虑生产环境的可观测性。面试官如果继续追问“消息大量堆积怎么办”“如何保证顺序消费”“RabbitMQ 节点宕机怎么办”你只需要沿着本文的链路逐层往下分析堆积问题看消费者并发度和队列拆分顺序问题看单队列消费和一致性哈希宕机问题看镜像队列和仲裁队列。这些知识点都可以在掌握本文的基础上继续深入。最后建议你动手做一件事把本文的 demo 代码跑起来故意发送一条路由键错误的消息打开管理控制台观察消息行为再故意让消费者抛异常观察死信队列的变化。只有亲手做过一遍你才能真正理解“消息失联”的每一个环节。如果本文对你有帮助欢迎收藏备用。也欢迎在评论区聊聊你遇到过的 RabbitMQ 消息丢失场景。
返回列表