消息队列如何保证消息不丢失:生产端确认、Broker 持久化与消费端手动 ACK
订单支付成功后发一条 MQ 消息,通知积分系统加积分——结果消息丢了,用户积分没到账,客诉 +1。消息丢失是 MQ 生产事故的高发区,而且它可能发生在三段链路中的任何一段:生产者发不出去、Broker 存不住、消费者没消费到。本文把每一段怎么丢、怎么防一次讲透。
先把链路拆开:消息一生要过三关
一条消息从业务发出到消费完成,经历:
第一关 生产者 -> Broker:网络抖动、发送失败、发送成功但回执丢失。
第二关 Broker 存储:写入内存后进程宕机没来得及落盘、单副本机器磁盘损坏、主从切换时 Leader 数据还没同步给 Follower。
第三关 Broker -> 消费者:消费者拉取后没处理完就提交进度、进程崩溃、业务异常被吞掉。
想保证"一条不丢",必须三关全部守住,任何一段裸奔都可能丢消息。
第一关:生产者如何确认真的发出去了
生产者的核心是同步等待确认 + 失败重试,绝不能"发完就当成功"。
Kafka:Producer 设置 acks=all,配合 min.insync.replicas=2,Leader 写入并等所有 ISR 副本确认后才返回成功;开启 retries 自动重试瞬时故障;关键业务建议同步 send().get() 或回调里检查异常并做补偿。
RocketMQ:用同步发送,检查 SendResult 的 sendStatus 是否为 SEND_OK;失败可重试,或记录待发送表由定时任务兜底补偿。
RabbitMQ:开启 Publisher Confirm 模式,channel.confirmSelect() 后每条消息 broker 会回 ack/nack;配合事务或备份交换机。
核心代码习惯(以 RocketMQ 为例):
SendResult result = producer.send(msg); // 同步发送,阻塞等回执
if (!result.getSendStatus().equals(SendStatus.SEND_OK)) {
// 发送失败:写入本地重试表,由补偿任务重发
retryService.save(msg);
}
再往前想一层:如果"本地业务提交了,但发消息失败",就会造成业务成功但消息缺失。这种"本地事务与发消息不一致"问题要靠事务消息(RocketMQ 半消息 + 回查)或本地消息表 + 定时补偿兜底——先写业务表和消息表(同一本地事务),再异步发消息,确认后改消息表状态,没发出去的定时扫描重发。
第二关:Broker 如何保证存得住
Broker 防丢三板斧:落盘 + 多副本 + 合理刷盘策略。
落盘(持久化):RocketMQ 有同步刷盘(SYNC_FLUSH)和异步刷盘(ASYNC_FLUSH),同步刷盘消息写入磁盘才返回,最稳但吞吐略降;Kafka 靠 OS 页缓存 + 定期 flush,看似"没写盘"其实是数据先进页缓存,由副本保证不丢。
多副本:RocketMQ 主从(同步双写/异步复制),Kafka 分区多副本 + ISR 机制。Kafka 关键配置:acks=all + min.insync.replicas=2,并且关闭 unclean.leader.election,防止落后太多的副本被选为 Leader 造成数据截断丢失。
副本同步模式:同步复制最稳但慢;异步复制快但主挂瞬间可能丢一小段。对账、幂等、可重放的业务消息可以接受异步复制 + 少量丢失窗口,但金融核心链路务必用同步或事务保障。
第三关:消费者如何确认真的消费完
消费端丢消息最常见的原因是提前提交消费进度:消息拉下来还没处理完,进程重启,offset 已经提交,Broker 认为消费过了,消息永久丢失。
正确姿势是手动 ack 且先业务后提交:
// 伪代码:处理业务成功后才提交消费位点
consumer.registerMessageListener((msgs, ctx) -> {
for (MessageExt msg : msgs) {
try {
businessService.doBiz(msg); // 1. 先执行业务
consumeService.markDone(msg); // 2. 记录处理成功(幂等标记)
} catch (Exception e) {
// 3. 业务失败:记录日志,稍后重试或进死信
// 不要在这里提交消费进度
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 4. 全部成功后提交
});
同时要注意:消费逻辑里异常不能静默吞掉(吞异常 = 假装成功),重试要有上限,超过上限进死信队列由人工或定时任务处理,避免消息无限重试阻塞后续消息。
防丢的终极答案:至少一次 + 幂等
要认清现实:绝大多数 MQ 只能保证至少一次(at-least-once)投递,即"不丢但可能重复"。想做到"既保证不丢、又保证只处理一次",正确做法不是让 MQ 恰好一次,而是:
可靠性兜底:生产者确认 + Broker 多副本落盘 + 消费者手动 ack——确保消息不丢。
业务幂等兜底:消费端用业务唯一键(订单号、消息 id)做去重(唯一索引 / Redis SETNX / 去重表),重复投递最多重复消费,但不会产生重复副作用。
两件事配合,才能对外宣称"消息可靠且业务只处理一次"。
分场景总结表
| 链路 | 风险点 | 防御手段 | 关键配置 |
|---|---|---|---|
| 生产者 | 发送失败/回执丢失 | 同步确认 + 重试 + 本地消息表补偿 | RocketMQ 同步发送;Kafka acks=all |
| Broker | 宕机未落盘/单副本损坏 | 刷盘策略 + 主从/多副本 + ISR | SYNC_FLUSH;min.insync.replicas=2 |
| 消费者 | 提前提交/异常吞掉 | 手动 ack + 先业务后提交 + 死信队列 | enable.auto.commit=false |
| 业务侧 | 重复投递产生重复副作用 | 唯一索引/去重表幂等 | 以业务单号做幂等键 |
面试常问"如何保证消息不丢失",记住这句答题主线:生产端确认不丢、Broker 副本不丢、消费端手动 ack 不丢,最后用幂等把重复消费变成无害的重复。