🛡️ 消息可靠性
丢失/重复/顺序/堆积——四大经典问题与解决方案
1. 消息丢失发生在哪三个阶段?怎么保证不丢?
三阶段全链路分析:
- 生产阶段丢失:
- 原因:网络异常发送失败、发送后 Broker 未确认
- 解决:同步发送 + 捕获异常重试;Kafka 设 acks=all(见 Kafka 页);发消息前先落库/本地消息表(终极兜底:消息+业务同一事务,定时任务扫描重发)
- Broker 存储阶段丢失:
- 原因:Broker 宕机时数据未刷盘/未同步副本
- 解决:多副本 + acks=all + min.insync.replicas;刷盘参数合理(Linux 上刷盘交给 OS 也行,但要求别丢就 flush)
- 消费阶段丢失:
- 原因:先提交 offset 后处理业务——处理中崩溃,offset 已提交,消息不再投递
- 解决:先处理业务,成功后再提交 offset(at-least-once 语义);消费失败重试,重试失败进死信队列
本地消息表(终极不丢方案)
// ① 业务写库 + 消息状态表 同一本地事务
@Transactional
public void createOrder(Order o) {
orderDao.insert(o); // 业务
msgDao.insert(new MsgRecord("order-created", o.id, STATUS_PENDING)); // 消息记录
}
// ② 定时任务扫描 PENDING 消息发送,成功改 SENT(幂等重发)
// 保证:消息一定发出(最多重复,不丢)
🎯 面试要点
- 分三段答(生产/存储/消费)是最佳结构,每段答"原因 + 对策"
- at-most-once(可能丢)/ at-least-once(可能重复)/ exactly-once(精确一次,代价大)——先说明语义再选方案
2. 消息重复消费如何解决?(幂等)
为什么重复:at-least-once 语义下——消费成功但提交 offset 失败、Rebalance、网络重试都会导致同一条消息被消费多次。重复是常态,必须幂等设计。
幂等方案(按场景):
- 唯一约束/唯一索引:消息带业务唯一键(订单号),插入冲突则跳过——最可靠
- Redis SETNX:以消息 ID 为 key 去重(带 TTL,注意分布式原子性)
- 状态机校验:处理前检查业务状态(如订单已支付则跳过)
- 数据库乐观锁:UPDATE ... WHERE status = 旧状态,影响 0 行则说明已处理
幂等消费模板
public void onMessage(Msg msg) {
String bizKey = msg.getBizId(); // 业务幂等键
int rows = orderDao.updateStatusIfPending(bizKey, "PAID", "PENDING");
if (rows == 0) return; // 已处理过,幂等跳过
// 继续业务...
}
🎯 面试要点
- 幂等的本质:用业务键的约束代替"只消费一次"的假设
- 消息 ID 幂等 vs 业务键幂等:消息 ID 只防"同一条消息重复",业务键防"重复请求"——业务键更本质
3. 如何保证消息的顺序性?
顺序性来源:MQ 的并行(多分区/多消费者)破坏顺序。保证方案分三层:
- 生产端:分区有序——同一业务键(如 orderId)的消息发到同一分区(key 哈希路由),分区内天然有序(Kafka)
- 消费端:单消费者串行——一个分区只被一个消费者处理(组内分配保证),且该消费者单线程顺序处理(不要内部开线程池并发处理同分区)
- RocketMQ 专门方案:MessageQueueSelector 按业务键选队列 + 消费端顺序消费(串行)
🎯 面试要点
- 顺序 vs 并行的取舍:全局顺序 = 全局单线程 = 吞吐低;通常只需要局部顺序(按订单/用户维度),保留并行度
- 重试会打乱顺序:失败消息要特殊处理(挂起后续或阻塞重试),别直接重投
4. 消息堆积了怎么排查和解决?
判断:消费速度 < 生产速度,offset 落后越来越大(Kafka 看 lag;RocketMQ 看积压数)。
排查:
- 消费者是否正常(进程挂/卡死/频繁 Rebalance)
- 消费逻辑是否变慢(慢 SQL、外部依赖超时、死循环)——看消费者耗时指标
- 是否分区数太少(并发不足)
解决:
- 提升单机消费能力:优化消费逻辑(并行化内部处理、批量操作)
- 扩容消费者:加消费者实例 + 分区数(Kafka 消费者数 ≤ 分区数)
- 临时降级:堆积严重时丢弃低价值消息(先保实时性)或写离线存储后慢慢补
- 根治:评估消费瓶颈(DB 压力?外部调用?)
🎯 面试要点
- 监控 lag 是 MQ 运维核心指标:lag 突然上涨 = 报警信号
- 区分"堆积"与"延迟":消息有 TTL/延迟队列的积压是设计的(削峰场景)