把消息发出去不代表消息一定能到达,RabbitMQ 默认的工作方式在多个环节都可能造成消息丢失:生产者发出的消息可能在网络传输中丢失,Broker 收到后如果没落盘,宕机重启消息就没了,消费者拿到消息还没处理完就挂掉,消息同样会被误认为已消费。要做到消息不丢,必须针对链路上的每个节点单独配置保障手段,其中消费端的手动 ACK 是最容易被忽视、也最容易踩坑的一环。本文把整条可靠性链路拆开来逐一实现。

一、先搞清楚消息在哪些环节会丢
一条消息从生产者到消费者要经历三个阶段,每个阶段都有丢失风险。第一阶段是生产者到 Exchange 的传输过程,RabbitMQ 默认采用“发后即忘”模式,消息发出去后生产者不知道 Broker 有没有收到,如果网络抖动或者 Exchange 不存在,消息就无声无息地消失了。第二阶段是 Exchange 到 Queue 以及 Queue 内部存储,默认情况下消息保存在内存里,一旦 RabbitMQ 重启,未持久化的消息全部清空;另外如果路由键写错,消息会被 Exchange 直接丢弃,连日志都很难查到。
第三阶段是消费者处理环节,这是最隐蔽的一种丢失。RabbitMQ 默认开启自动 ACK,消费者一拿到消息,队列就立刻把消息删除,至于消费者的业务代码有没有执行成功,Broker 完全不关心。假如消费者收到消息后还没写数据库就抛了异常,这条消息已经从队列里消失,业务数据就永久丢了。所以可靠性投递的完整方案必须三段齐抓:生产端用 Confirm 机制确认到达,Broker 端开启持久化,消费端改手动 ACK 并配合重试。
二、生产端 Confirm 确认与消息持久化配置
生产端的可靠性依赖两个回调:ConfirmCallback 用来确认消息是否成功到达 Exchange,ReturnCallback 用来捕获路由失败的情况(消息到了 Exchange 但找不到匹配的 Queue)。注意 Return 回调生效有个前提,必须设置 mandatory=true,否则路由失败的消息会被直接丢弃。下面是 Spring Boot 中的完整配置。
@Configuration
public class RabbitConfig {
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory factory) {
RabbitTemplate template = new RabbitTemplate(factory);
// 开启 mandatory,路由失败时触发 ReturnCallback 而不是静默丢弃
template.setMandatory(true);
// 消息成功或失败到达 Exchange 时触发
template.setConfirmCallback((correlationData, ack, cause) -> {
if (ack) {
System.out.println("消息已到达 Exchange: " + correlationData.getId());
} else {
// ack 为 false 说明 Exchange 没收到,需要记录并补偿重发
System.out.println("消息发送失败: " + cause);
// 实际项目中写入补偿表,由定时任务重发
}
});
// 消息无法路由到任何队列时触发
template.setReturnsCallback(returned -> {
System.out.println("路由失败: " + returned.getMessage()
+ " replyText: " + returned.getReplyText());
});
return template;
}
// 队列必须声明为持久化,durable=true
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue").build();
}
// 交换机同样要持久化
@Bean
public DirectExchange orderExchange() {
return ExchangeBuilder.directExchange("order.exchange").durable(true).build();
}
}application.yml 里还要打开 publisher-confirm,以及开启发布退回:
spring:
rabbitmq:
host: 192.168.0.1
port: 5672
username: admin
password: admin
publisher-confirm-type: correlated # 异步确认,性能好且能拿到回执
publisher-returns: truepublisher-confirm-type 有三个可选值:none 表示关闭确认;simple 是同步阻塞等待回执,吞吐量很低,一般不用;correlated 是异步回调方式,推荐生产环境使用。发送消息时记得带上 CorrelationData 作为唯一标识,确认回调里靠它定位是哪条消息失败,否则失败重发就无从谈起。发消息时也要设置消息本身持久化:MessageProperties.setDeliveryMode(MessageDeliveryMode.PERSISTENT),使用 RabbitTemplate 的 convertAndSend 时默认就是持久化,但如果手动构造 Message 就要格外留意别漏掉。
三、消费端手动 ACK 配置与三种确认方式
消费端是消息安全的最后一道闸门。开启手动 ACK 只需要两步:配置文件里把 acknowledge-mode 改成 manual,然后监听方法里注入 Channel 参数,在业务逻辑执行完之后手动调用确认方法。
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # 手动确认,默认是 auto 自动确认
prefetch: 1 # 一次只取一条,处理完再取下一条
@Component
public class OrderConsumer {
@RabbitListener(queues = "order.queue")
public void handle(String content, Message message, Channel channel) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
// 业务处理:写库、调用下游服务等
processOrder(content);
// 处理成功,确认消息,false 表示只确认当前这条
channel.basicAck(tag, false);
} catch (Exception e) {
// 处理失败,拒绝消息
// 第二个参数 requeue:true 表示重回队列,false 表示丢弃或进入死信队列
channel.basicNack(tag, false, false);
}
}
}手动 ACK 有三种确认手段,语义各不相同。basicAck 表示消费成功,Broker 可以删除消息;basicNack 是否定确认,支持批量拒绝,第三个参数 requeue 决定消息是重回队列还是进入死信;basicReject 与 basicNack 类似但不支持批量,一次只能拒绝一条。这里要特别提醒一个经典死循环陷阱:如果消息本身有问题(比如报文格式错乱导致永远解析失败),又设置 requeue=true,消息会立刻重回队列被再次投递,再次失败再次重回,消费者会在毫秒级循环中打满 CPU。正确做法是限制重试次数,超过阈值后 requeue 设为 false,配合死信队列承接这些处理不了的消息,人工排查后补偿。
更稳妥的重试控制可以借助消息头计数:每次 nack 前检查 x-death 头里的重试次数,或者干脆在消费者里自己维护重试计数,超限就 ack 掉当前消息并转存到失败表中。另外别忘了 prefetch 的设置,它的含义是消费者最多预取多少条未 ack 的消息,设太大一旦消费者崩溃,这批消息都要重新投递,设成 1 能保证消息绝对安全但吞吐量下降,需要根据业务在可靠性和性能之间权衡。
四、幂等消费与死信队列兜底
开启手动 ACK 之后消息基本不会丢,但代价是可能重复投递:消费者处理完业务、还没来得及 ack 就宕机,Broker 会把消息重新发给其他消费者,业务逻辑就会执行两次。所以消费端必须做幂等处理。常用方案有几种:利用数据库唯一索引,把消息 ID 作为唯一键,重复插入直接失败;用 Redis 的 setnx 记录消息 ID,设置合理过期时间;或者维护一张已消费消息表,处理前先查一遍。生产环境推荐数据库唯一索引方案,因为它和业务操作在同一个事务里,天然原子,不会出现幂等记录写成功但业务失败的错位情况。
死信队列是整个链路的兜底方案。消息被 basicReject 或 basicNack(requeue=false)、消息过期、队列长度超限这三种情况都会变成死信,如果队列声明时配置了死信交换机,死信会被转发到指定队列,方便后续人工处理或定时补偿。配置方式如下:
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
// 指定死信交换机
.deadLetterExchange("dlx.exchange")
// 死信路由键
.deadLetterRoutingKey("dlx.order")
.build();
}
@Bean
public Queue deadLetterQueue() {
return QueueBuilder.durable("order.dlx.queue").build();
}把这几个环节串起来看,完整的可靠性链路是:生产端 ConfirmCallback 确认消息到达 Exchange,ReturnCallback 拦截路由失败,消息和队列双重持久化防止 Broker 宕机丢数据,消费端手动 ACK 保证处理成功才确认,处理不了的消息进死信队列,配合幂等消费消化重复投递。每一层都不是可选项,少配任何一层,消息就可能在那个环节悄悄消失。上线前建议用杀消费者进程和重启 RabbitMQ 两个场景做演练,验证消息确实不丢、不重,再对外承诺可靠性指标。