1. 消息丢失发生在哪三个阶段?怎么保证不丢?

三阶段全链路分析:

  1. 生产阶段丢失:
    • 原因:网络异常发送失败、发送后 Broker 未确认
    • 解决:同步发送 + 捕获异常重试;Kafka 设 acks=all(见 Kafka 页);发消息前先落库/本地消息表(终极兜底:消息+业务同一事务,定时任务扫描重发)
  2. Broker 存储阶段丢失:
    • 原因:Broker 宕机时数据未刷盘/未同步副本
    • 解决:多副本 + acks=all + min.insync.replicas;刷盘参数合理(Linux 上刷盘交给 OS 也行,但要求别丢就 flush)
  3. 消费阶段丢失:
    • 原因:先提交 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、网络重试都会导致同一条消息被消费多次。重复是常态,必须幂等设计。

幂等方案(按场景):

  1. 唯一约束/唯一索引:消息带业务唯一键(订单号),插入冲突则跳过——最可靠
  2. Redis SETNX:以消息 ID 为 key 去重(带 TTL,注意分布式原子性)
  3. 状态机校验:处理前检查业务状态(如订单已支付则跳过)
  4. 数据库乐观锁: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 的并行(多分区/多消费者)破坏顺序。保证方案分三层:

  1. 生产端:分区有序——同一业务键(如 orderId)的消息发到同一分区(key 哈希路由),分区内天然有序(Kafka)
  2. 消费端:单消费者串行——一个分区只被一个消费者处理(组内分配保证),且该消费者单线程顺序处理(不要内部开线程池并发处理同分区)
  3. RocketMQ 专门方案:MessageQueueSelector 按业务键选队列 + 消费端顺序消费(串行)

🎯 面试要点

  • 顺序 vs 并行的取舍:全局顺序 = 全局单线程 = 吞吐低;通常只需要局部顺序(按订单/用户维度),保留并行度
  • 重试会打乱顺序:失败消息要特殊处理(挂起后续或阻塞重试),别直接重投

4. 消息堆积了怎么排查和解决?

判断:消费速度 < 生产速度,offset 落后越来越大(Kafka 看 lag;RocketMQ 看积压数)。

排查:

  1. 消费者是否正常(进程挂/卡死/频繁 Rebalance)
  2. 消费逻辑是否变慢(慢 SQL、外部依赖超时、死循环)——看消费者耗时指标
  3. 是否分区数太少(并发不足)

解决:

  1. 提升单机消费能力:优化消费逻辑(并行化内部处理、批量操作)
  2. 扩容消费者:加消费者实例 + 分区数(Kafka 消费者数 ≤ 分区数)
  3. 临时降级:堆积严重时丢弃低价值消息(先保实时性)或写离线存储后慢慢补
  4. 根治:评估消费瓶颈(DB 压力?外部调用?)

🎯 面试要点

  • 监控 lag 是 MQ 运维核心指标:lag 突然上涨 = 报警信号
  • 区分"堆积"与"延迟":消息有 TTL/延迟队列的积压是设计的(削峰场景)

🎤 常见面试追问

  1. 消息丢失的三阶段分析(必背)?——生产阶段:同步发送+重试、acks=all、本地消息表兜底;存储阶段:多副本+min.insync.replicas;消费阶段:先处理业务成功再提交 offset。
  2. 为什么会重复消费?——at-least-once 语义下:消费成功但 offset 提交失败、Rebalance、网络重试。重复是常态 → 业务必须幂等(唯一索引/状态机/Redis SETNX)。
  3. 本地消息表和事务消息的区别?——本地消息表:业务+消息记录同一本地事务,定时任务重发(自己实现);事务消息:RocketMQ 内置半消息+回查(框架实现)。本质同思路,后者更优雅。
  4. 消息堆积了怎么办?——先查原因:消费者挂/慢/频繁 Rebalance?再解决:优化消费逻辑、扩容消费者+分区、临时降级(低价值消息丢弃或离线补)。
  5. exactly-once 能做到吗?——精确一次 = at-least-once + 幂等消费(业务层去重)。Kafka 的 exactly-once 是流处理层面(事务性生产者),消费侧仍要靠幂等。

📖 名词解释(本页术语)

术语 大白话解释
at-most-once / at-least-once / exactly-once投递语义三档:最多一次(可能丢)/ 至少一次(可能重复,主流)/ 精确一次(不丢不重,代价大)。
幂等消费同一消息消费多次结果一样——用业务唯一键去重(唯一索引/状态机/Redis SETNX)。
本地消息表业务写库 + 消息记录同一事务,定时任务扫描重发——保证"消息一定发出"的兜底方案。
死信队列(DLQ)消费重试多次仍失败的消息送进"死信",人工/脚本处理——可靠性的最后一环。
消息堆积(lag)消费速度跟不上生产,积压消息——监控 lag 是 MQ 运维核心指标。
顺序消息按业务键路由到同一队列/分区 + 消费端串行处理——保证相关消息按序消费。
重试 / 补偿消费失败自动重投(带次数上限);补偿 = 业务层面的反向操作(对账、人工处理)。
⚠️ 本页面由 AI 生成,内容仅供参考,请以官方文档和实际源码为准。