Skip to content

第四章 消息队列(第79-96题) ​

79. 架构评审:引入 MQ 的理由与代价(MQ 价值与代价) ​

【考察内容】MQ 核心价值与引入代价(解耦/异步/削峰)

【题目】架构评审会上,你提出要引入消息队列,有人反对说“多一个组件多一份故障”。请讲清楚:引入 MQ 到底解决什么痛点(削峰/解耦/异步),不引入会怎样,以及引入后新增了哪些风险与成本?

【参考答案】三大核心作用:

  1. 解耦:A 系统产生订单事件,需要通知 B/C/D 系统——同步调用会“新增一个下游就改一次代码”;通过 MQ,A 只发消息,下游自行订阅,新增下游不动 A 的代码;
  2. 异步:下单链路同步调用“发短信、送积分、更新推荐”导致接口 RT 长——MQ 异步化后主链路只做核心操作,RT 从 2s 降到 200ms;
  3. 削峰:秒杀瞬间 10 万请求,直接打到 DB 会挂——MQ 缓冲,消费端按 DB 可承受的速率慢慢消费(如每秒 1000),保护下游;
  4. 代价:引入新组件(运维成本)、消息可靠性问题(丢失/重复/乱序需要治理)、最终一致性(业务要接受异步延迟)。
  5. 容量估算(步进):大促峰值 10 万下单请求/10s ≈ 1 万 QPS 瞬时写入;DB 稳态只能扛约 1000–2000 TPS,需 MQ 削峰约 5–10 倍;积压缓冲要先判过载形态,两种算法差两个量级:持续过载(流量长期高于处理能力,如大促 30 分钟)→ 缓冲 ≈(到达率−处理率)×持续时长=(1 万−1500)×1800 ≈ 1530 万条;瞬时尖峰(总量有限、峰后即回落,如题干 10 万请求压在 10 秒内)→ 缓冲上限=突发总量−期间已处理量≈ 10 万−1500×10 ≈ 8.5 万条,约 60–100 秒即可排空。本题是尖峰型:误按 30 分钟持续算会把缓冲高估约 160 倍、照此扩容就是浪费;工程上取 min((到达率−处理率)×持续时长, 突发总量),再按消息体折算磁盘与内存。异步化后主链路 RT:串行慢下游 2s → 只保留本地事务 200ms 量级;通知/积分类 SLA 常约定秒级,对账 T+0/T+1。
  6. 失败与降级:MQ 集群不可用时,下单主事务不依赖 MQ 成功与否——本地消息表/事务消息暂存,恢复后补发——但降级通道不能把尖峰原样压回同一套 DB:本库对 MySQL 单机的口径是写 2000–5000 TPS,而 MQ 挂时原本要被异步吸收的瞬时写入可达 1 万 QPS(5–10 倍超载),照此答等于用降级打垮 DB。降级三件套要同时给出:① 只留核心事件、② 本地表批量合并写、③ 入口按 DB 实测余量限流(与第 91 题同口径);非核心副作用(推荐、运营通知)直接关闭;生产端发送失败写本地重试表并告警,禁止静默丢弃;消费端失败进死信人工介入。引入前必须书面确认:哪些业务可接受最终一致、延迟窗口多长、降级开关负责人。

【原理溯源】

  • 为什么解耦必须靠“事件广播”而不是接口调用? 同步调用是编译期耦合:调用方必须知道下游地址、协议、签名,新增下游 = 改上游代码并重新发布。MQ 把“调用”变成“发布事实”:生产者只声明“发生了什么”,谁关心谁订阅。上游与下游在时间、空间、接口三维度解耦——时间上不必同时在线,空间上不必互相知道地址,接口上只约定消息 schema。
  • 为什么异步能显著压 RT? 同步链路的 RT = 核心操作 + 所有下游串行/并行耗时的最慢者。发短信(外部 SMS 网关 300ms)、积分(跨服务 50ms)、推荐更新(重计算 800ms)都挂在主链路上时,用户等的是整条链。异步化后主链路只保留“落库成功”这一正确性关键路径,其余变成可延迟的副作用——RT 由“最慢下游”决定变成由“本地事务”决定。
  • 削峰的本质是什么? 把突发写入速率与稳定处理速率解耦。DB 的稳态吞吐是相对固定的(连接池、索引、锁竞争),瞬时 10 倍流量不是“处理慢一点”而是直接雪崩。MQ 以磁盘/内存缓冲吸收峰值,消费端按自身容量拉取——代价是端到端延迟变成可变量,积压时延迟线性上涨。
  • 引入 MQ 新增了哪些不可回避的复杂度? ① 可靠性从“同步返回码”变成“分段确认”(生产/Broker/消费三段都要治理);② 一致性从 ACID 变成最终一致(业务要接受“支付成功但积分晚几秒”);③ 顺序性从“调用顺序天然保证”变成“需要按 key 路由 + 单线程消费”;④ 运维多一类有状态中间件(磁盘、副本、Rebalance、监控)。这些不是“配置一下就好”,而是要落到代码与预案里。
  • 什么场景不该上 MQ? 强一致且同步可完成的操作(扣款余额)、实时性要求毫秒级的查询、调用链只有 1 个下游且低频、需要同步拿到下游结果才能继续的分支逻辑——这些场景上 MQ 是为了异步而异步,增加故障面却不换来收益。

【选型判断树】

要不要引入 MQ?
├─ 是否存在“生产快于消费”的瞬时峰值?
│   └─ 是(秒杀/推送/日志洪峰)→ 削峰价值成立
├─ 下游是否 ≥2 个且还会继续增加?
│   └─ 是 → 解耦价值成立
├─ 下游操作是否可延迟(秒级/分钟级可接受)?
│   └─ 是 → 异步价值成立
├─ 以上三条全否 → 不要上 MQ(本地调用更简单)
└─ 以上命中 ≥1 条
    ├─ 强一致 + 必须同步拿到结果 → 同步 RPC,MQ 只做旁路通知
    └─ 允许最终一致
        ├─ 需要事务消息/金融级可靠 → RocketMQ
        ├─ 超高吞吐日志/大数据 → Kafka
        └─ 中小流量复杂路由 → RabbitMQ

判断口诀:先问“峰值、多下游、可延迟”三问,全否就别上;上了就要预算可靠性治理成本。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“这题核心是权衡:MQ 用最终一致和运维复杂度,换解耦、异步、削峰三类收益”
0:30–2:30三价值展开每条价值配一个“不引入会怎样”的反例:新增下游改代码、RT 被慢下游拖死、峰值打挂 DB
2:30–3:30代价清单主动列:可靠性三段治理、最终一致、顺序问题、运维成本——主动谈代价是加分项
3:30–4:30边界场景说明什么场景不该上:强一致同步操作、毫秒级实时、单下游低频
4:30–5:00收尾“一句话:MQ 不是银弹,是用‘异步 + 最终一致’换‘扩展性 + 抗峰值’的架构交易”

【关键数字】

参数经验值说明
同步 RT 优化幅度2s → 200ms 量级异步化砍掉慢下游后的典型落差
削峰消费比峰值:稳态 ≈ 5–10:1秒杀常见;消费端按 DB 容量定速
MQ 单机吞吐量级Kafka 十万级;RocketMQ 十万级;RabbitMQ 万级选型时的粗粒度锚点
异步延迟可接受窗口积分/通知 秒级;对账 T+0/T+1业务必须书面确认 SLA
引入后故障面多 1 个有状态组件需独立监控、备份、演练
峰值:稳态消费比5–10:1秒杀常见;MQ 缓冲 + 定速消费
削峰后消费速率按 DB/下游容量反推如 DB 1000 TPS → 消费端限流 1000
缓冲积压容量(先分过载形态,两式差两个量级)持续过载:(到达率−处理率)×持续时长;瞬时尖峰:突发总量−期间已处理量;工程取 min(两式)秒杀属尖峰型(≈8.5 万条),误按持续型×30 分钟会高估约 160 倍;磁盘/保留策略按算出的值装下

【追问链】(三层)

L1|“什么场景不适合 MQ?” → 三类:① 强一致且必须同步完成(扣余额);② 实时性要求极高(毫秒级查询/推送确认);③ 简单调用、单下游、低频——上 MQ 是过度设计。核心判断:下游结果是否阻塞主流程、延迟是否业务可接受。

L2|“异步后用户点了支付,怎么知道成功了?” → 主链路同步保证核心正确性(订单落库、支付回调处理成功)立即返回;副作用(积分、通知)异步。体验上返回“处理中/成功”,提供订单状态查询与推送;资金类以支付渠道回调+订单状态机为准。面试可补数字:若下游积分 SLA 为 5s,用户无感;若超 30s 未到账触发补偿对账。

L3|“MQ 整个集群挂了,下单还能做吗?” → 必须能。方案:① 核心链路与 MQ 解耦——下单主事务不依赖 MQ 是否可用;② 降级:本地消息表暂存、MQ 恢复后补发(见 91 题)——但本地表仍写同一套 DB,必须同时「只留核心事件+批量合并写+入口按 DB 实测余量限流」,否则 5–10 倍尖峰会把降级做成二次故障;③ 非核心异步功能直接关闭。原则:MQ 是增强可用性与扩展性的旁路,不能成为下单成功与否的同步依赖。

【评分标准】

档位答案特征
60 分能背出解耦/异步/削峰三个词,各给一句定义
80 分每条价值都有“不引入会怎样”的反例;能主动列出可靠性/最终一致/运维三类代价
95 分说清削峰=速率解耦、解耦=三维度解耦;给出“不该上 MQ”的边界;MQ 挂掉时的降级要说三件套(本地消息表+只留核心事件/批量合并写+入口按 DB 余量限流),只答「写本地消息表」等于把尖峰压回同一套 DB

【关联题】

  • 同簇: 第 80 题(不丢)→ 第 81 题(不重)→ 第 82 题(有序)→ 第 94 题(ack)——MQ 可靠性四件套
  • 选型落地: 第 84 题(三 MQ 对比)
  • 降级兜底: 第 91 题(MQ 故障降级)、第 93 题(本地消息表)
  • 上游理论: 第 98 题(BASE)——异步最终一致的理论基础

【自测】

  1. 判断对错:只要 QPS 高,就该上 MQ 削峰。 参考答案: 错。若消费端能力与生产端匹配、无突发峰值、且需要同步结果,上 MQ 只增加复杂度。削峰只解决“峰值 ≫ 稳态”问题。
  2. 架构评审时,反对者说“多一个组件多一份故障”,你怎么回应? 参考答案: 承认代价,再算收益:MQ 把故障域从“全链路同步串联”变成“主链路短 + 异步可重试”;同时给出可靠性三段防护与降级预案,证明引入后整体可用性可更高,而不是简单多一个单点。
  3. 为什么说解耦是“发布事实”而不是“调用接口”? 参考答案: 发布事实只声明事件与 schema,不绑定下游地址与存在性;新增订阅方无需改生产者。接口调用则把下游身份硬编码进上游,形成编译期耦合。

80. 业务方投诉“消息丢了”,订单状态没推进(消息不丢失) ​

【考察内容】消息不丢失是 MQ 最高频题

【题目】业务方投诉:异步任务经常“消息丢了”,比如订单支付成功但积分没到账。消息在“生产者发送→Broker 存储→消费者消费”三个阶段都可能丢失。请给出每个阶段不丢的配置与机制(ack、持久化、重试、确认机制)。

【参考答案】逐环节防护:

  1. 生产端:
    • Kafka:acks=all(等所有 ISR 副本确认)+ 重试机制;同步发送并在失败时记录告警;
    • RabbitMQ:开启 publisher-confirm(Broker 返回 ack/nack,nack 重发);
    • RocketMQ:同步发送(syncSend),不丢才返回成功;
  2. Broker 端:
    • 持久化:Kafka topic 副本数≥2 + min.insync.replicas=2 + unclean.leader.election 关闭(缺 mir 时 acks=all 只需 leader 落盘即算成功,leader 崩仍丢已 ack 消息);RabbitMQ 队列持久化 + 消息 delivery_mode=2;RocketMQ 刷盘策略(同步刷盘更强);
    • 集群高可用(副本机制),防止单节点宕机丢数据;
  3. 消费端:
    • 手动 ack(而不是自动):业务处理成功后才提交 offset/ack,处理失败不 ack,消息重新投递;
    • 关闭“收到即确认”(RabbitMQ autoAck=false);
  4. 兜底:发送失败写本地重试表/MQ 自身重试;消费失败进死信队列人工处理;对账(业务侧核对消息与结果)。
  5. 容量估算(步进):日订单 100 万、每单平均 3 条异步消息 → 日消息约 300 万条,均值 ≈ 300万/86400 ≈ 35 msg/s,大促峰值 ×50 ≈ 1750 msg/s(约 0.18 万)。若单条消息 1KB、保留 3 天:存储 ≈ 300万×1KB ≈ 3GB/日,保留 3 天约 9GB,再乘副本×3 约 27GB(磁盘与保留策略要按这个量预留)。生产端 acks=all 同步确认时,单分区写入 RT 约 2–10ms,吞吐低于 acks=1,但换来不丢。
  6. 失败与降级:生产端 confirm 失败 → 重试 N 次后落本地重试表/告警,禁止业务假装成功;Broker 单副本故障 → 副本接管,若 ISR 为空宁可拒绝写也不丢已确认数据;消费失败 → 不 ack,重试后进死信;业务方仍报丢 → 启动对账(消息流水 vs 业务结果),差错走补偿任务。绝对不丢不可达,工程目标是“可发现 + 可补偿”。

【原理溯源】

  • 为什么“消息会丢”不是 bug 而是默认设计? MQ 默认走性能优先路径:生产者 fire-and-forget、Broker 异步刷盘、消费者 auto-ack——任一环节“先确认再干活”或“不确认”都会在进程崩溃时丢数据。不丢是配置 + 代码约定换来的,不是开箱即用。
  • 生产端丢消息的根因是什么? 发送方把“写入本地缓冲/网络发出”当成成功。网络分区、Broker 未落盘、leader 切换时,消息只存在于发送方内存或已发出但未持久化。所以必须有同步确认:Kafka acks=all 等 ISR 全部写入;RabbitMQ publisher-confirm 等 Broker 回调;RocketMQ syncSend 等写入结果。失败要重试 + 告警,不能静默吞掉。
  • Broker 端为什么会丢? ① 未持久化:消息只在 page cache/内存,进程崩溃即丢;② 副本不足:leader 挂且 follower 数据落后;③ 非干净选举:允许落后副本成为 leader,丢已确认消息。对策是“写盘策略 + 副本数 + 禁止 unclean election”三件套。Kafka 关闭 unclean.leader.election 后,ISR 空时宁可不可用也不选落后副本。
  • 消费端丢消息的根因是什么? auto-ack / 过早提交 offset。消费者收到消息就确认,但业务处理中崩溃——Broker 以为消费成功,消息不再投递。正确姿势:处理成功后再 ack/commit offset;失败不 ack,依赖重投。代价是可能重复(at-least-once),需幂等(见 81 题)。
  • 为什么理论上仍难做到“绝对不丢”? 任一确认信号都可能在“已持久化但 ack 丢失”或“已 ack 但下游业务失败”之间存在窗口。工程上用事务消息/本地消息表消掉生产侧窗口(92/93 题),用对账消掉消费侧与业务侧不一致——所以答案永远是“三段防护 + 对账兜底”,而不是“配一个参数就完事”。

【选型判断树】

用户报“消息丢了”
├─ 先定位丢失段
│   ├─ 生产日志无发送失败、Broker 无消息 → 生产端未开 confirm/acks,或异步发送失败被吞
│   ├─ Broker 有消息、消费者无拉取记录 → 消费未启动/订阅错误/Rebalance 异常
│   └─ 消费者拉到了但业务未生效 → auto-ack 过早 / 处理中崩溃 / 业务异常未回滚
├─ 配置核对清单
│   ├─ 生产:Kafka acks=all + 重试;Rabbit confirm;RocketMQ sync
│   ├─ Broker:持久化 + 副本≥2 + 同步刷盘(金融) + 禁 unclean 选举
│   └─ 消费:手动 ack + 处理成功再提交 + 死信兜底
└─ 最终兜底:对账任务(消息表 vs 业务结果)发现静默丢失

口诀:三段都要确认;确认越晚越安全,但重复越多——不丢与不重必须成对设计。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“消息丢失必须按生产、Broker、消费三段拆,每段有独立的确认机制”
0:30–2:30三段展开每段:怎么丢的因果 + 对应 MQ 的关键配置(acks/confirm/持久化/手动 ack)
2:30–3:30因果闭环强调:越安全的确认越慢、越容易重复,所以必须接幂等
3:30–4:30兜底死信、重试表、对账——承认“极端仍可能丢”,工程靠对账发现
4:30–5:00收尾“不丢 = 三段确认 + 兜底;不丢与不重是一对 trade-off”

【关键数字】

参数经验值说明
Kafka acks0 / 1 / allall 最安全;1 在 leader 故障时可能丢
Kafka 副本数≥2,金融 ≥3ISR 同步副本
RabbitMQ delivery_mode2(persistent)配合队列持久化
RocketMQ 刷盘同步刷盘 > 异步刷盘同步更安全,吞吐更低
消费确认时机业务成功后auto-ack 仅适合可容忍丢失的日志类
对账周期T+0 准实时或 T+1资金/积分建议 T+0
日消息量估算日订单×每单消息数100万单×3 ≈ 300万条/日
峰值 msg/s均值×峰值系数均值35 → 大促×50 ≈ 1750(约0.18万)
存量存储日量×消息体×保留天×副本注意磁盘水位与保留策略

【追问链】(三层)

L1|“开了 acks=all 是不是就绝对不丢?” → 生产端接近不丢,但 Broker 宕机、消费 auto-ack、业务失败仍可能丢。acks=all 只覆盖“生产者 → ISR 副本”这一段,不是端到端保证。

L2|“消费者处理成功了,但提交 offset 前挂了,消息会怎样?” → offset 未提交,重启后从上次提交点重投 → 重复消费,不会丢。这正是 at-least-once:宁可重复不可丢失,所以消费必须幂等。可补:Kafka 提交 offset 默认 5s 一批,崩溃窗口内未提交的消息都会重投;用 enable.auto.commit=false + 业务成功后手动 commit,可把重复面缩到“业务成功但 commit 前挂掉”这一极窄窗口。

L3|“如果业务方坚持要‘恰好一次’呢?” → 严格 exactly-once 在分布式中昂贵。实践拆成:生产端幂等生产者(Kafka enable.idempotence)+ 事务(跨分区原子写)+ 消费端幂等。对多数业务,at-least-once + 消费幂等已达到业务上的“效果恰好一次”,成本远低于端到端 exactly-once。

【评分标准】

档位答案特征
60 分知道要开 ack,但分不清生产 ack 与消费 ack
80 分能按生产/Broker/消费三段给出配置与机制,点名 acks=all、confirm、手动 ack
95 分说清每段“为什么会丢”的因果;主动提对账兜底;指出不丢必然带来重复、需幂等;能对比三家 MQ 的实现差异

【关联题】

  • 同簇: 第 81 题(重复消费)、第 94 题(ack 机制)、第 88 题(死信)、第 90 题(重试)
  • 生产侧一致: 第 92 题(事务消息)、第 93 题(本地消息表)
  • 下游保障: 第 110 题(接口幂等)
  • 对比参照: 第 32 题(Redis 单线程)——同样是“默认性能优先 vs 要什么就得显式开什么”

【自测】

  1. Kafka 的 acks=1 在什么情况下会丢消息? 参考答案: leader 写入成功并 ack,但 follower 尚未同步时 leader 宕机且发生非干净选举,已 ack 消息丢失。acks=all 可避免(等 ISR 全部写入)。
  2. 为什么消费端要“处理成功再 ack”? 参考答案: 若收到即 ack,处理中崩溃则消息被标记成功不再投递 → 丢失。处理后再 ack,崩溃时未 ack 会重投 → 最多重复,不丢。
  3. 判断对错:只要 Broker 做了副本,消息就一定不丢。 参考答案: 错。还需生产端确认、消费端正确 ack、禁止不干净选举,否则任一段仍可能丢。

81. 网络抖动导致同一条消息被消费两次(重复消费) ​

【考察内容】消息重复消费是 MQ 必考题

【题目】网络抖动或消费者宕机重平衡,导致同一条消息被投递了两次,业务被重复处理(比如重复发奖)。怎么保证“消息只被处理一次”?消费端的幂等设计怎么做?

【参考答案】

  1. 认知:重复消费在至少一次(at-least-once)投递语义下必然存在(ack 丢失、重试、Rebalance、进程崩溃),只能靠“消费幂等”解决;
  2. 幂等方案:
    • 业务唯一键 + 数据库唯一索引:处理前先插入唯一键(如 order_id),重复插入报错则跳过——最可靠;
    • 状态机幂等:消息处理推进状态机(如订单:待支付→已支付),重复消息发现状态已推进则忽略;
    • Redis SETNX 去重:以消息 ID/业务 ID 为 key,SETNX 成功才处理(配合过期时间),Redis 挂了由 DB 唯一索引兜底;
    • 消息表去重:本地记录已处理消息 ID 表,查重后处理;
  3. 关键:幂等要基于业务唯一键(订单号、支付流水号)而非消息本身 ID(同业务可能产生多条消息);
  4. 设计原则:消费逻辑本身要“可重入”(同参数执行多次结果一致)。
  5. 容量/一致性窗口估算:幂等存储(DB 唯一索引/Redis SETNX)要按峰值消费 TPS 设计——若峰值 5000 msg/s,Redis SETNX 需能扛同量级写(单实例通常 10 万+ ops,足够);DB 唯一索引插入在热点业务键上可能成瓶颈,建议按业务键分表或走 Redis 预检 + DB 最终落库。去重 key 的 TTL:至少覆盖最大重试窗口 + 时钟偏差 + 业务最长处理时间,常见 24h~7 天;过短会在延迟投递时失效导致重复。
  6. 失败与降级:Redis 去重层宕机 → 降级只靠 DB 唯一索引(可扛较低 QPS)或本地缓存布隆预筛;唯一键冲突 → 视为“已处理”直接 ack,不再进死信(否则死信被打爆);状态机发现已推进 → 同样静默成功。原则:幂等失败要“安全成功”而不是“安全失败”,否则重复流量会拖垮核心库。

【原理溯源】

  • 为什么 at-least-once 必然重复? 确认信号与业务副作用不在同一个原子单元里:消费者先做业务再 ack,若 ack 丢失/超时/进程在 ack 前崩溃,Broker 认为未成功会重投。反之若先 ack 再做业务,崩溃就丢消息。分布式里没有“既确认又保证副作用恰好一次且不丢”的免费午餐——要么可能丢,要么可能重。at-least-once 选择了“宁重不丢”。
  • 哪些操作会触发重复? ① 消费者 ack 丢失后重投;② 生产者超时重发(发送成功但 ack 未收到);③ Kafka Rebalance:分区再分配后新消费者从上次提交 offset 重放;④ 消费者处理成功但提交 offset 前崩溃;⑤ 网络分区导致的重复投递。这些都不是“偶发 bug”,而是协议允许的正常路径。
  • 为什么幂等键要用业务唯一键而不是消息 ID? 同一业务事件可能被重发成多条消息(不同 msgId),用 msgId 去重会漏;同一 msgId 也可能对应需要多次处理的语义边界不清。业务唯一键(订单号 + 业务动作)才能表达“这件事只应发生一次”。
  • 为什么唯一索引是最可靠的幂等手段? 数据库唯一约束是与业务写入同一个本地事务的原子检查:插入成功 = 第一次处理;冲突 = 已处理过。不存在 Redis 那种“检查与执行不原子”的窗口。Redis SETNX 只能做前置过滤减少冲突,最终正确性靠 DB。
  • 状态机为什么天然幂等? 状态迁移是“从状态 A 到 B”,若已是 B,重复消息判定为 no-op。前提是迁移条件明确、单向,且并发下用乐观锁/唯一约束保护(WHERE status='A' 更新)。

【选型判断树】

如何做消费幂等?
├─ 消息是否可映射到稳定业务唯一键(订单号/流水号)?
│   ├─ 是 → DB 唯一索引 / 状态机(正确性优先)
│   └─ 否 → 先补齐业务键,或退化为消息表按 msgId 去重(有漏重风险)
├─ 吞吐要求?
│   ├─ 高 → Redis SETNX 前置过滤 + DB 唯一约束兜底
│   └─ 中低 → 直接 DB 唯一索引
└─ 业务是否本身就是状态推进?
    └─ 是 → 状态机 + 乐观锁(天然幂等)

口诀:唯一约束保正确,Redis 只挡流量,状态机管推进。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“重复消费在 at-least-once 下必然存在,解法不是消灭重复,而是幂等”
0:30–2:00根因讲清 ack 与副作用非原子 → 要么丢要么重;列 Rebalance/重试等触发场景
2:00–3:30方案分层唯一索引(正确性)→ 状态机(业务语义)→ Redis(性能)→ 消息表(兜底)
3:30–4:30关键细节强调业务唯一键 > msgId;Redis 必须设 TTL;冲突异常要能识别为“已处理”
4:30–5:00收尾“幂等是 MQ 可靠性闭环的最后一块拼图,与不丢、有序成套设计”

【关键数字】

参数经验值说明
Redis 去重 TTL与消息最大重试窗口对齐(如 24h)过短漏防,过长占内存
重复率网络抖动时 0.1%–1% 不罕见必须工程化处理,不能当小概率
唯一索引冲突耗时毫秒级高并发下需 Redis 前置
SETNX + DB 双层Redis 挡 90%+ 重复DB 仍可能被穿透
峰值幂等 QPS= 峰值消费 msg/sRedis 轻松;DB 需关注热点键
唯一键选择业务键优先于消息 ID同业务多条消息时消息 ID 无效

【追问链】(三层)

L1|“加个消息 ID 去重不就行了?” → 不够。同一业务可能产生多条 msgId(生产端重发),会漏;且 msgId 与业务正确性无关。应基于订单号/流水号等业务唯一键。

L2|“Redis SETNX 和 DB 唯一索引都用,先查哪个?” → 先 Redis 后 DB。Redis SET key NX EX 是 O(1)、单实例能扛 10 万+ ops,先过滤掉绝大多数重复,避免每条消息都去打 DB 唯一索引(热点业务键上会形成索引/行锁争用,正是本参考答案第 5 条点名的瓶颈)。但 Redis 不能当唯一真相——重启、主从切换、key 过期都会让去重记录消失,所以最终兜底是 DB 唯一索引:插入冲突即判定重复并放弃业务。顺序反过来等于把去重压力全压到 DB。

L3|“消费逻辑里调用了下游接口,怎么保证重试时不重复扣款?” → 下游接口本身要提供幂等(幂等号=业务唯一键)。消费端重试时传同一幂等号;下游用唯一约束/状态机防重。若下游无幂等能力,则消费端需本地记账“已调用成功”后再跳过,或改为“先落地意图再异步调用”并接受对账补偿。

【评分标准】

档位答案特征
60 分知道要“去重”,方案只有 Redis set
80 分承认 at-least-once 必然重复;给出唯一索引/状态机/Redis 分层;强调业务唯一键
95 分讲清 ack 与副作用非原子的根因;说明 Redis 只做前置、DB 保正确;能处理下游接口幂等与 Rebalance 场景

【关联题】

  • 同簇: 第 80 题(不丢)、第 94 题(ack)、第 87 题(Rebalance)、第 90 题(重试)
  • 通用幂等: 第 110 题(接口幂等)
  • 生产侧: 第 92/93 题(发送端也可能重复)

【自测】

  1. 为什么说重复消费“无法避免,只能幂等”? 参考答案: 因为 ack 与业务副作用不在同一原子单元;at-least-once 协议下,崩溃/网络问题会导致“业务已做但未确认”从而重投。
  2. 幂等键为什么不用 msgId? 参考答案: 同一业务事件可能对应多条 msgId(生产端重发),用 msgId 会漏重;业务唯一键才能表达“只应发生一次”。
  3. 判断对错:Redis SETNX 成功就可以不落库唯一索引。 参考答案: 错。Redis 非持久、可能丢、与业务非原子,只能前置过滤;正确性必须靠 DB 唯一约束或状态机。

82. 订单的“创建→支付→发货”消息乱序了(消息顺序性) ​

【考察内容】消息顺序是 MQ 三大可靠性问题之一

【题目】订单系统的状态流转依赖消息顺序:必须先消费“创建”、再“支付”、再“发货”,一旦乱序订单状态就错。消息队列怎么保证这种顺序性?Kafka/RocketMQ 分别怎么实现?代价是什么?

【参考答案】

  1. 认知:全局有序成本极高(只能单分区单消费者,吞吐骤降),生产上只做局部有序(同一业务 key 的消息有序);
  2. 方案:
    • Kafka:同一业务 key(如 order_id)作为消息 key,hash 到同一个分区(同一分区内消息有序);消费者组内一个分区只能被一个消费者线程消费,单线程按序处理;
    • RocketMQ:顺序消息(MessageQueueSelector 按业务 id 选队列)+ 单消费者;
    • RabbitMQ:同一条队列 + 单消费者(队列天然 FIFO,但要避免多消费者并发消费同一队列);
  3. 代价:同一 key 的消息串行处理,吞吐受限——拆 key 粒度(用户级/店铺级)平衡;热点 key 可在允许乱序的业务维度上拆(如按用户/店铺分队列);不能按「订单子流程」拆——那会把同一个 order_id 的创建/支付/发货打散到不同队列,直接破坏题干要求的同单有序;真要提速只能改成状态机版本号+条件更新来容忍乱序(另见原理溯源的旁路 topic 方案,同样必须配版本裁决);
  4. 失败处理:消息处理失败不能跳过继续(会乱序),需重试阻塞或进死信并暂停后续(业务补偿);
  5. 兜底:对账/重放机制纠正极端乱序。
  6. 容量与并行度估算:订单状态消息峰值假设 1 万 msg/s,若 order_id 哈希到 N 个分区,则有序消费最大并行度 = N。要保证同一订单串行:同 order_id 固定分区;不同订单可并行。若单消费者处理 500 msg/s,则需要至少 20 个分区 + 20 个消费者实例才能扛 1 万/s。代价:热点商户/热点订单单分区可能过热,需在业务键中加入合理分散因子(如 seller_id+order_id)——前提是同一订单的创建/支付/发货三条消息上该因子恒等(同一卖家、同一订单),否则同一 order_id 会因字段取值不同落到不同分区,顺序照样破;所以分散因子只能用下单时就固定、与事件类型无关的稳定字段(如订单号自带的分片号),绝不能用「消息类型」「当前状态」这类会变的属性。也不能破坏“必须有序的最小键”。
  7. 失败与降级:单分区内消息失败 → 该分区阻塞(保序必然牺牲局部可用);处理策略:① 可跳过的非关键乱序消息进死信;② 关键状态走状态机幂等,旧状态消息直接忽略;③ 极端时对该分区做“人工跳过 + 对账修复”。禁止为提速把同 key 打散到多分区——那会重新引入乱序。Rebalance 期间可能短暂重复/乱序,消费端必须幂等。

【原理溯源】

  • 为什么全局有序不可行? 全局 FIFO 要求所有消息进入同一队列/分区,并由单线程消费——并行度降为 1,吞吐塌到单线程处理能力。订单场景真正需要的是同一订单内有序,订单之间无顺序依赖。用全局有序去满足局部需求,是用 1% 的需求杀掉 99% 的吞吐。
  • 为什么“同 key 同分区”能保证局部有序? Kafka 分区内消息以 append-only 日志存储,offset 单调递增,分区内天然 FIFO。key hash 决定分区 → 同一订单的创建/支付/发货落入同一分区 → 按 offset 消费即按序。分区是顺序的边界。
  • 为什么一个分区不能被组内多个消费者并行消费? 两个原因:① 并发处理会打乱“先创建后支付”的处理顺序;② offset 以“组 + 分区”为单位提交,多消费者提交会互相覆盖,进度错乱。所以组内是分区独占分配。
  • 重试为什么会破坏顺序? 失败消息若被转发到重试 topic 与后续消息并行,或跳过失败继续消费后续——创建失败被跳过、支付先被处理,状态机直接错乱。顺序场景下失败必须阻塞当前 key 的后续消息(或暂停该分区),直到成功或进死信并人工介入。
  • 代价从哪来? 同 key 串行 = 单 key 吞吐上限是单线程处理速度。热点 key(秒杀单品、大 V 直播)会把该分区打满。缓解:拆 key 粒度(用户/店铺)、把可乱序的旁路事件另建 topic(创建/支付/发货这类同单主链路禁止拆子流程——会把同 order_id 打散,与本条第 3 项的约束相反),或业务层允许最终一致 + 状态机防乱序(乱序到达时用版本号/状态前置条件拒绝过期消息)。

【选型判断树】

需要什么级别的顺序?
├─ 业务上真需要全局 FIFO?
│   └─ 几乎从不 → 不要选全局有序
├─ 是否只要求同一实体(订单/用户)内有序?
│   └─ 是 → 局部有序:同 key 同分区 + 分区独占消费
├─ 是否存在热点 key?
│   ├─ 否 → 直接 hash 分区即可
│   └─ 是 → 拆 key 粒度 / 旁路 topic 分流(主链路禁拆子流程)/ 状态机版本号防乱序
└─ 消费失败怎么办?
    ├─ 阻塞重试(保序)
    └─ 进死信 + 暂停该 key 后续(人工/补偿)

口诀:顺序以分区为界;同 key 同分区;失败不能跳。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“全局有序吞吐不可接受,生产只做同 key 局部有序”
0:30–2:30机制Kafka:key hash → 同分区 → 分区内 FIFO + 分区独占;RocketMQ MessageQueueSelector;Rabbit 单消费者
2:30–3:30代价与热点同 key 串行;热点 key 拆分策略
3:30–4:30失败处理不能跳过;阻塞重试或死信+暂停;对账兜底
4:30–5:00收尾“顺序、吞吐、失败处理三者要一起设计”

【关键数字】

参数经验值说明
全局有序吞吐降为单线程消费能力通常不可接受
局部有序并行度= 分区数(非热点 key 时)分区数决定上限
热点 key 判定单分区 lag 持续上涨需拆 key
顺序失败策略阻塞 ≤ 数次重试后进死信无限阻塞会堵分区
有序并行度上限= 分区数(同 key 单分区内串行)1万 msg/s ÷ 单消费者500 = 约20分区
单分区消费能力单线程/单实例慢逻辑会成为整 key 热点
热点 key 策略业务键粒度取舍不能为分散而破坏最小有序键

【追问链】(三层)

L1|“Kafka 不是能保证顺序吗?” → 只保证单分区内有序。多分区时跨分区无序。要业务顺序必须同 key 同分区,并保证组内分区独占。

L2|“消费者多线程消费同一个分区行不行?” → 组内不行。一个分区在同一消费组里只能分给一个消费者;两个线程共享同一分区就得共享同一份消费进度,offset 由谁提交、提交到哪条都说不清,Rebalance 时还会互相覆盖,直接造成重复或丢失。要提并行度只有两条正路:① 加分区(保持同 key 同分区,顺序不破);② 消费者实例数 ≤ 分区数,多出来的只会空闲。若坚持在单进程内多线程处理同一分区,只能自己拉下来再按 key 派发到工作线程,这会破坏分区内有序(完成顺序不定),且 offset 必须等最慢的一条才能提交。

L3|“消息已经乱序到达了,消费端还能救吗?” → 能。用状态机前置条件 + 版本号:处理“支付”时校验当前状态必须是“已创建”;版本号更小的旧消息直接丢弃或补偿。这把“传输层顺序保证”降级为“业务层状态一致性”,对最终一致系统更稳健。

【评分标准】

档位答案特征
60 分说“Kafka 保证顺序”
80 分知道局部有序 vs 全局有序;说清同 key 同分区 + 单消费者
95 分解释分区是顺序边界;讲清失败阻塞与热点 key;能提状态机兜底乱序

【关联题】

  • 同簇: 第 87 题(分区与消费者)、第 83 题(积压)、第 90 题(重试)、第 81 题(幂等)
  • 业务侧: 订单状态机设计(可关联第 110 题幂等)

【自测】

  1. 为什么生产环境几乎不用全局有序? 参考答案: 全局有序要求单分区单消费者,并行度为 1,吞吐无法支撑业务;实际只需同业务实体内有序。
  2. 同一订单的三条消息 hash 到不同分区会怎样? 参考答案: 跨分区无顺序保证,可能先消费支付再消费创建,状态错乱。必须保证同 order_id 同分区。
  3. 判断对错:顺序消息消费失败应该跳过,先处理后面的。 参考答案: 错。跳过会破坏顺序,应阻塞重试或进死信并暂停该 key 后续消息。

83. 一觉醒来 MQ 积压了几百万条消息(消息积压) ​

【考察内容】消息积压是生产事故高频题

【题目】凌晨上线的一个 bug 让消费者全挂了,早上发现 MQ 积压了几百万条消息,业务链路大面积延迟。先恢复消费还是先修 bug?积压的消息怎么快速消化(扩容、临时消费者、跳过)?会不会挤垮下游?

【参考答案】

  1. 先定位原因:
    • 消费能力不足(消费者实例少、消费逻辑慢);
    • 下游故障(DB/依赖服务挂了,消费一直失败重试);
    • 生产端突发流量(大促)超过消费速率;
  2. 处理流程:
    • 修复下游故障:如果是下游问题,先恢复下游(否则扩容也没用);
    • 第 0 步:先回滚或热修(题干根因就是「凌晨上线的 bug 让消费者全挂」):扩容、临时 topic、转发都跑在同一份有 bug 的代码上,消费照样失败,重试还会把下游再压一遍;所以下列手段都以「消费代码已能正常处理」为前提;
    • 临时扩容消费者:增加消费者实例(注意:Kafka 单分区只能一个消费者消费,分区数=最大并行度——可临时增加分区并重新分布,或用“多线程消费”);更快的办法:把积压消息转发到新建的临时 topic(分区数扩大),用更多消费者并行消费,消费完再处理结果;
    • 降级策略:非核心消息(如日志、统计)在积压严重时可跳过/丢弃(业务可接受);核心消息不能丢;
    • 限流生产端:暂停/减缓非关键生产,先清积压;
  3. 根治:优化消费逻辑(批量消费、异步化)、评估分区数与消费者数匹配、监控堆积量告警(堆积超过阈值自动告警/扩容);
  4. 关键防坑:不要盲目重置位点、或用新消费组从 auto.offset.reset=latest 起(会把积压直接跳过,等效丢消息);单纯重启不会丢已提交 offset——位点存在 Broker 上;扩容要保证消费幂等(重复消费仍安全)。
  5. 容量估算(步进):积压 L=500 万条,下游恢复后安全消费速率 C=500 msg/s,则清空时间 ≈ 500万/500 = 10000s ≈ 2.8 小时(未计生产新增)。若可扩到 5000 msg/s(10 实例×500,且分区≥10),约 17 分钟(500万÷5000=1000s)。但这个数与本题的限流口径必须二选一:本条后半与第 6 条又按「DB 只能扛 1000 TPS」限流,按 1000/s 算是 500万÷1000=5000s≈1.4 小时,恢复期 20%→50%→100% 的爬坡也不能越过该阈值——答题要先声明用的是哪个口径,否则「17 分钟」与「限流 1000」互相抵消。生产端若仍以 P=2000 msg/s 写入、消费 C=500,则积压还在以 1500 msg/s 上涨——必须同步限流生产。磁盘:500 万×1KB×副本3 ≈ 15GB 量级,一般可承受;但保留策略与磁盘水位要盯住。下游 DB:清积压时写入冲到峰值,若 DB 只能扛 1000 TPS,消费限流就不得超过 1000。
  6. 失败与降级:下游未恢复前禁止盲目扩容;非核心 topic(日志/统计/埋点)可丢弃或采样;核心消息不可丢,但可“只处理最近 N 小时 + 更早走对账补偿”;消费端必须幂等(临时 topic 转发会带来重复);恢复期并发按 20%→50%→100% 爬坡,观察下游错误率。监控:lag、消费 RT、下游错误率、磁盘水位四条曲线一起看。

【原理溯源】

  • 积压的数学本质是什么? 积压量 L 的变化:dL/dt = 生产速率 P − 消费速率 C。当 P > C 持续,L 线性上涨。清积压只有三条路:提高 C、降低 P、或对部分消息“取消消费”(丢弃/跳过非核心)。
  • 为什么“下游挂了”时扩容消费者是白扩? 消费瓶颈不在消费者线程,而在下游。扩容只会让更多线程同时打挂掉的下游 → 更多重试 → 更大压力 → 可能把下游彻底打死。先修下游,再谈扩容。
  • Kafka 为什么“加机器不加速”? 组内并行度 = 分区数。消费者数 > 分区数时多余消费者空闲。所以扩容前先看分区数;不够就扩分区或用临时 topic 转发。临时 topic 方案避免了直接改原 topic 分区带来的 Rebalance 与迁移风险。
  • 为什么“直接改线上分区数”要谨慎? 增分区会改变 key→partition 映射,同 key 消息可能进入新分区,破坏局部有序;并触发大规模 Rebalance,短暂影响可用性。事故中更稳妥的是新 topic 转发,原 topic 保持不动。
  • 清积压会不会挤垮下游? 会。积压解除时消费速率会冲到峰值,若下游刚恢复仍脆弱,可能二次打挂。正确做法:恢复下游后限流爬坡(逐步放开消费并发),配合批量与背压,而不是全速冲。

【选型判断树】

发现积压
├─ 根因?
│   ├─ **上线 bug/代码异常 → 第 0 步回滚或热修(扩容跑在同一份坏代码上是白扩)**
│   ├─ 下游故障 → 先恢复下游(禁盲目扩容)
│   ├─ 消费逻辑慢 → 优化逻辑/批量/异步;再扩分区与实例
│   └─ 生产突增 → 限流非关键生产 + 扩容消费
├─ 如何加速?
│   ├─ 分区数充足 → 加消费者实例
│   ├─ 分区数不足 → 新建临时大分区 topic 转发
│   └─ 非核心消息 → 丢弃/跳过(业务确认)
└─ 恢复期
    ├─ 限流爬坡,防止二次打挂下游
    └─ 全程幂等 + 监控 lag 回落

口诀:先根因后扩容;分区=并行度;恢复要限速。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“积压是 P>C 的流量问题,先分清是下游挂、消费慢还是生产突增”
0:30–2:30处置顺序先回滚/热修(题干根因是凌晨上线的 bug) → 恢复下游 → 扩容/临时 topic → 非核心丢弃 → 限流生产
2:30–3:30Kafka 特殊点分区=最大并行度;临时 topic 而不是直接改分区
3:30–4:30防二次事故恢复限流爬坡;幂等;不乱动 offset
4:30–5:00收尾“积压处置 = 根因 + 并行度 + 下游保护”

【关键数字】

参数经验值说明
消费并行度上限= Kafka 分区数超出的消费者空闲
积压告警阈值绝对量(如 10 万)或 lag 持续上涨时长(5–10min)按业务 SLA 定
临时 topic 分区原分区数 × 2–10视下游容量
恢复爬坡并发 20% → 50% → 100%每档观察下游 RT/错误率
非核心可丢日志/统计/部分埋点需业务事先确认
清空时间积压量 ÷ 安全消费速率500万÷500/s≈2.8h;÷5000/s≈17min;若下游只能 1000/s 则≈1.4h——17min 与「限流 1000」二选一,答题要先声明口径
下游保护限流≤ 下游稳态容量如 DB 1000 TPS → 消费限 1000
积压磁盘条数×消息体×副本500万×1KB×3≈15GB

【追问链】(三层)

L1|“积压了先加消费者行不行?” → 看根因。下游挂了加消费者是白扩甚至有害;代码有 bug 时扩容照样失败,还会用重试把下游再压一遍,必须先回滚/热修;只有「消费能力不足且分区够」时才有效。先看错误率与下游健康。

L2|“能不能直接把分区数从 12 调到 60?” → 技术上可以,但:① 触发全量 Rebalance;② key 映射变化可能破坏同 key 顺序;③ 迁移期间可用性抖动。事故中更稳:新建分区更多的临时 topic,把积压消息转发过去并行消费,原 topic 不动。若业务要求严格有序,扩分区前必须确认状态机/对账能否吸收短期乱序,否则只扩“可乱序”的旁路 topic。

L3|“几百万条积压,下游只恢复到每秒 500,要跑很久,业务等不了怎么办?” → 分层处理:① 非核心丢弃;② 核心中可合并的做批量聚合(如同订单多条合并);③ 与业务确认“只处理最近 N 小时”的降级窗口,更早的走对账补偿。目标从“全部实时消费”改为“核心正确 + 可接受延迟”。

【评分标准】

档位答案特征
60 分说“加消费者扩容”
80 分先分根因;知道 Kafka 分区=并行度;会用临时 topic
95 分先答「根因是发布引入的 bug 就先回滚/热修」,再谈修下游与扩容;恢复限流爬坡、非核心丢弃策略、幂等与 offset 防坑

【关联题】

  • 同簇: 第 87 题(消费者组)、第 95 题(消费流控)、第 96 题(监控)、第 90 题(重试)
  • 对照: 第 179 题(突发流量应急)——同样是“先止血再根治”

【自测】

  1. 下游 DB 挂了导致积压,扩容消费者有用吗? 参考答案: 没用且有害。瓶颈在下游,应先恢复下游,再限流爬坡恢复消费。
  2. 为什么 Kafka 加消费者实例有时速度不变? 参考答案: 消费者数超过分区数后,多余实例分不到分区,空闲。并行度上限=分区数。
  3. 判断对错:清积压时应全速消费尽快追平。 参考答案: 错。全速可能二次打挂刚恢复的下游,应限流爬坡。

84. 选型会:Kafka/RabbitMQ/RocketMQ 用哪个(MQ 选型) ​

【考察内容】MQ 选型是高频综合题

【题目】公司要做消息中间件选型,候选 Kafka、RabbitMQ、RocketMQ。按吞吐量、可靠性、顺序性、事务消息、延迟、运维成本几个维度对比,并给出你的业务场景(比如日志收集 vs 交易订单)该怎么选?

【参考答案】

  1. Kafka(顺序性:key→分区内原生局部有序,跨分区无序;组内分区独占消费,无原生死信/延迟):优点——超高吞吐(顺序写+零拷贝+批量)、天然分布式分区、消息可回溯(按 offset 重放)、生态好(大数据/日志);缺点——功能简单(无死信/延迟等高级特性)、消息可能重复(at-least-once)、管理较复杂。适合:日志、埋点、大数据、削峰场景;
  2. RabbitMQ(顺序性:三家最弱——只有单队列+单消费者才有序,多消费者或调大 prefetch 即乱序):优点——功能丰富(延迟队列、死信、优先级、灵活路由)、Erlang 稳定可靠、社区成熟;缺点——吞吐中等(万级)、分布式扩展(镜像队列)不如 Kafka。适合:业务解耦、复杂路由、中小流量;
  3. RocketMQ(顺序性:官方顺序消息 API=MessageQueueSelector 队列选择器+失败阻塞重试,最省事):优点——阿里出品、吞吐高(十万级)、事务消息/延迟消息/死信内置、Java 生态友好;缺点——社区相对小、功能多但运维复杂。适合:电商交易、金融业务消息(需要事务消息和可靠性的场景);
  4. 选型结论:日志/大数据→Kafka;交易/事务/延迟消息→RocketMQ;轻量业务/复杂路由→RabbitMQ;技术栈统一优先考虑团队熟悉度。
  5. 容量锚点(选型用):日志/埋点类常见日增 TB 级、峰值 数十万~百万 msg/s → Kafka;交易订单消息峰值 1万~10万 msg/s、单条 KB 级、需事务/延迟 → RocketMQ;企业内业务解耦 数百~数千 msg/s、复杂路由 → RabbitMQ。容量粗算:峰值 msg/s × 消息体 × 86400 × 副本 × 保留天 = 磁盘(别漏掉「×86400 秒/天」这一步);例:5万 msg/s×1KB ≈ 50 MB/s ≈ 4.32 TB/日,×3副本×3天 ≈ 约 39 TB。消费端能力要 ≥ 生产峰值,否则必须削峰策略。选型时同时估:集群规模、运维人力、机房/跨城复制带宽。
  6. 失败与降级:任何 MQ 都要做生产 confirm、Broker 持久化副本、消费手动 ack+死信;集群故障时业务侧本地消息表/降级开关;多 MQ 并存时要定义“哪条链路依赖谁”以及故障时的隔离边界。Kafka 无原生延迟消息时用“时间轮/延迟表+二次投递”兜底;RocketMQ 半消息回查超时要有本地事务状态表,避免悬挂。

【原理溯源】

  • 为什么 Kafka 吞吐最高? 架构目标是“高吞吐日志流”:顺序写磁盘、页缓存、零拷贝 sendfile、生产/消费批量与压缩、分区并行。它用功能做减法换吞吐——没有内置事务消息(早期)、延迟队列、复杂路由。选型时要明白:Kafka 的强项来自“把 MQ 当日志存储系统”。
  • 为什么 RabbitMQ 适合复杂路由但吞吐有限? 每条消息经过 exchange 路由到队列,Broker 做更多逻辑;镜像队列的复制开销大。它强在协议与路由模型(topic/fanout/direct、死信、优先级、TTL),弱在水平扩展。中小流量业务解耦非常合适。
  • 为什么交易场景偏 RocketMQ? 电商/金融需要:事务消息(本地事务与发消息原子)、延迟消息(关单)、死信、顺序消息 API 这些 RocketMQ 内置,Kafka 要自研;但「按业务 key 局部有序」Kafka 原生就有(key→分区 + 分区内 FIFO,见第 82 题),RocketMQ 多出的是队列选择器与失败阻塞重试。对交易团队,“开箱即用的事务消息”价值远大于“再多 5 倍吞吐”。
  • “可靠性”三家差在哪? 都能通过配置做到接近不丢(副本、刷盘、confirm)。差异在语义完备性:RocketMQ 事务消息提供生产端与本地事务的原子性;Kafka 提供幂等生产者 + 事务(跨分区原子,但不等于业务事务消息);RabbitMQ 有 confirm + 事务(channel tx,性能差,少用)。
  • 运维成本如何估? Kafka/RocketMQ 集群组件多(ZK/KRaft、NameServer),需要专人;RabbitMQ 简单但扩展天花板低。还要看团队语言栈:Java 团队 RocketMQ 顺手;大数据生态 Kafka 必选。

【选型判断树】

场景是什么?
├─ 日志 / 埋点 / 大数据管道 / 超高吞吐削峰 → Kafka
├─ 交易 / 支付 / 需要事务消息或延迟消息 → RocketMQ
├─ 中小流量业务解耦 / 复杂路由 / 优先级 / 快速上手 → RabbitMQ
└─ 多场景并存
    ├─ 允许多 MQ → 按场景分用(常见:Kafka 做日志,RocketMQ 做交易)
    └─ 必须统一 → 交易优先则 RocketMQ;数据平台优先则 Kafka + 自研事务层

附加检查:
- 团队是否熟悉运维?
- 是否需要跨机房 / 多租户?
- 生态绑定(Flink/Spark → Kafka)

口诀:吞吐与生态选 Kafka,交易特性选 RocketMQ,路由与轻量选 RabbitMQ。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“选型不是比参数表,是按场景匹配能力边界”
0:30–2:30三家对比每家:架构取向 + 优势 + 硬伤 + 典型场景,各 30–40 秒
2:30–3:30落地结论日志→Kafka,交易→RocketMQ,轻量路由→RabbitMQ
3:30–4:30工程因素运维、团队栈、生态、是否允许多 MQ 并存
4:30–5:00收尾“没有最好,只有最匹配;交易特性与吞吐不可兼得时按业务分级”

【关键数字】

参数经验值说明
Kafka 吞吐单机数十万–百万级 msg/s批量 + 压缩时
RocketMQ 吞吐十万级 msg/s交易特性更全
RabbitMQ 吞吐万级 msg/s镜像队列更低
顺序性Kafka:key→分区内原生局部有序,跨分区无序;RocketMQ:队列选择器+阻塞重试,官方 API 最直接;RabbitMQ:单队列+单消费者才有序,最弱题干点名维度,不能只比吞吐/功能/生态
事务消息RocketMQ 原生;Kafka 事务≠业务事务消息选型分水岭
副本数≥2,生产常 3可靠性
Kafka 吞吐锚点单机数十万–百万 msg/s批量+压缩时
RocketMQ 吞吐锚点十万级 msg/s交易特性更全
RabbitMQ 吞吐锚点万级 msg/s镜像队列更低
磁盘粗算峰值msg/s×体×86400×副本×天5万/s×1KB≈4.32TB/日,×3×3≈39TB

【追问链】(三层)

L1|“Kafka 能保证不重复吗?” → 默认 at-least-once,可能重复。可开幂等生产者(enable.idempotence)避免生产端重复;消费端仍要幂等。跨分区 exactly-once 用事务 API,但业务上的“恰好一次”仍推荐消费幂等。

L2|“RabbitMQ 高吞吐场景行不行?” → 万级 QPS 内可以;更高吞吐、要水平扩展、或已绑定大数据生态时换 Kafka/RocketMQ。镜像队列复制开销大,日志洪峰不适合。补充判断:若峰值稳定在 3k–5k msg/s 且路由复杂,RabbitMQ 仍可胜任;一旦出现“大促 10 倍尖峰 + 需要回放”,应提前按 Kafka/RocketMQ 规划分区与容量,而不是上线后再换。

L3|“已经用了 Kafka,但业务要事务消息,怎么办?” → 三条路:① 业务侧本地消息表(93 题)——最通用;② 换/并用 RocketMQ 做交易通道;③ Kafka 事务 + 消费端幂等做“写多分区原子”,但解决不了“本地 DB 事务与发消息”的原子,仍需消息表或对账。不要误以为 Kafka 事务 = RocketMQ 事务消息。

【评分标准】

档位答案特征
60 分能说出三家名字和大概用途
80 分按吞吐/可靠性/顺序性/事务消息/延迟/运维成本六个维度对比,给出场景结论
95 分解释架构取向(Kafka 日志流、RocketMQ 交易特性、Rabbit 路由);指出 Kafka 事务≠事务消息;考虑运维与多 MQ 并存策略

【关联题】

  • 同簇: 第 79 题(价值)、第 86 题(Kafka 吞吐)、第 89 题(延迟)、第 92 题(事务消息)
  • 落地: 第 91 题(降级)、第 96 题(监控)

【自测】

  1. 日志收集与订单状态通知,分别选什么?为什么? 参考答案: 日志选 Kafka(高吞吐、回放、大数据生态);订单通知选 RocketMQ(可靠、事务/延迟/死信内置,Java 交易栈)。
  2. 为什么不能说“Kafka 有事务所以等价于 RocketMQ 事务消息”? 参考答案: Kafka 事务保证跨分区读写的原子性与幂等生产,不解决“本地数据库事务与发消息”的两阶段问题;后者需要 RocketMQ 半消息或本地消息表。
  3. 判断对错:功能越多的 MQ 越好,应优先选功能全的。 参考答案: 错。功能多常伴随吞吐/复杂度代价。应按场景选:日志不需要事务消息,交易不需要百万 TPS 日志管道。

85. 自研消息队列技术预研,核心能力有哪些(MQ 设计) ​

【考察内容】“设计一个 MQ”是考察系统设计深度的经典题

【题目】团队做技术预研:想自研一个消息队列。从零设计,需要哪些核心能力?存储模型(日志/队列)、生产消费模型、ack 与重试、顺序、堆积、集群高可用分别怎么设计?

【参考答案】按“最小可用→逐步补齐”的顺序:

  1. 基础模型:生产者 → Broker(队列/分区)→ 消费者,消息的发送、存储、消费基本流程;
  2. 持久化:消息落磁盘(顺序写),解决重启丢消息——顺序写文件+批量刷盘;
  3. 高可用:副本机制(leader/follower),主故障自动选举,解决单点;
  4. 可靠性:生产端 ack、Broker 持久化、消费端 offset/ack 机制(处理成功才提交);
  5. 顺序性:分区/队列模型 + 同一 key 路由同分区 + 单消费者消费;
  6. 负载均衡:消费者组内分区分配(一个分区一个消费者)、Rebalance;
  7. 高级能力:延迟消息(时间轮/分级)、重试设计(独立重试队列/Topic+指数或分级退避+最大重试次数+不可重试异常直投死信;重试按延迟等级分档存放,避免定时器轮询全量消息;顺序 topic 的重试必须阻塞该 key 而不是并发重投,否则重新引入乱序)、死信队列(重试耗尽转人工,DLQ 要可查询、可回放)、事务消息(半消息+确认)、消息回溯(按 offset 重放)、削峰(消费限流);
  8. 其他:监控(积压量、消费 lag)、权限、压缩、多租户。
  9. 容量估算(自研 MQ 预研目标可写进设计文档):目标单机写入 5万–10万 msg/s(1KB 消息、3 副本、异步刷盘);同步刷盘场景降到 1万–3万 msg/s 量级。分区数规划:生产峰值 / 单分区消费能力;例峰值 5 万、单消费者 500/s → 至少 100 分区(再留 50% 余量)。存储:日增消息条数×体×副本×保留天;堆外内存/页缓存要能覆盖热数据。延迟消息:时间轮精度 vs 扫描频率的 trade-off(秒级精度可接受则定时扫描更简单)。
  10. 失败与降级:副本同步失败/ISR 不足时的写入策略(拒绝写 vs 允许落后副本);Broker 宕机选主期间生产端重试与本地缓存;消费 Rebalance 期间的重复与暂停;磁盘打满时的拒绝策略与告警;自研组件必须有旁路降级——业务可切本地消息表或第三方 MQ。没有监控(lag、刷盘耗时、副本延迟、拒绝写次数)的自研 MQ 不允许上生产。

【原理溯源】

  • 为什么存储模型选“追加日志”而不是“逐条队列”? 追加写(append-only)把随机 IO 变成顺序 IO,磁盘顺序写接近内存速度;同时天然支持“按 offset 读取/回放”。传统“消息出队即删”的队列模型难以高吞吐持久化,也不好做多订阅回放。Kafka/RocketMQ CommitLog 都走这条路。
  • 为什么必须先有持久化再谈高可用? 副本复制的是“已持久化的日志段”。若消息只在内存,副本同步窗口内宕机即丢。顺序写 + 批量刷盘是在“延迟、吞吐、安全”之间的经典折中(异步刷盘丢最后几毫秒,同步刷盘更安全)。
  • ack 体系为什么必须两端都有? 生产端 ack 解决“发出是否被持久化”;消费端 offset/ack 解决“是否被业务处理”。只做一端就会在另一端丢或重。设计时要把 at-least-once 当默认语义,并把幂等责任明确交给消费方。
  • 分区模型为什么同时解决并行与顺序? 分区是并行单元(可多线程/多实例),又是顺序边界(分区内 FIFO)。同 key 路由同分区得到局部有序;分区数决定消费者组最大并行度。这是“吞吐与顺序”的平衡点,而不是二选一。
  • Rebalance 为什么难? 消费者增减时要把分区重新分配,期间可能:分区暂停消费(延迟)、重复消费(offset 交接)、顺序抖动。设计要尽量缩短 rebalance 时间(增量协作协议)、并要求消费幂等。自研时这是最容易出事故的模块之一。

【选型判断树】(自研能力优先级)

自研 MQ 能力路线图
├─ P0 最小闭环
│   ├─ 追加日志存储 + 持久化
│   ├─ 生产 ack / 消费 offset
│   └─ 分区 + 单消费者
├─ P1 生产可用
│   ├─ 副本与选主
│   ├─ Rebalance
│   └─ 死信 + 重试
├─ P2 业务增强
│   ├─ 延迟消息
│   ├─ 事务消息
│   └─ 回放 / 压缩 / 监控
└─ 先问:为什么不直接用 Kafka/RocketMQ?
    除非有强定制(合规、特殊存储、深度嵌入平台),否则自研 ROI 极低。

口诀:先日志后副本,先 ack 后高级特性;没有强理由不要自研。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“设计 MQ 的核心是追加日志 + 分区并行 + 分段确认”
0:30–2:30P0/P1存储、持久化、ack/offset、分区顺序、副本
2:30–3:30难点深挖选 1–2 个展开:顺序写为什么快、Rebalance 代价、事务消息半消息
3:30–4:30工程完整度死信、延迟、监控、权限、多租户
4:30–5:00收尾 + 反问“自研前先论证为何不用现成组件——多数时候答案是不该自研”

【关键数字】

参数经验值说明
顺序写吞吐数百 MB/s 级vs 随机写差 1–2 个数量级
批量刷盘攒批数十–数百条或数 ms平衡延迟与 IO
副本数2–3ISR 同步
默认语义at-least-once消费幂等必须
延迟精度离散档位(1s–2h 共 18 级)vs 连续(时间轮)档位实现简单,最细到 1s
Rebalance 目标秒级完成越长积压越久
单机写入目标5–10万 msg/s(异步刷盘)/ 1–3万(同步)1KB、3副本经验值
分区数≥ 峰值 / 单消费者能力 × 1.5例:5万/500×1.5=150
存储日量×体×副本×天提前估磁盘与保留策略

【追问链】(三层)

L1|“为什么顺序写比随机写快?” → 机械盘避免寻道;SSD 减少写放大与擦除。顺序写可近乎满带宽,随机写受 IOPS 限制。日志型存储把“写路径”全部变成 append。

L2|“offset 存在哪?” → 分两层。消费进度(committed offset)存在 Broker 内部的压缩型 topic __consumer_offsets(键为 group + topic + partition),老版本曾放 ZooKeeper,现在客户端与 broker 之间的 OffsetCommit/OffsetFetch 也走 broker;ZK 只留 broker/consumer 的注册与再均衡协调(新版也已搬到 broker)。消费者进程本地只持有下一次要拉的 nextFetch,真正「谁消费到哪」以 broker 上的提交记录为准。所以 auto.commit.interval.ms 决定写盘频率,崩溃时未提交的那段会被重放——这就是 at-least-once 的重复来源(与第 80 题同一因果链)。

L3|“事务消息为什么要半消息?” → 先发完整消息,本地事务失败就产生脏消息;先提交本地事务再发,宕机则丢消息。半消息对消费者不可见,等本地事务结果再 commit/rollback,未知状态靠回查兜底——用“暂存 + 二次确认”换原子性。

【评分标准】

档位答案特征
60 分能列生产-Broker-消费、持久化、ack
80 分有优先级路线图;讲清分区与顺序、副本
95 分深挖顺序写/offset/重试设计(独立重试队列+退避+次数上限+不可重试直投死信,顺序 topic 的重试须阻塞该 key)/半消息/Rebalance 之一;主动质疑“该不该自研”;覆盖监控与死信

【关联题】

  • 同簇: 第 86 题(Kafka 高吞吐)、第 80/81/82 题(可靠性)、第 92/93 题(事务)、第 87 题(消费模型)
  • 方法论: 第 99 题(选型思维)

【自测】

  1. 为什么 MQ 存储多用追加日志? 参考答案: 顺序写高吞吐、实现简单、支持按 offset 回放;出队即删的队列模型难以高吞吐持久化与多订阅。
  2. 消费端 offset 提交的正确时机? 参考答案: 业务处理成功之后。过早提交会丢,过晚会重复(可接受)。
  3. 判断对错:自研 MQ 应优先实现延迟消息和事务消息。 参考答案: 错。应先完成持久化、ack、分区、副本的最小闭环;高级特性建立在可靠核心之上。

86. 单机几十万 QPS 写入,Kafka 凭什么(Kafka 高吞吐原理) ​

【考察内容】Kafka 高吞吐原理是大数据/后端双高频题

【题目】压测显示 Kafka 单机写入能到几十万 QPS,而传统 MQ 只有几万。Kafka 靠什么做到?请从磁盘顺序写、页缓存、零拷贝、批量、分区并行几个角度解释。

【参考答案】多维度设计:

  1. 顺序写磁盘:消息 append-only 追加到分区文件末尾,顺序写接近内存速度(600MB/s 级),避开随机 IO;
  2. 零拷贝:消费时用 sendfile(内核态直接发到网卡),避免用户态拷贝;生产时页缓存(page cache)直接写;
  3. 批量与压缩:生产者批量发送(batch)、Broker 批量写、消费批量拉取;支持压缩(lz4/zstd)减少网络与磁盘;
  4. 分区并行:topic 分成多个分区,分区是并行单元,多磁盘/多线程并行处理;
  5. 页缓存而非应用缓存:利用 OS page cache,读写都命中缓存;数据已落盘,进程重启后 page cache 通常仍可命中;但 page cache 在机器重启/断电后必然丢失,此时需回读磁盘;
  6. 稀疏索引:每个 segment 只索引少量偏移,查找快,索引内存占用小;
  7. 常数级开销:O(1) 磁盘读写(append 与读 offset 附近数据),不受数据总量影响。
  8. 容量估算:Kafka 单机“几十万 QPS”通常指 小消息、批量、压缩、异步发送 条件。粗算:单 broker 顺序写磁盘可达数百 MB/s~GB/s 级;若消息 200B、批量后压缩比 5:1,有效写入放大后仍可支撑数十万条/s。集群规划:目标峰值 TPS ÷ 单 broker 有效 TPS × 副本系数 = broker 数;例目标 30 万 msg/s、单机有效 10 万、3 副本 → 先按写入吞吐算 broker 数:30万 ÷ 10万 = 3 台;再乘副本得到的是 9 个副本位(leader+follower 共 9 份数据/9 段写流量),这 9 个副本位要摊到更多机器上才不会被复制写放大拖垮——3 台是下限、按副本写流量预留到 5–9 台才是可落地的起步规模(先说清"单机有效 TPS"是否已含 follower 复制开销,两者会差 3 倍)。分区数:生产峰值 / 单分区吞吐,并保证 ≥ 消费者实例数。
  9. 失败与降级:开启高吞吐配置(acks=1、异步刷盘、大 batch)会提高丢消息风险——金融/交易链路必须回到 acks=all + 足够 ISR;磁盘/页缓存不足时顺序写优势消失,表现为吞吐骤降;生产端缓冲打满开始阻塞/丢弃,必须监控 buffer 与 request timeout。高吞吐场景的降级:允许日志类消息采样丢弃,交易类不可降级为“丢了也行”。

【原理溯源】

  • 为什么顺序写能接近内存速度? 避免了磁盘最贵的操作——寻道/随机写。HDD 顺序写可达数百 MB/s,SSD 更高;随机写受 IOPS 限制差 1–2 个数量级。Kafka 把所有写变成 segment 文件 append,把“磁盘很慢”的直觉改写成“磁盘顺序写很快”。
  • 零拷贝省了哪几次拷贝? 传统:磁盘 → 内核 page cache → 用户态缓冲 → Socket 缓冲 → 网卡。sendfile/transferTo 让内核直接从 page cache 送到网卡,跳过用户态往返,CPU 与延迟都下降。这是消费路径吞吐的关键。
  • 批量化为什么放大吞吐? 网络 RTT、系统调用、刷盘固定成本被摊到多条消息上。批量越大,单条消息成本越低(延迟上升)。压缩在批量之上进一步减少字节数——CPU 换带宽,通常划算。
  • 为什么用 OS page cache 而不是 JVM 堆内缓存? ① 堆内缓存受 GC 影响,大堆更痛;② page cache 进程重启后仍在(除非机器重启);③ 与 sendfile 无缝配合。代价:机器冷启动后要回盘,所以持久化与顺序写仍不可省。
  • 稀疏索引为什么够用? 消费是顺序近似顺序读,只需定位到大致 offset 附近再扫描少量数据。每 4KB 左右建一个索引点即可,索引可全部装入内存,查找 O(1) 级。
  • 与“传统 MQ 几万 QPS”差在哪? 传统 Broker 逐条路由、可能随机写、每消息一次确认、无批量/零拷贝深度优化。Kafka 把自己设计成分布式提交日志,用日志模型换吞吐,而不是功能丰富的逐条队列。

【选型判断树】

要更高写入吞吐
├─ 生产端
│   ├─ 开 batch.size / linger.ms
│   ├─ 压缩 lz4/zstd
│   └─ 幂等+适当 acks(all 会略降吞吐)
├─ Broker
│   ├─ 分区与磁盘分散
│   ├─ 页缓存充足(内存)
│   └─ 段文件与刷盘参数
└─ 消费端
    ├─ 批量拉取
    └─ 零拷贝(默认行为,勿破坏)

口诀:批量摊固定成本,顺序写打满磁盘,零拷贝省 CPU,分区换并行。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“Kafka 把自己做成分布式提交日志,用日志模型换吞吐”
0:30–2:30四件套顺序写、零拷贝、批量压缩、分区并行——每项一句机制一句收益
2:30–3:30页缓存与索引page cache vs 应用缓存;稀疏索引
3:30–4:30数据路径从生产者到消费者的路径串起来
4:30–5:00收尾“高吞吐来自放弃部分功能复杂度,这是架构取舍”

【关键数字】

参数经验值说明
顺序写带宽数百 MB/s 级远高于随机写 IOPS
单机写入数十万 msg/s(批量+压缩)与消息大小强相关
索引间隔约 4KB 一条索引项稀疏索引
批量linger.ms 5–20ms 常见用延迟换吞吐
压缩lz4 / zstd网络与磁盘双赢
segment默认 1GB滚动删除与索引
单机吞吐前提小消息+批量+压缩+顺序写参数一变数字差一个数量级
集群规模粗算峰值÷单机有效TPS×副本30万/10万=3 台下限,×3 副本=9 副本位(按复制写放大预留 5–9 台)
分区数≥ 消费者数,且匹配峰值分区=并行度上限

【追问链】(三层)

L1|“页缓存和直接写磁盘有什么区别?” → 写 page cache 立即返回,由内核异步刷盘;读常命中缓存。比每条消息 fsync 快得多。代价:掉电可能丢最后未刷盘数据(可配置刷盘策略)。

L2|“零拷贝对生产端有用吗?” → 基本不沾边,零拷贝吃的是读路径的红利。消费者(以及 follower 副本同步)从磁盘取数据时,Kafka 用 sendfile 让数据从 page cache 直接到网卡,省掉「内核→用户态→内核」两次拷贝和一次上下文切换;而生产端走的是写路径,收益来自顺序追加写 + page cache + 批量与压缩,与零拷贝无关。所以给生产端提吞吐要调 batch.size/linger.ms/压缩算法,而不是指望零拷贝;只有当 broker 需要把消息原样转发出去(消费、副本复制)时零拷贝才起作用——这也正是 Kafka「写一遍日志、多方各自读」的架构前提。

L3|“acks=all 会不会把吞吐打回原形?” → 会降低,因为要等 ISR 副本网络往返。但配合批量、幂等、足够 ISR,仍可保持很高吞吐。金融/核心业务值得用 all;日志类可 acks=1 换吞吐——这是可靠性与吞吐的显式 trade-off。

【评分标准】

档位答案特征
60 分背出顺序写、零拷贝、分区
80 分能解释批量化摊薄成本、页缓存角色
95 分画清生产→Broker→消费路径;说清 sendfile 省了什么;指出 page cache 掉电丢失与持久化关系;能谈 acks 对吞吐影响

【关联题】

  • 同簇: 第 84 题(选型)、第 85 题(MQ 设计)、第 87 题(消费模型)
  • 对比: 第 32 题(Redis 为什么快)——同属“单机为什么能扛”方法论

【自测】

  1. Kafka 零拷贝主要优化的是哪条路径? 参考答案: 消费路径。Broker 用 sendfile 将 page cache 数据直接送网卡,避免用户态拷贝。
  2. 为什么批量化能提高吞吐? 参考答案: 网络 RTT、系统调用、刷盘等固定成本被多条消息分摊,单条成本下降。
  3. 判断对错:用了页缓存就不用落盘了。 参考答案: 错。page cache 机器重启/断电会丢;仍需顺序写落盘保证持久性。

87. 组内加消费者想加速,结果速度没变(Kafka 消费者组) ​

【考察内容】Kafka 消费模型

【题目】Kafka 消费速度不够,你往消费者组里加了几台机器,结果速度几乎没变。为什么?请讲清消费者组与分区的关系:为什么同一个分区不能被组内多个消费者同时消费?并行度由什么决定?

【参考答案】

  1. 概念:分区是 Kafka 的并行存储单元;消费者组(group)内的消费者共同消费一个 topic,分区分配给组内消费者(一个分区分配给一个消费者,一个消费者可消费多个分区);
  2. 为什么同组不能多消费者共用一个分区:
    • 顺序性:分区内消息有序,若多消费者并发消费同分区,处理顺序无法保证;
    • offset 管理:分区进度(offset)以“消费者组+分区”为单位提交,多消费者并发会导致 offset 提交互相覆盖(A 处理到 10,B 处理到 5,提交错乱);
  3. 推论:分区数=组内最大并行度——消费者数超过分区数时,多余的消费者空闲;分区数决定吞吐上限;
  4. 不同组可以消费同一分区(各组独立 offset,互不影响)——用于多业务各自消费同一份数据;
  5. Rebalance:消费者增减/分区变化触发再平衡,分区重新分配(可能引起重复消费,需幂等)。
  6. 容量/并行度估算:消费者组内有效并行度 = min(消费者实例数, 分区数)。若 topic 有 8 个分区,拉起 20 个消费者也只有 8 个在干活,空转 12 个。目标吞吐反推:峰值 2 万 msg/s、单消费者稳定 2000/s → 需要 ≥10 个分区 + ≥10 个消费者。加机器前先看:分区数、单消费者处理 RT、是否有共享下游瓶颈(DB 连接池打满则扩消费者无效)。Rebalance 频率过高也会吃掉吞吐——扩缩容要平滑。
  7. 失败与降级:消费者实例故障触发 Rebalance,期间分区短暂无主,lag 上涨;处理:会话超时调合理、避免频繁滚动发布、用协作式再均衡。若下游容量有限,正确做法是消费端限流而不是无限加消费者。监控:每消费者分配到的分区数、lag、再均衡次数、单条处理耗时。发现“加机器不加速”先查分区分配再查下游。

【原理溯源】

  • 为什么并行度由分区数决定? 分区是 Kafka 的互斥消费单元:组内对每个分区只有一个活跃消费者。消费者只是“执行体”,没有更多分区可领就只能空转。所以加机器前先看 num.partitions。
  • 为什么 offset 以“组+分区”为单位? 消费进度必须能确定性地回答“这个组对这个分区读到哪了”。若多个消费者共写同一进度,会出现覆盖与回退。独占分区让进度语义简单可靠。
  • 不同组为什么可以并行消费同一分区? 各组维护独立 offset,互不干扰。这实现了“一份数据、多个业务各自消费”(如计费组、风控组、数仓组),是 pub/sub 的核心能力。
  • Rebalance 的代价是什么? 暂停消费(延迟)、可能从旧 offset 重放(重复)、分区迁移抖动。频繁 Rebalance(心跳超时、处理太慢)会恶化积压。优化:提高 max.poll.interval、减少单次处理耗时、使用静态成员/协作协议。
  • 单消费者多线程行不行? 可以提升处理并行,但必须自己保证:同分区(或同 key)消息有序与 offset 提交正确。常见做法:拉取线程 + 按 key 分片的工作线程池 + 按分区提交。复杂度高,先优先加分区。

【选型判断树】

消费慢如何加速?
├─ 消费者数 < 分区数 → 直接加消费者
├─ 消费者数 ≥ 分区数
│   ├─ 允许改 topic → 加分区(评估顺序影响)+ 加消费者
│   ├─ 不能改分区 → 单消费者内多线程(自行保序与提交)
│   └─ 消费逻辑本身慢 → 优化逻辑/批量/异步下游
└─ Rebalance 频繁 → 先修稳定性再谈扩容

口诀:分区是并行度上限;加机器不加分区常常无效。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“速度没变通常是因为消费者数已超过分区数”
0:30–2:30模型组内分区独占分配;一个消费者可多分区;分区数=最大并行度
2:30–3:30两个为什么顺序性 + offset 管理
3:30–4:30扩展路径加分区、多线程、优化逻辑;Rebalance 与幂等
4:30–5:00收尾“扩容前先看分区数与消费瓶颈位置”

【关键数字】

参数经验值说明
最大并行度= 分区数组内
加分区影响key 映射变化,可能乱序需评估
Rebalance 耗时秒级~数十秒期间可能停顿
常见配比消费者数 = 分区数再多无益
单消费者线程拉取 1 + 处理 NN 需自管顺序
有效并行度min(消费者数, 分区数)多出的消费者空闲
扩容前置条件分区数足够且下游能扛否则先扩分区/限流
反推分区峰值÷单消费者能力2万/2000=至少10分区

【追问链】(三层)

L1|“如何提高消费速度?” → 加分区 + 加消费者;优化单条处理耗时(批量、异步、缓存);确认下游不是瓶颈。不要只加机器。

L2|“为什么一个分区不能被组内多个消费者同时消费?” → 因为分区是组内并行的最小分配单位,消费进度就是一个单调递增的 offset。若组内两个消费者同时持有同一分区,就会出现两份独立进度、两次提交:谁都无法判断哪些消息已被处理,Rebalance 时两个 offset 还会互相覆盖,结果不是重复就是丢失。代价是组内最大并行度 = 分区数,消费者多于分区数的部分只会空闲。注意这条限制只作用于同一个组:不同消费组可以各自消费同一分区、各记各的 offset。

L3|“消费太慢导致 poll 超时被踢出组,怎么处理?” → 提高单次 poll 条数后的处理时间上限(max.poll.interval.ms)、减少单次处理耗时、改用异步处理并控制提交点。根因是“拉取循环被长处理阻塞”,而不是网络问题。

【评分标准】

档位答案特征
60 分说“要加分区”
80 分讲清组与分区分配、并行度=分区数
95 分解释顺序与 offset 两个原因;谈 Rebalance、多线程消费的自管成本

【关联题】

  • 同簇: 第 82 题(顺序)、第 83 题(积压)、第 86 题(吞吐)、第 81 题(Rebalance 重复)

【自测】

  1. 12 个分区,加到 20 个消费者,有几个在干活? 参考答案: 最多 12 个分到分区,其余空闲。
  2. 不同业务组能否消费同一分区? 参考答案: 可以。各组独立 offset,互不影响。
  3. 判断对错:消费慢就加消费者实例总能提速。 参考答案: 错。受分区数上限约束,且若瓶颈在下游或单条逻辑,加实例无效。

88. 消息重试了 10 次还是失败,一直卡住队列(死信队列) ​

【考察内容】死信队列是 MQ 可靠性治理高频题

【题目】有一条消息消费一直失败(比如数据格式非法),重试了 10 次还失败,占着队列头,后面的消息全被堵住。怎么处理?死信队列怎么设计?死信消息如何告警、人工介入、回放?

【参考答案】

  1. 概念:消息消费失败超过重试次数上限(如 3 次)后,转入死信队列(DLQ)——专门存放“处理不了”的消息,与正常队列隔离;
  2. 触发场景:业务异常(数据不合法)、依赖故障持续超时、消息格式错误;
  3. 处理机制:
    • RabbitMQ:死信交换机(DLX)+ 死信路由键,消息在队列中被拒绝/过期/超限后转发;
    • RocketMQ:重试队列(%RETRY%)消费失败自动重试,超过次数进死信队列(%DLQ%);
    • Kafka:无内置死信,需自建——消费失败写“失败 topic”或落库;
  4. 死信处理:
    • 告警通知:死信产生即告警(钉钉/短信),人工介入;
    • 死信消费端:读取死信,分析失败原因(日志/消息体),修复后重新投递(重建消息回原队列);
    • 定时重放:对可恢复的死信定时重试(如依赖恢复后);
    • 兜底:无法恢复的登记后人工/对账处理;
  5. 设计要点:死信消息要保留原始信息(原 topic、消息体、失败原因、重试次数)。
  6. 容量估算:若业务消息峰值 5000 msg/s,失败率 1%,则重试流量约 50 msg/s;若单消费者重试逻辑更慢(如人工接口 1s/次),死信与重试队列会快速堆积。重试队列容量按“峰值失败量 × 最大重试次数 × 消息体 × 保留天”估算。例:50/s×10 次×1KB,保留 7 天 ≈ 约 300GB 量级(50×10×1KB ≈ 500KB/s × 604800s ≈ 302GB,再乘副本)。生产上重试 topic 要与主 topic 隔离分区与消费者,避免重试流量拖垮主链路。
  7. 失败与降级:重试耗尽 → 进死信,主队列必须立刻继续消费,不能被单条毒消息卡死;死信可人工修复后重新投递,或对账补偿;对可丢弃的非核心消息,业务确认后可定期清理死信。监控死信堆积量与增长速率,超过阈值告警。禁止:无限原地重试、或为了“清死信”无限制并发重放把下游打挂。

【原理溯源】

  • 为什么要死信,而不是一直重试? 无限重试有三宗罪:① 毒消息(格式错误)永远失败,空耗资源;② 占住队列头/分区,阻塞后续消息(顺序消费时尤其致命);③ 重试风暴打挂下游。死信把“处理不了”从主通道剥离,主通道保持流动。
  • 为什么毒消息必须隔离? 数据非法类失败与“下游暂时超时”不同:重试一万次也一样。若不识别异常类型并快速进死信,会把瞬时故障策略误用到永久故障上。
  • Kafka 为什么没有内置 DLQ? Kafka 定位日志流,消费失败的语义交给应用。常见自建:失败发到 xxx.DLT topic、或写入失败表;由专门消费者处理。这也说明死信是消费端模式,不是存储端魔法。
  • 死信闭环为什么必须“告警 + 原因 + 可回放”? 只堆积不告警 = 静默数据丢失(业务以为在处理)。必须保留原 topic、消息体、异常栈、重试次数,才能修复后回放。回放要幂等,避免修好后重复副作用。
  • 与重试策略的关系? 死信是重试的终点。合理策略:可重试异常有限次指数退避;不可重试立即死信;总次数超限进 DLQ(见 90 题)。

【选型判断树】

消费失败
├─ 异常类型?
│   ├─ 参数/格式/业务非法 → 不重试,直接死信
│   ├─ 下游超时/临时故障 → 有限次退避重试
│   └─ 未知 → 先按可重试,但设较低上限
├─ 超过重试上限 → 死信队列
│   ├─ 立即告警
│   ├─ 保留上下文(原 topic、body、异常、次数)
│   └─ 人工/自动修复后回放(幂等)
└─ 监控:死信量、失败率、重试次数分布

口诀:毒消息快进死信;瞬时故障才重试;死信必须可告警可回放。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“死信是重试的终点隔离区,防止毒消息堵死主通道”
0:30–2:00机制三家 MQ 实现差异(Rabbit DLX、Rocket %DLQ%、Kafka 自建)
2:00–3:30处理闭环告警 → 分析 → 修复回放 → 人工兜底
3:30–4:30设计要点上下文完整、幂等回放、与重试策略衔接
4:30–5:00收尾“没有死信的重试系统,等于把事故藏起来”

【关键数字】

参数经验值说明
重试上限3–5 次(业务);Rocket 默认 16之后进死信
退避指数 1s→2s→4s…防重试风暴
死信告警非 0 即告警不要等积压
原文保留必须含原 topic + 异常否则无法回放
回放幂等必须修好后可能重复投递
重试流量峰值×失败率5000×1%=50 msg/s
死信/重试存储失败量×重试次数×体×天×副本注意与主 topic 隔离
毒消息策略快速进死信,不阻塞主队列主链路可用性优先

【追问链】(三层)

L1|“死信消息还能再消费吗?” → 能。分析修复后重新投递原队列或专用回放通道。回放消费必须幂等,因为业务可能已部分执行。

L2|“顺序消费时,一条失败进死信,后面的怎么办?” → 不能直接跳过——跳过就等于让后面的消息越过失败的那条,顺序当场就破了。顺序消费的正确做法是阻塞当前队列、原地重试:RocketMQ 的 OrderListener 失败后挂起该队列并按最大次数重试,重试耗尽才投死信并人工介入;Kafka 侧要保序就不能并发重试,通常把失败消息写入侧链并冻结该 key/分区后续消息的处理。工程上更常见的选择是放弃严格顺序,改用状态机前置校验 + 版本号(版本号更小的旧消息直接丢弃),让乱序变成可识别、可丢弃的,而不是错账。

L3|“死信量突然暴涨说明什么?” → 可能:下游故障、上游 schema 变更不兼容、脏数据洪峰、或错误地把可重试问题当毒消息。先看失败原因分布,再决定回滚变更、修下游还是批量清洗回放。死信是系统健康度信号,不只是垃圾场。

【评分标准】

档位答案特征
60 分知道失败多次进死信
80 分说清触发条件、三家实现、告警回放
95 分区分毒消息与瞬时故障;讲清阻塞 vs 跳过对顺序影响;死信作为监控信号

【关联题】

  • 同簇: 第 90 题(重试)、第 82 题(顺序)、第 83 题(积压)、第 96 题(监控)
  • 可靠: 第 80 题(不丢)

【自测】

  1. 为什么要设死信而不是无限重试? 参考答案: 防止毒消息永久失败、阻塞通道、造成重试风暴;把不可处理消息隔离出来人工治理。
  2. Kafka 的死信一般怎么实现? 参考答案: 无内置,消费失败写入失败 topic 或落库,由专门消费者告警与回放。
  3. 判断对错:死信队列里的消息可以忽略。 参考答案: 错。死信常代表数据丢失或业务未完成,必须告警、分析、回放或对账补偿。

89. 订单 30 分钟未支付自动关单,怎么定时触发(延迟消息) ​

【考察内容】延迟任务与 MQ 结合

【题目】订单 30 分钟未支付自动关闭。用 MQ 的延迟消息能力怎么做?RocketMQ 延迟消息怎么实现(延迟级别)、RabbitMQ 的 TTL+死信方案怎么搭?各自的延迟精度和限制是什么?

【参考答案】

  1. RocketMQ 延迟消息:内置延迟等级(1s/5s/10s/30s/1m/.../2h 等 18 级),发送时设 delayTimeLevel,到点投递——简单可靠,但只支持固定档位;
  2. RabbitMQ 延迟队列:TTL(消息过期时间)+ 死信交换机——消息进“延迟队列”设 TTL,过期后被死信路由转发到真正消费队列;TTL 可设任意时长(需注意队列级/消息级 TTL 语义差异);
  3. Kafka:无内置延迟,需自建:方案 A——发消息时带“目标执行时间”,消费端判断未到时间则重新投递/暂存(可用分层时间轮+重投);方案 B——时间轮(Netty HashedWheelTimer)在应用层管理;
  4. 对比:RocketMQ 延迟消息最省事;RabbitMQ 用 TTL+DLX 组合;Kafka 生态需自研或用时间轮;
  5. 生产建议:MQ 延迟消息为主 + 定时任务扫表兜底(延迟消息丢失/超时的订单由定时任务扫描“待支付超时”订单关闭),双保险。
  6. 容量估算:日订单 50 万、约 30% 会超时关单 → 日延迟消息约 15 万条,均值 ≈ 1.7 msg/s,峰值下单时(如晚 8 点)可能瞬时 数百 msg/s 触发关单扫描/消费——「数百」要带前提才算得出来:日 15 万条均匀只有 1.7/s,「晚 8 点这一小时占 20%」也才 8.3/s;要想到数百必须假设秒级突刺(如限量开抢:开抢 10 分钟内下单占日量 20% → 10 万单÷600s ≈ 167/s,延迟 30 分钟后集中到期)。所以消费能力按「突刺系数 × 峰值小时速率」配,并把突刺假设写出来,否则要么配空要么配爆。RocketMQ 延迟消息或 Redis ZSet+轮询都要按峰值到期量设计消费能力。存储要先定发法:方案 A 每单都发(下单即投一条延迟消息,到期再判是否超时)——日发 50 万条,在途量 ≈ 近 30 分钟下单量;方案 B 只对超时单发(日 15 万条)——但下单那一刻无法预知它会不会超时,所以 B 实际只能靠「先写状态 + 扫表兜底」实现,延迟消息省下的只是消息量,扫描与索引成本要加回来。本库正文与【关键数字】按 A 计;用 B 就把 15 万带进所有下游算式(消费能力、Broker 存储、扫描 QPS)保持一致。A 下 30min 窗口内的待投递量 ≈ 近 30 分钟的下单量,按本题日订单 50 万、峰值小时占 20% 折算 ≈ 10 万单/h ÷ 60 = 1670 单/min,30 分钟约 5 万条 在途(其中约 30% 会真正到期关单 ≈ 1.5 万条)。关单消费要幂等,且与支付回调竞争同一状态机。
  7. 失败与降级:延迟消息丢失/投递失败 → 依赖兜底定时扫描(扫描 created_at 超过 30min 且未支付的订单);MQ 延迟精度误差(RocketMQ 预设延迟级别非任意时间)→ 业务用“到期后再校验是否已支付”;消费者宕机 → 消息重投,状态机保证只关一次;支付与关单并发 → 数据库状态条件更新(仅当 status=待支付 才关单)。监控:延迟消息堆积、关单成功率、超时未关单对账差异。

【原理溯源】

  • 为什么延迟消息不能“Broker 睡眠 30 分钟”? Broker 要服务海量消息,不能为每条延迟消息挂起线程。实现要么到期时间索引 + 调度投递(RocketMQ 定时/延迟服务),要么暂存后过期转发(Rabbit TTL+DLX),本质都是“先隔离,到点再进入可消费通道”。
  • 为什么 RocketMQ 只给固定延迟级别? 用有限级别把延迟消息映射到内部调度结构,实现简单、吞吐高。任意精度需要更复杂的时间轮/持久化定时器。业务上 30 分钟可对齐到最近档位(如 30m 或用扫表精确关单)。
  • RabbitMQ TTL+DLX 的坑在哪? ① 队列级 TTL 与消息级 TTL 并存时,以先到期者为准,易误伤;② 消息只有到达队列头才检查过期,若前面有更长 TTL 的消息,后面的短 TTL 会被“卡住”——所以要每消息一队列或保证 TTL 单调。这是面试经典坑。
  • 为什么必须“延迟消息 + 扫表兜底”? 延迟消息可能丢(Broker 故障、消费失败、级别设错),而关单是资金相关正确性动作,不能只靠 MQ。定时任务扫描“创建超过 30 分钟且未支付”的订单做二次校验关闭,是最终一致的安全网。
  • 关单动作本身为什么要幂等? 延迟消息与扫表可能都触发关单;支付回调也可能并发。关单必须状态机幂等(仅“待支付”可关)。

【选型判断树】

延迟/定时任务怎么做?
├─ 精度要求
│   ├─ 分钟级可接受 → MQ 延迟消息
│   └─ 秒级/任意时刻 → 时间轮或调度中心(如 XXL-JOB)+ 扫表
├─ MQ 能力
│   ├─ RocketMQ → delayLevel
│   ├─ RabbitMQ → TTL + DLX(注意头阻塞)
│   └─ Kafka → 自建时间轮/重投或不用 MQ 延迟
└─ 一律加扫表兜底(资金/状态类)

口诀:延迟消息是快路径,扫表是正确性路径,两者都要。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“30 分钟关单是延迟投递问题,主流用 MQ 延迟消息 + 扫表双保险”
0:30–2:30三家实现Rocket 等级、Rabbit TTL+DLX、Kafka 自建
2:30–3:30坑与限制固定档位、TTL 头阻塞、Kafka 无内置
3:30–4:30工程闭环扫表兜底、状态机幂等、与支付回调并发
4:30–5:00收尾“能用内置不用自研;资金动作必须有第二路径”

【关键数字】

参数经验值说明
RocketMQ 延迟级别18 级,1s–2h固定档位
订单支付窗口常见 15–30 分钟与库存预占一致
扫表周期1–5 分钟兜底而非主路径
TTL 头阻塞短 TTL 可能被长 TTL 卡住Rabbit 经典坑
时间轮精度毫秒级内存方案,注意持久化
延迟消息量每单都发=日订单量(50 万/日,本题口径);只对超时单发=日订单×超时比例(15 万/日,需配扫表兜底)两式别混用
峰值到期下单峰×超时比例须给突刺假设:均匀只 1.7/s、晚高峰一小时占 20% 也才 8.3/s,「数百 msg/s」要按开抢级突刺算(10 分钟占日量 20% → ≈167/s)才成立

【追问链】(三层)

L1|“Kafka 说有延迟消息吗?” → 无内置。需消费端判断时间未到则重投/暂存,或应用层时间轮。多数团队在有 RocketMQ 时用 Rocket,或干脆用调度系统+扫表。

L2|“RabbitMQ 队列 TTL 和消息 TTL 混用有什么坑?” → 坑在过期检查只在队首触发。RabbitMQ 只判断队首消息是否过期,队列中间的消息即便 TTL 已到也不会被清理;一旦给每条消息设不同的 x-message-ttl,长 TTL 的消息压在队首,后面已过期的短 TTL 消息就迟迟不生效。队列级 TTL 是统一的、行为可预期;两者同时设置时按较短的算,但排查时很难看清是哪一层在生效。更要分清队列本身的空闲回收 x-expires(整个队列闲置多久被删)与消息 TTL 是两回事,面试里把它们混叫「队列 TTL」就会答错。实践口径:同一队列只用一种粒度且优先用队列级;要真延迟就换 RocketMQ 的延迟等级或外部时间轮(见本题 L1/L3 与第 92 题)。

L3|“延迟消息丢了,订单一直挂着怎么办?” → 扫表兜底:按创建时间扫描超时未支付订单,状态机校验后关闭并回滚库存。这就是为什么不能只依赖 MQ 延迟。同时对延迟消息本身做监控(发送量 vs 到期消费量)。

【评分标准】

档位答案特征
60 分知道 RocketMQ 有延迟消息
80 分三家方案都能说;知道 TTL+DLX
95 分指出固定级别限制与 TTL 头阻塞;强调扫表兜底与关单幂等

【关联题】

  • 同簇: 第 84 题(选型)、第 88 题(死信)、第 92 题(事务)、第 83 题(积压)
  • 业务: 订单状态机、库存回滚

【自测】

  1. RocketMQ 延迟消息的主要限制? 参考答案: 只支持固定延迟级别,任意时长需对齐或配合扫表。
  2. 为什么关单要扫表兜底? 参考答案: 延迟消息可能丢失或消费失败;关单影响库存与资金正确性,需要数据库扫描作为最终一致保障。
  3. 判断对错:RabbitMQ 中给每条消息设不同 TTL 就能实现任意延迟。 参考答案: 不完全对。存在头阻塞与队列/消息 TTL 交互问题,需谨慎设计。

90. 消费者一抛异常就无限重试,队列被堵死(消费重试策略) ​

【考察内容】消息消费的容错设计

【题目】消费者处理消息抛异常后,如果无限重试会把队列堵死;不重试又会丢消息。重试策略怎么设计才合理(重试次数、退避间隔、区分可重试/不可重试异常、超过次数进死信)?

【参考答案】

  1. 失败分类处理:
    • 可重试错误(下游超时、临时故障):自动重试,采用指数退避(1s、2s、4s...)或固定间隔,设最大次数(3~5 次);
    • 不可重试错误(消息格式错误、业务数据非法):重试无意义,直接进死信/记录告警(避免无限重试浪费);
  2. 实现方式:
    • RocketMQ:默认重试 16 次(延迟等级递增),可配最大重试次数;
    • RabbitMQ:消息 reject/nack + 重试队列(或延迟重试);
    • Kafka:消费端手动控制——捕获异常后不提交 offset,等待重投;或用“重试 topic”(失败消息发到重试 topic,延迟后重新消费);
  3. 关键防坑:
    • 重试风暴:大量消息同时失败重试,打爆下游——指数退避+熔断(下游故障时暂停消费,恢复后再继续);
    • 重试消息乱序:顺序消息重试要阻塞后续(或单线程消费);
    • 重试与幂等配合:重投的消息消费时要幂等(去重);
  4. 监控:失败率、重试次数分布、死信量告警。
  5. 容量/重试风暴估算:假设峰值消费 2000 msg/s,其中 5% 下游异常 → 100 msg/s 进入重试;若策略是“立即无限重试”,等效把失败流量以乘数放大,可能在秒级打满下游与线程池。合理预算:最大重试 3~5 次、指数退避(1s/2s/4s)、重试 topic 与主 topic 隔离。线程池:重试消费与主消费隔离,重试线程数 ≤ 下游可承受并发的一定比例(如 20%)。观察指标:失败率、重试队列 lag、下游 RT/错误率。
  6. 失败与降级:识别不可重试异常(参数错误、业务拒绝)→ 直接死信/失败回调,不再重试;可重试异常(超时、限流、短暂不可用)→ 退避重试;达到上限 → 死信 + 告警 + 补偿任务。业务侧开关:极端时可暂停该 topic 消费,先保下游。原则:重试是放大器,必须限次、限流、可熔断。

【原理溯源】

  • 无限重试为什么会堵死队列? 顺序消费/单线程下,毒消息失败后原地重试,后续消息无法前进;即使多线程,持续重试也占用消费容量并放大下游压力。重试解决的是瞬时故障,不是永久错误——不分类就必然踩坑。
  • 为什么要指数退避? 下游故障恢复需要时间。固定短间隔会在下游最脆弱时持续打击(重试风暴)。指数退避 + 抖动能拉开间隔,给恢复窗口,降低惊群。
  • 为什么要有熔断/暂停消费? 当失败率飙升,说明下游整体不可用,此时“继续消费继续失败”毫无意义,还占用资源。暂停消费(消息留在 Broker)是背压:把压力挡在 MQ 外,恢复后再追。比死磕重试安全。
  • Kafka 为什么重试要更“手工”? 不提交 offset 即重投,但会阻塞该分区后续消息;更常见是失败转发到重试 topic(带延迟),主 topic 继续。RocketMQ 把重试队列做成产品能力,Kafka 交给应用。
  • 重试与顺序、幂等的三角关系? 保序 → 失败必须阻塞;重试 → 必然可能重复 → 必须幂等;要吞吐 → 不能无限阻塞。设计时三选二的组合决定架构:交易常用“阻塞重试 + 幂等”,日志常用“跳过/死信 + 最终一致”。

【选型判断树】

消费失败
├─ 异常分类
│   ├─ 不可重试(格式/业务非法)→ 立即死信
│   ├─ 可重试(超时/限流/临时)→ 指数退避,上限 3–5 次
│   └─ 下游整体故障(失败率高)→ 熔断暂停消费
├─ 超过上限 → 死信 + 告警
├─ 顺序要求?
│   ├─ 要求保序 → 阻塞重试,不跳过
│   └─ 无顺序 → 可异步重试 topic
└─ 全程幂等

口诀:先分类再重试;退避防风暴;熔断防雪崩;死信保通道。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“重试策略核心是异常分类,而不是次数”
0:30–2:30分类与参数可重试/不可重试;次数、指数退避
2:30–3:30三家实现Rocket 16 次、Rabbit nack、Kafka 不提交/重试 topic
3:30–4:30三大坑重试风暴、乱序、与幂等配合
4:30–5:00收尾“重试是补偿,不是正确性来源;正确性靠幂等与死信闭环”

【关键数字】

参数经验值说明
业务重试上限3–5 次之后死信
RocketMQ 默认16 次,等级递增可改
退避1s,2s,4s… 可加抖动防惊群
熔断阈值失败率 >50% 持续 N 秒按下游 SLA
重试 topic 延迟10s–5min给下游恢复时间
失败重试流量峰值×失败率2000×5%=100 msg/s
重试次数3–5 次 + 死信无限重试=事故

【追问链】(三层)

L1|“下游持续故障怎么办?” → 熔断暂停消费,消息在 MQ 保留;恢复后限流爬坡继续。不要无限重试打下游。

L2|“顺序消息重试怎么不乱序?” → 关键是重试不能把消息扔回公共池与其它消息竞争。要点四条:① 按 key/队列串行重试(RocketMQ 顺序消息用本地阻塞重试 + 挂起当前队列,而不是异步重投);② 重试期间该 key 的后续消息不得先行处理,越过失败者就是乱序;③ 达到重试上限才转死信,并保留原始位点(topic/queue/offset)以便按序回放;④ 消费端仍要做状态机或版本号校验,把晚到的旧消息识别并丢弃,作为最后兜底。

L3|“重试次数和延迟怎么定?” → 参考下游恢复时间分布:多数秒级故障用 3 次指数退避覆盖;分钟级依赖故障靠熔断而不是更长重试。重试总时长应 < 业务可接受延迟,超过就进死信并告警。

【评分标准】

档位答案特征
60 分知道要设置最大重试次数
80 分区分可重试/不可重试;指数退避;进死信
95 分谈重试风暴与熔断;顺序阻塞;与幂等联动;给监控指标

【关联题】

  • 同簇: 第 88 题(死信)、第 81 题(幂等)、第 82 题(顺序)、第 95 题(流控背压)

【自测】

  1. 消息格式错误为什么要直接进死信? 参考答案: 永久性失败,重试无意义且浪费资源、可能阻塞通道。
  2. 指数退避解决什么问题? 参考答案: 避免在下游故障期高频重试造成重试风暴,给恢复时间。
  3. 判断对错:重试次数越多越可靠。 参考答案: 错。过多重试放大负载、延迟暴露问题;应分类 + 上限 + 死信 + 告警。

91. 核心链路依赖的 MQ 突然不可用(MQ 故障降级) ​

【考察内容】依赖故障的容灾设计

【题目】MQ 集群故障不可用,而核心下单链路依赖它做异步化。业务不能停,怎么办?有没有降级方案(开关切同步、本地暂存、DB 轮询兜底)?恢复后积压数据怎么处理?

【参考答案】

  1. 认知:MQ 是核心依赖,要先做高可用(集群、副本、监控)降低挂的概率;但必须准备“MQ 挂掉”的降级预案;
  2. 降级方案:
    • 本地消息表兜底:业务与消息同库同事务写入本地消息表,MQ 恢复后定时任务把积压消息补发(消息表是“永不丢失”的持久层)——最稳妥;
    • 同步调用降级:MQ 不可用时,异步链路降级为同步调用(如通知/积分直接同步执行,牺牲 RT 保正确性)——适合低频轻量操作;
    • 服务降级:非核心功能(推荐、统计)直接关闭,核心链路(下单、支付)优先保障;
    • 消息落本地文件/磁盘队列:轻量场景可写本地队列文件,恢复后重放(有丢失风险);
  3. 恢复流程:MQ 恢复后,先补发积压(本地消息表/重放),再恢复消费,监控积压量回落;
  4. 预案要求:降级方案要提前演练(混沌工程),不能事发时临场设计。
  5. 容量估算(降级通道):MQ 全挂时,若下单峰值 5000 TPS、每单平均 2 条异步事件,则本地消息表要承接 1 万行/秒 写入——这通常不现实,必须配合:① 只保留核心事件(订单状态),非核心关闭;② 本地表批量写入;③ 业务限流到 DB/本地表可承受范围(如 500–1000 TPS)。口径必须连着说清,否则自相矛盾:题干峰值 5000 TPS 限到 1000 即只放行 20%,所以【关键数字】里「降级期间 ≥99%」只能指被放行的那部分流量的成功率,其余是排队或明确拒绝;不能同时承诺「全量下单 99% 成功」和「限流 1000 TPS」。降级窗口内用户侧表现:主流程可下单,积分/通知延迟;恢复后补发速度按下游容量限流,避免二次打挂。
  6. 失败与降级:预案分三级——① MQ 部分节点故障:自动切其他 broker/集群;② 集群整体不可用:核心链路切本地消息表/事务性 outbox,异步功能降级关闭;③ 长时间不可用:人工确认哪些副作用可丢、哪些必须对账补齐。生产端发送超时要 fail-fast 并写本地,不能让线程池阻塞在 MQ 客户端。演练:定期做 MQ 摘除演练,验证本地表与补发任务。

【原理溯源】

  • 为什么“核心下单”不能同步依赖 MQ? 同步依赖会把 MQ 的可用性直接乘进下单成功率。MQ 再高可用也是额外故障域。正确架构:主事务成功与否只取决于 DB/支付,MQ 只承载可异步、可补偿的副作用。
  • 本地消息表为什么是最佳兜底? 它把“要发的消息”写进与业务同一个本地事务,业务成功则消息记录必然存在。MQ 挂了只是“投递暂停”,不是“消息丢失”。恢复后扫表补发,正确性由数据库事务保证(见 93 题)。
  • 为什么“切同步”要慎用? 同步化会把下游 RT 与错误率拉进主链路,MQ 故障时可能从“异步延迟”变成“下单超时”。只适合低频、轻量、高可用的下游;高 QPS 时应关闭非核心而非全部同步。
  • 恢复期为什么要限流补发? 故障期间堆积的发送/消费会在恢复瞬间形成洪峰,可能再次打挂 MQ 或下游。正确顺序:验证 MQ 健康 → 按速率补发生产 → 消费端限流爬坡 → 观察 lag 与错误率。
  • 降级开关为什么必须提前演练? 事发时没有时间设计与验证。开关粒度(按功能/按接口)、自动触发条件(MQ 健康检查失败 N 次)、回滚条件都要预置。混沌工程定期注入 MQ 故障,验证下单成功率。

【选型判断树】

MQ 不可用
├─ 下单主链路是否依赖 MQ 才能成功?
│   ├─ 是(架构问题)→ 立即改:主事务不依赖 MQ
│   └─ 否 → 启动降级
├─ 降级动作
│   ├─ 非核心异步(推荐/统计)→ 直接关闭
│   ├─ 核心副作用(积分/通知)→ 本地消息表暂存,恢复补发
│   └─ 低频关键 → 评估切同步
└─ 恢复
    ├─ 补发限流爬坡
    ├─ 对账确认
    └─ 演练归档

口诀:主链路不绑 MQ;消息表保正确;恢复要限速。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“先保证主链路不依赖 MQ 同步可用,再谈降级”
0:30–2:30降级矩阵本地消息表、切同步、关非核心、本地暂存
2:30–3:30恢复流程补发顺序与限流爬坡
3:30–4:30预案工程开关、监控、混沌演练
4:30–5:00收尾“容灾不是加机器,是失败路径可走”

【关键数字】

参数经验值说明
健康检查失败阈值连续 3–5 次触发降级开关
补发速率恢复稳态的 20%→100%防二次洪峰
本地消息表保留至少覆盖最大故障窗口如 7 天
演练频率季度/半年混沌工程
下单成功率目标降级期间 ≥99% 核心路径口径=被放行的那部分流量;与下一行「限流到 500–1000 TPS」不能同时承诺「全量下单 99% 成功」
降级写入预算峰值事件×核心比例5000×2→限流到500–1000 TPS
补发限流≤ 下游容量防恢复后二次故障

【追问链】(三层)

L1|“MQ 挂了订单会不会丢单?” → 若主事务不依赖 MQ,订单在 DB,不丢。异步副作用进本地消息表,恢复后补发。若设计上“不发消息就不能提交订单”,那是架构缺陷,必须改。

L2|“本地消息表会不会把数据库拖垮?” → 会,这正是降级方案最容易忽略的一点。本参考答案第 5 条算过:MQ 全挂时若下单峰值 5000 TPS、每单 2 条异步事件,本地消息表就要承接 1 万行/秒的写入,远超单库能力。所以降级必须同时做三件事:① 只保留核心事件(订单状态变更),非核心直接关闭;② 批量合并写入,攒 100~500 条一次 insert;③ 把业务限流到本地表能承受的量级(500~1000 TPS),宁可排队也不能压垮 DB。恢复后按 (status, 时间) 索引分批补发并限速爬坡,避免补发本身形成二次洪峰。

L3|“如何自动发现 MQ 故障并自动降级?” → 生产探活(发送心跳消息/检查集群状态)+ 错误率熔断。连续失败达到阈值自动切降级开关,并告警。恢复后自动或人工切回。开关要灰度,避免抖动来回切。

【评分标准】

档位答案特征
60 分说“多部署节点做高可用”
80 分给出本地消息表/切同步/关非核心降级
95 分强调主链路不依赖 MQ;恢复限流;预案演练与自动开关

【关联题】

  • 同簇: 第 93 题(本地消息表)、第 79 题(价值与代价)、第 83 题(积压)、第 96 题(监控)
  • 理论: 第 97 题(CAP)、第 98 题(BASE)

【自测】

  1. MQ 完全不可用时,本地消息表如何保证消息不丢? 参考答案: 消息与业务同一本地事务落库;MQ 只是投递通道,恢复后定时任务扫“待发送”补发。
  2. 为什么不适合把所有异步都改成同步降级? 参考答案: 同步会把下游 RT/故障引入主链路,高并发下可能导致下单超时雪崩;应只保留轻量低频,其余暂存或关闭。
  3. 判断对错:降级方案上线前不用演练,真挂了再按文档操作。 参考答案: 错。未演练的预案不可信;需混沌工程验证开关、补发、回切。

92. “先发消息还是先写库”才能两全(RocketMQ 事务消息) ​

【考察内容】事务消息是 RocketMQ 标志性能力,分布式事务专题高频题

【题目】下单要“写订单库 + 发消息给下游”,先写库可能消息没发出去,先发消息可能库没写成。RocketMQ 事务消息怎么解决这个两难?半消息、本地事务、回查机制分别是什么?

【参考答案】

  1. 解决的问题:先发消息再提交事务(事务失败消息白发了)/先提交事务再发消息(消息发送失败业务成功但下游不知道)——本地操作与发消息无法原子;
  2. 半消息机制:
    • 生产者发送“半消息”(half message)——Broker 收到但不投递给消费者;
    • 执行本地事务(如订单落库);
    • 根据本地事务结果向 Broker 提交 commit(消息可见可投递)或 rollback(删除消息);
  3. 事务状态未知兜底(关键):若生产者本地事务执行中挂了/超时,Broker 长期持有半消息——Broker 会回查(check)生产者(回调 checkLocalTransaction),生产者查询本地事务状态后回复 commit/rollback;
  4. 保证:半消息+提交/回滚+回查,使“本地事务成功⇔消息最终可见”达到最终一致;
  5. 与本地消息表对比:事务消息把“消息表”搬到 MQ 内部,业务少维护一张表,但回查逻辑仍需业务实现。
  6. 容量估算:事务消息比普通消息多一次半消息写入 + 状态回查,吞吐通常低于普通发送(经验值约为普通 sync 的 50%–80%,视回查比例)。若交易峰值需发送 5000 msg/s 事务消息,按 80% 效率折算为 5000÷0.8=6250、按 50% 折算为 10000,故集群容量要按 6250–10000 普通消息当量预估。回查压力:本地事务未决比例越高,回查 QPS 越高;状态表要有索引且查询足够快(按 transactionId/msId)。半消息在 Commit 前对消费者不可见,会增加端到端延迟约一次确认往返(通常毫秒~十毫秒级)。
  7. 失败与降级:半消息发送成功但本地事务失败 → Broker 回查,业务返回回滚,半消息丢弃;本地事务成功但确认消息失败 → Broker 定时回查,业务返回提交,消息最终投递;回查服务不可用 → 依赖事务状态表 + 超时策略,极端时告警人工。若 RocketMQ 不可用 → 降级本地消息表。禁止:无状态表导致回查悬挂/误提交。

【原理溯源】

  • 为什么“先写库”和“先发消息”都不对? 两者不在同一事务域,中间任一崩溃都会造成:库成功消息丢,或消息在库回滚后变成脏消息。这是分布式写“无法两阶段原子”的根本问题,必须引入暂存 + 二次确认或本地同事务记录。
  • 半消息解决了什么? 把消息变成“对消费者不可见的预写”。只有 commit 后才投递,rollback 则删除。这样“发送”动作可以发生在本地事务之前而不产生脏消息。
  • 回查为什么是关键? 生产者可能在本地事务后、commit 前崩溃,Broker 若一直等就永久悬挂半消息。回查让 Broker 主动问“你本地事务到底成没成”,生产者查库答 commit/rollback。没有回查,事务消息就只是“两阶段的一半”。
  • 与本地消息表的同构性? 本质同一思想:把“待发消息”与业务绑定在同一一致性边界。消息表在业务库,由扫表投递;事务消息在 MQ,由回查确认。前者兼容所有 MQ,后者依赖 RocketMQ 能力但少维护一张表。
  • 消费端失败与事务消息无关。 事务消息只保证“生产端本地事务与消息可见”的一致;消费失败仍走重试/死信/幂等。不要混为一谈。

【选型判断树】

如何保证“业务+发消息”原子?
├─ MQ 是否支持事务消息(RocketMQ)?
│   ├─ 是 → 半消息 + 本地事务 + 回查
│   └─ 否 → 本地消息表(Kafka/Rabbit 通用)
├─ 是否强一致到消费完成?
│   └─ 否(最终一致即可)→ 上述方案足够
│   └─ 是 → 还需消费侧对账/TCC,事务消息不够
└─ 回查必须实现:查本地库真实状态,不可写死 commit

口诀:半消息暂存,回查兜底,消费另算。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“两难本质是本地事务与发消息跨系统,需要暂存+确认”
0:30–2:30时序半消息 → 本地事务 → commit/rollback → 回查
2:30–3:30为什么一致每一步失败都有路径:rollback 删消息;未知靠回查
3:30–4:30对比与边界vs 本地消息表;消费失败另论
4:30–5:00收尾“事务消息解决的是生产端一致,不是端到端 exactly-once”

【关键数字】

参数经验值说明
回查次数默认约 15 次间隔递增
回查间隔10s→30s→1m→…→2h量级参考
半消息对消费者不可见commit 后才投递
本地事务超时秒级完成为佳过长占用半消息
适用一致性最终一致非消费端强一致
事务消息吞吐约为普通 sync 的 50%–80%含半消息与回查开销
回查 RT毫秒~十毫秒依赖状态表查询性能
端到端延迟多一次确认/回查窗口通常仍远小于同步跨服务事务

【追问链】(三层)

L1|“本地事务成功但 commit 消息丢了怎么办?” → Broker 超时后回查,生产者查库发现成功则回 commit,消息最终投递。这就是回查存在的意义。

L2|“回查时生产者服务也挂了呢?” → 半消息会一直停在「未提交」状态,Broker 按 transactionCheckMax(默认 15 次)周期性回查;生产者始终不在线、回查次数耗尽时,Broker 默认回滚这条半消息——宁可不发一条本地事务状态未知的消息。由此有三个必做配套:① 回查接口要幂等且真能查到本地事务状态(按 transactionId/msId 建索引,别在回查里做重活);② 生产者重启后要能立刻恢复回查能力;③ 业务侧必须有对账兜底——「本地事务已提交但消息被回滚」只能靠对账发现并补发,不能指望 MQ 自己纠正。

L3|“事务消息和本地消息表怎么选?” → 有 RocketMQ 且团队接受其 API/运维 → 事务消息,少维护表;需要多 MQ 统一、或已有消息表基建 → 本地消息表。正确性等价,工程取舍不同。

【评分标准】

档位答案特征
60 分知道有半消息
80 分完整时序:半消息→本地事务→提交/回滚
95 分讲透回查;对比本地消息表;说明消费端不在保证范围

【关联题】

  • 同簇: 第 93 题(本地消息表)、第 99 题(分布式事务选型)、第 80 题(不丢)、第 58 题(跨库一致)
  • 理论: 第 98 题(BASE)

【自测】

  1. 半消息与普通消息的区别? 参考答案: Broker 暂存且不对消费者投递,直到 commit;rollback 则删除。
  2. 回查机制解决什么故障? 参考答案: 生产者在本地事务后未能发送 commit/rollback 时,Broker 主动查询本地事务状态以完成提交或回滚。
  3. 判断对错:用了事务消息,消费端就不会重复消费。 参考答案: 错。事务消息只管生产端一致性;消费仍可能重复,需幂等。

93. 不用事务消息,怎么保证“业务和发消息”不丢(本地消息表) ​

【考察内容】本地消息表是分布式事务最经典落地方案

【题目】如果 MQ 不支持事务消息,怎么保证“本地业务操作”和“发消息”要么都成功要么都补偿?本地消息表方案怎么设计?定时任务怎么扫描投递?重复投递怎么幂等?

【参考答案】

  1. 核心思想:业务操作与消息记录在同一个本地事务中完成,消息不丢;投递失败可重试;
  2. 流程:
    • 生产者:开启本地事务 → 执行业务操作(如订单表插入)→ 同时写消息表(msg 状态=待发送)→ 提交事务(原子保证);
    • 发送服务:扫描消息表中“待发送”消息(定时任务)→ 投递 MQ → 收到 ack 后更新状态=已发送;
    • 消费者:消费消息处理业务(幂等),处理成功返回 ack;
  3. 可靠性保证:
    • 不丢:消息与业务同事务落库,投递失败/宕机后定时任务继续扫发(最终一定发出);
    • 不重:消费者幂等(唯一键)+ 发送端状态机(已发送不重发,或重发但消费幂等);
  4. 缺点:每套业务要多维护消息表、扫表有延迟、消息表与业务表耦合;与 RocketMQ 事务消息是同一思路的两种实现(一个在业务侧、一个在 MQ 侧);
  5. 适用:无事务消息功能的 MQ(Kafka/RabbitMQ)场景,或已有消息表基础设施。
  6. 容量估算:本地消息表扫描线程按轮询间隔 T 与批量大小 B 设计。若峰值需投递 2000 msg/s,轮询间隔 500ms,则每批需投递约 1000 条——批量发送比逐条快,但单批不宜过大(避免长时间事务与内存)。表增长:日订单 50 万×平均 2 消息 = 100 万行/日;保留 7 天 ≈ 700 万行,要按时间分区/归档。扫描索引:(status, next_retry_time) 必须有,否则全表扫描会拖垮 DB。
  7. 失败与降级:MQ 暂不可用 → 消息停留在待发送,任务退避重试,不影响业务主事务;消息表所在库压力大 → 批量发送、异步扫描、读写分离;补发风暴 → 限流发送;业务与消息不一致(极少数)→ 对账任务按业务 ID 比对,缺失则补偿发送,多发靠消费幂等。监控:待发送积压量、最老消息年龄、发送失败率。

【原理溯源】

  • 原子性从哪来? 唯一来源是数据库本地事务。把业务行和消息行放进同一事务,要么都提交要么都回滚,消除了“业务成功但消息没记下”的窗口。这比“先调 MQ 再写库”或反之都强。
  • 为什么还要扫表投递? 事务只保证“消息已记录”,不保证“已送到 MQ”。投递是网络 IO,可能失败。扫表把投递变成可重试的最终动作:状态=待发送就继续发,直到 ack。
  • 状态机为什么重要? 待发送→发送中→已发送→(可选)已确认。防止并发扫表重复发送;发送中需超时回退,避免卡死。多实例扫表要用分布式锁或按 id 分片。
  • 重复投递为什么必然存在? 发送成功但更新状态前崩溃 → 下次扫到仍会再发。所以消费端必须幂等——与 MQ at-least-once 同一逻辑,不能假设“发一次”。
  • 与事务消息对比。 同构:半消息 ≈ 待发送行;回查 ≈ 扫表。事务消息把复杂度收到 MQ 内部;消息表通用但要自己维护表、任务、清理。

【选型判断树】

MQ 无事务消息
└─ 本地消息表
    ├─ 同事务写业务+消息行
    ├─ 扫表投递 + 状态机
    ├─ 消费幂等
    └─ 清理已成功消息(TTL/归档)
有 RocketMQ 且愿意绑定
└─ 事务消息(92 题)

口诀:同库事务保不丢,扫表保投递,幂等保不重。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“本地消息表用数据库事务绑定业务与待发消息,再扫表投递”
0:30–2:30流程写表 → 扫描投递 → ack 更新状态 → 消费幂等
2:30–3:30为什么可靠原子性来自同事务;最终发出来自扫表
3:30–4:30工程细节状态机、锁、清理、延迟
4:30–5:00收尾“与事务消息同构,是通用底座”

【关键数字】

参数经验值说明
扫表间隔1s–1min延迟 vs 压力
分页大小100–500防大事务
发送中锁超时30s–2min卡死回退
已发送保留1–7 天再清理/归档
消费幂等键业务唯一键非 msgId
批量投递峰值×轮询间隔2000/s×0.5s≈1000条/批
消息表增长日订单×均消息数50万×2=100万行/日
关键索引(status, next_retry_time)无索引会全表扫

【追问链】(三层)

L1|“把发消息写在事务外会怎样?” → 事务提交后、发消息前崩溃,消息永久丢失。必须同事务写消息表。

L2|“多个实例同时扫表会不会重复发?” → 若每个实例都各自扫同一批 status=待发送 的行,就会重复发。防重靠「抢占」而不是「各扫各的」:① UPDATE ... SET status='处理中', owner=?, version=version+1 WHERE status='待发送' LIMIT n(或 SELECT ... FOR UPDATE SKIP LOCKED)先抢到再发;② 或把消息表按 id % N 静态分片给不同实例,互不重叠;③ 接收端仍要保持幂等(消息 ID/业务键唯一索引),因为整条链路仍是至少一次投递。另外扫表要限批量并走 (status, 时间) 索引,否则消息表本身会变成慢查询源。

L3|“消息表越来越大怎么办?” → 定期删除/归档已成功且超过保留期的数据;大表可分区;扫表务必走索引(status + 时间)。

【评分标准】

档位答案特征
60 分知道要有个消息表
80 分同事务写入 + 扫表投递 + 状态机
95 分讲清原子性来源、重复投递必然、清理与锁;能对比事务消息

【关联题】

  • 同簇: 第 92 题(事务消息)、第 91 题(降级)、第 99 题(选型)、第 110 题(幂等)

【自测】

  1. 本地消息表的原子性靠什么保证? 参考答案: 业务数据与消息记录在同一数据库本地事务中提交。
  2. 为什么消费端仍要幂等? 参考答案: 扫表可能重复投递,MQ 也可能重复,必须保证业务效果只发生一次。
  3. 判断对错:本地消息表可以完全替代对账。 参考答案: 错。消费失败、业务部分成功等仍需对账兜底。

94. 一条消息 ack 没配对,消费者重启后消息全丢了(ack 机制) ​

【考察内容】MQ 投递语义体系

【题目】有同事把消费者的 auto-ack 打开后,消费者一处理完(甚至还没处理完)消息就被确认,重启后消息丢失。请讲清生产端 ack 与消费端 ack 的机制,ack 时机配错分别会带来什么问题?

【参考答案】

  1. 生产端 ack(Broker 确认收到):
    • Kafka:acks 参数——acks=0 不等确认(可能丢);acks=1 leader 落盘即确认(leader 挂丢数据);acks=all 等所有 ISR 副本确认(最可靠,延迟最高);
    • RabbitMQ:publisher-confirm,Broker 落盘后回调 ack,未确认/nack 可重发;
    • 意义:生产端确认“消息已持久化”,避免“发了就当成功”导致丢失;
  2. 消费端 ack(消费者确认处理完):
    • 自动 ack:收到消息即 ack——若处理中崩溃,消息已确认=丢失(默认禁用安全场景);
    • 手动 ack:处理成功才 ack;处理失败不 ack(消息重新投递);
    • Kafka:offset 提交机制(enable.auto.commit=false,处理成功后手动提交);
  3. ack 丢失场景:
    • 生产 ack 丢失:消息实际已到 Broker 但生产者没收到 ack → 重发 → 重复消息(消费端幂等解决);
    • 消费 ack 丢失:处理成功但 ack 没送达 → 消息重投 → 重复消费(幂等解决);
  4. 结论:ack 机制解决“丢”与“不丢”,但引入了“重复”——所以幂等是 MQ 可靠性闭环的最后一块拼图。
  5. 容量/确认开销估算:手动 ack 会让消费 RT 从“拉取即可”变成“处理完才确认”。若单条业务处理 10ms、单线程,则上限 100 msg/s/线程;要扛 5000 msg/s 需要 50 线程或 50 分区并行。同步刷盘/acks=all 会增加生产与 Broker 开销,吞吐下降但更安全。缓冲:in-flight 消息过多会导致 Rebalance 时重复面变大;建议限制 max.poll.records 与处理批量,使单次 poll 在会话超时内处理完。
  6. 失败与降级:处理失败不 ack,依赖重投;多次失败进死信;消费者重启后从 last committed offset 继续,可能重复不会丢;若错误地 auto-ack/先 commit 后处理,则崩溃窗口会丢消息——这是本题要消灭的 bug。监控:commit offset 与 end offset 的 lag、消费失败率、死信量。发布时注意 graceful shutdown:处理完 in-flight 再退出,减少重复。

【原理溯源】

  • ack 的本质是什么? 分布式中的进度确认。生产 ack = “Broker 已持久化”;消费 ack = “业务已处理”。两者确认的对象不同,配错会系统性丢消息。
  • auto-ack 为什么会丢? 确认发生在“收到”而非“处理成功”。崩溃窗口在收到与处理完成之间,消息却被视为已消费。手动 ack 把确认点移到副作用之后,崩溃则重投。
  • Kafka acks=all 为什么更安全? 等待 ISR 中所有副本写入后再向生产者确认,leader 故障时数据仍在其他副本。acks=1 只等 leader,acks=0 不等——性能递增,安全性递减。
  • 为什么 ack 可能带来重复? 确认信号本身会丢。生产端超时重发、消费端未收到 ack 结果的重投,都是“不确定就重试”。at-least-once 由此而来。
  • 幂等为什么是最后一块拼图? 不丢(确认尽量晚)与不重(确认尽量早)互斥。系统选择不丢,重复必须由消费幂等吸收,否则“不丢”换来的重复会破坏业务正确性。

【选型判断树】

ack 怎么配?
├─ 消息丢失可容忍?(日志/统计)
│   ├─ 是 → 可 auto-ack + acks=1,换吞吐
│   └─ 否 → 手动 ack + 处理后提交 + acks=all
├─ 生产端
│   ├─ 核心业务 → sync + confirm/acks=all + 重试
│   └─ 高吞吐日志 → 可异步 + 批量
└─ 无论哪种,消费端尽量幂等

口诀:确认点在副作用之后才不丢;不丢必可能重。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“ack 分生产与消费两端,配错分别导致丢或重”
0:30–2:30两端机制Kafka acks 三档;confirm;auto vs manual ack
2:30–3:30因果链早 ack→丢;晚 ack→重;重复靠幂等
3:30–4:30事故复盘式本题 auto-ack 导致重启丢失的机理
4:30–5:00收尾“ack 是可靠性开关,不是细节”

【关键数字】

参数经验值说明
Kafka acks0/1/allall 最安全
enable.auto.commit生产建议 false手动提交
autoAck核心业务 falseRabbitMQ
重复率网络抖动时可见必须幂等
ISR副本同步集合acks=all 等待对象
单线程消费上限1/单条RT10ms→100 msg/s
并行需求峰值÷单线程能力5000/100=50 并行(分区/线程)
确认时机业务成功后auto-ack 仅日志类可容忍

【追问链】(三层)

L1|“生产 ack 和消费 ack 分别确认什么?” → 生产:Broker 是否持久化;消费:业务是否处理完成。两端都要正确配置。

L2|“acks=all 会慢多少?” → 慢在多等一次「ISR 副本拉取并确认」的网络往返加选主判定:生产端延迟通常从本地写盘的亚毫秒级涨到几毫秒到几十毫秒。但吞吐不一定同比例掉——Kafka 靠批量摊薄这次等待:linger.ms/batch.size 开起来后,acks=all 相比 acks=1 的吞吐一般只降 10%~30%;若逐条同步发送则会掉一个数量级。所以结论是:开批量与幂等生产者、保持足够大的 ISR,核心业务用 acks=all(再配 min.insync.replicas=2、unclean.leader.election.enable=false)这个代价是值得的。

L3|“offset 提交和 ack 是什么关系?” → Kafka 中提交 offset 即消费进度确认。应处理成功后再 commit;auto-commit 定时提交可能在处理中提交,导致崩溃丢消息——与 auto-ack 同类错误。

【评分标准】

档位答案特征
60 分知道不要 auto-ack
80 分分清两端 ack;Kafka 三档;手动提交
95 分讲清“早确认丢、晚确认重”;幂等闭环;offset≈ack

【关联题】

  • 同簇: 第 80 题(不丢)、第 81 题(幂等)、第 90 题(重试)
  • 对照: 第 92/93 题(生产端一致)

【自测】

  1. auto-ack 在什么时点确认? 参考答案: 收到消息时(或按配置极早确认),处理中崩溃即丢失。
  2. Kafka acks=1 与 all 的差别? 参考答案: 1 只等 leader;all 等全部 ISR 副本,更安全更慢。
  3. 判断对错:ack 越早系统越高效,应尽量早确认。 参考答案: 错。过早确认会丢消息。应在业务成功后确认,用幂等处理重复。

95. 消息每秒 1000 条,下游接口只扛得住 100(消费端流控) ​

【考察内容】消费端容量治理

【题目】消费端每处理一条消息要调一个下游接口,该接口每秒最多承受 100 次调用;但消息流入量每秒 1000 条。怎么设计消费端(限速、分批聚合、削峰填谷、背压)才能既不丢消息又不压垮下游?

【参考答案】

  1. 需求本质:消费速率必须与下游处理能力匹配,防止把下游打挂、防止自己 OOM;
  2. 方案:
    • 限流消费:消费端加限流器(Semaphore/Guava RateLimiter),按下游阈值控制处理速率;
    • 批量消费:一次拉取/确认多条(Kafka 批量拉取、RabbitMQ prefetch 设置),批量处理(如攒 100 条调一次下游批量接口);
    • 控制并发:消费者并发数(线程数/实例数)与下游容量匹配,不盲目扩;
    • 背压/暂停:下游 RT 升高或错误率上升时,熔断暂停消费(消息留在 MQ),恢复后继续——比“继续消费硬重试”安全;
    • 预取限制:RabbitMQ 设 prefetch count(如 1~10),避免一次性拉太多到本地内存;
  3. 监控:消费 lag、处理耗时、下游错误率联动告警;
  4. 特殊场景:大促削峰——允许积压(MQ 天然缓冲),下游恢复后追平,但要设积压告警上限。
  5. 容量估算(步进):上游 1000 msg/s,下游仅 100 QPS,则缺口 900 msg/s 积压;若持续 1 小时,积压 ≈ 900×3600 = 324 万条。按消息体 1KB、副本 3,约 9.7GB。消化策略:① 消费端限流 100/s(积压线性涨,需扩下游)——但「既不丢又不垮」不能默认成立:MQ 的缓冲是有界的,恒定 900/s 缺口下 1 天就是 7776 万条,按 1KB×3 副本≈239GB/日,超过消息保留期或磁盘水位就会丢最老消息。要能给出「可容忍积压时长 = min(消息保留期, 可用磁盘 ÷ 日增量)」,超时必须合并处理/扩下游/按分级丢弃非核心,并把该时长本身做成告警阈值;② 下游扩容到 500/s 后积压降速;③ 非核心消息丢弃/合并(如同订单多条事件合并为 1 条)。批量消费:下游若支持批量接口,按每批 100 条 × 10 次/s = 等效 1000 条/s,正好吃掉 1000 msg/s 的流入,且只占 100 QPS 预算的 10%,其余留给重试。RPS 预算:重试流量要计入,避免 100 QPS 被重试吃掉。
  6. 失败与降级:限流阈值动态调整(下游 RT/错误率联动);下游过载 → 熔断消费或降低并发;积压超阈值 → 告警并启动扩容/丢弃非核心;禁止“不限流冲垮下游”。降级目标排序:资金/订单状态 > 通知 > 统计埋点。恢复时爬坡放量。

【原理溯源】

  • 为什么不能“有多少消费多少”? 消费速率若超过下游容量,结果是:下游过载雪崩、错误率暴涨、重试风暴,最终比“主动限速”更慢更糟。流控是保护下游以保总吞吐。
  • 积压为什么反而是正确状态? 在下游容量固定时,多出来的 900/s 必须存在某处。MQ 的磁盘缓冲就是为此设计——允许 lag 上涨,好过打挂下游。这是削峰填谷的消费侧镜像。
  • 批量为什么有效? 把 N 条消息合并为一次下游调用(若下游支持批量接口),或至少摊薄固定开销。1000/s 单条调用改为 10 次×100 条批量,下游压力与 RT 都下降。
  • 背压与限流的区别? 限流是稳态按阈值控制;背压是动态在下游异常时暂停。两者叠加:平时 RateLimiter 100/s,错误率升高时熔断暂停。
  • prefetch 为什么重要? 预取过大 → 本地堆积 → 内存压力,且崩溃时未 ack 的在途消息要整批重投(重复面变大、恢复变慢);过小 → 吞吐不足。按处理时延与并发调优,本质是“在途消息上限”。

【选型判断树】

生产 1000/s,下游 100/s
├─ 下游能否批量接口?
│   ├─ 能 → 聚合批量,提升有效吞吐
│   └─ 不能 → 硬限流 100/s,接受积压
├─ 积压是否可接受?
│   ├─ 可 → MQ 缓冲 + 告警阈值
│   └─ 不可 → 与业务谈降级(丢非核心/合并/近实时采样)
└─ 保护
    ├─ RateLimiter / 信号量
    ├─ 熔断背压
    └─ prefetch 控制在途

口诀:速率对齐下游;多余流量留给 MQ;异常时背压。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“消费速率必须 ≤ 下游容量,差额由 MQ 积压吸收”
0:30–2:30组合拳限流、批量、并发控制、prefetch
2:30–3:30背压下游异常时暂停而非硬重试
3:30–4:30积压策略可接受积压 + 告警;大促场景
4:30–5:00收尾“流控是保护下游,不是限制自己”

【关键数字】

参数经验值说明
下游限流按其压测阈值 80%留余量
Rabbit prefetch1–10 起步按 RT 调
批量大小下游批量接口上限内如 50–200
lag 告警持续上涨 5–10min或绝对阈值
熔断条件错误率/RT 双指标联动
积压速率上游−下游1000−100=900 msg/s
1小时积压缺口×3600≈324万条,约9.7GB(1KB×3副本)
等效扩容批量大小×每秒调用次数100条/批×10次/s=等效1000条/s

【追问链】(三层)

L1|“暂停消费会不会丢消息?” → 在缓冲的边界内不会:消息在 Broker 持久化,恢复后继续。但积压是有上界的——按消息保留期与磁盘水位算,恒定缺口(如 900/s ≈239GB/日×副本)超过保留期或写满磁盘就会丢最老消息。所以要说「在容忍时长内不丢」,并给出该时长=min(保留期, 可用磁盘÷日增量);同时消费端内存不要无限预取。

L2|“积压越滚越多怎么办?” → RocketMQ/Kafka 都可在消费端设置并发与拉取速率;更常见的是业务层限流:信号量/令牌桶限制同时调用下游的数量,配合批量接口降低有效 RPS。扩容下游时按 20%→50%→100% 爬坡。若下游是第三方(短信/支付),严格按对方限流协议配置,并预留重试预算(通常只用对方配额的 70%–80%)。

L3|“如何动态感知下游容量变化?” → 自适应限流:根据下游 RT/错误率调节并发与速率(类似 TCP 拥塞控制)。或配置中心热更新阈值。核心是把下游健康信号接进消费调度。

【评分标准】

档位答案特征
60 分说“加机器消费”
80 分限流 + 批量 + 控制并发
95 分背压熔断、prefetch、「允许积压」要带边界(积压有界:保留期/磁盘水位+容忍时长,超限必须扩下游或按分级丢弃)、动态调参

【关联题】

  • 同簇: 第 83 题(积压)、第 90 题(重试熔断)、第 87 题(并行度)、第 96 题(监控)

【自测】

  1. 为什么 1000/s 进、100/s 出时不应该全速消费? 参考答案: 会打挂下游并引发重试风暴;应限流到下游容量,积压留在 MQ。
  2. 批量消费的前提是什么? 参考答案: 下游支持批量或单次调用固定开销可被摊薄;且业务允许微批延迟。
  3. 判断对错:消费端 prefetch 越大吞吐越高越好。 参考答案: 错。过大导致内存压力与故障恢复变慢,应按处理能力设置。

96. MQ 上线后要盯哪些指标,才能提前发现风险(MQ 监控) ​

【考察内容】MQ 可观测性

【题目】MQ 接入核心链路后,运维要建立监控大盘。要盯哪些指标才能提前发现风险(积压、消费停滞、broker 负载、延迟)?哪些指标出现什么趋势就说明要出事了?

【参考答案】

  1. 积压量(Lag):消费 offset 落后生产 offset 多少——积压持续上涨=消费能力不足或下游故障,是最核心指标;按 topic/分区维度监控;
  2. 消费速率与生产速率:对比判断是否失衡;
  3. 消费失败率:失败率高=业务异常/消息异常,联动重试与死信监控;
  4. 死信队列量:死信持续产生=有系统性不可恢复问题,告警人工介入;
  5. 吞吐与延迟:生产/消费 TPS、端到端延迟(生产到消费的耗时);
  6. Broker 健康:CPU、磁盘(分区文件增长)、内存、网络、副本同步状态(ISR 收缩=副本异常,数据安全风险)、消费组 Rebalance 次数与时长;
  7. 重试次数分布:大量消息频繁重试=下游不稳定;
  8. 告警设计:积压量超过阈值(如 10 万)告警;失败率突增告警;死信非零告警;磁盘水位 75% 预警、85% 严重(与第 9 条同口径)。
  9. 容量相关监控阈值示例:生产 TPS 突增 > 基线 3 倍告警;消费 lag > 10 万或持续上涨 5 分钟告警;单 broker 磁盘水位 > 75% 预警、> 85% 严重;副本同步延迟 > 5s 关注;生产失败率 > 0.1% 告警;死信增长速率突变告警。容量巡检:按周评估峰值 TPS vs 集群水位,提前扩容分区/broker。大盘应同时展示:流入、流出、lag、失败、磁盘、CPU——只看 TPS 会漏掉“在堆积”和“在失败”。
  10. 失败与降级:监控本身挂了要有备用(多采集通道/云监控);告警要分级(P0 电话/P1 IM/P2 日报),避免告警疲劳;预设应急 runbook:积压怎么扩容、MQ 挂了怎么切本地消息表、磁盘满怎么处理。演练:季度故障演练验证监控能否在 5 分钟内发现问题。没有 runbook 的告警只是噪音。

【原理溯源】

  • 为什么 Lag 是第一指标? 它是生产与消费速率差的积分,直接对应“业务还要等多久”。Lag 持续上涨 = 系统在积累延迟债;Lag 为 0 但失败率高 = 消息可能在“假成功”。
  • 为什么要同时看速率与 Lag? 只看 Lag 会误判:短暂峰值 Lag 涨但 P>C 已恢复会回落;若 C 长期 < P,Lag 线性上涨。速率对比能区分“瞬时”与“趋势”。
  • 死信为什么必须非零告警? 死信是“最终处理不了”的信号。若只监控消费成功量,可能 lag 正常但业务数据已丢。死信与失败率是质量指标,lag 是进度指标。
  • ISR 收缩为什么危险? 同步副本减少,acks=all 可用性下降,leader 故障时丢数据或不可用风险上升。这是数据安全的前置信号,应在丢消息之前发现。
  • 端到端延迟与处理耗时的区别? 端到端含排队(积压);处理耗时是单条业务逻辑。前者高可能是积压,后者高是代码/下游慢——处置完全不同。
  • 监控的层次? ① 进度(lag、速率);② 质量(失败、死信、重试);③ 资源(CPU、磁盘、网络、ISR)。缺一层就有盲区。

【选型判断树】

监控大盘最小集
├─ 进度:生产 TPS、消费 TPS、Lag(topic/分区)
├─ 质量:失败率、重试分布、死信量
├─ 时延:端到端延迟、单条处理耗时
└─ 资源:Broker CPU/磁盘/网络、ISR、副本滞后

告警分级
├─ P0:死信非 0、ISR 收缩、磁盘 >85%、消费停滞
├─ P1:Lag 持续上涨、失败率突增
└─ P2:处理耗时缓慢劣化

口诀:进度质量资源三层都要;lag 为 0 未必安全。

【口述骨架】(5 分钟)

时间任务内容
0:00–0:30定性“监控分进度、质量、资源三层,Lag 最核心但不够”
0:30–2:30指标清单Lag、速率、失败、死信、延迟、Broker、ISR
2:30–3:30趋势解读哪些组合意味着要出事
3:30–4:30告警设计阈值与分级
4:30–5:00收尾“监控是为了在用户投诉前发现”

【关键数字】

参数经验值说明
Lag 告警绝对值(如 >10 万)或持续上涨 5–10min按业务 SLA 调
磁盘水位75% 预警 / 85% 严重预留清理时间,配合保留策略
死信非 0 告警人工介入
ISR收缩即告警数据安全
消费失败率突增或 >1%–5%看基线(正文第 3 条)
端到端延迟超 SLA 告警如 >30s
生产失败率>0.1% 告警确认链路是否静默丢(正文第 9 条)
TPS 突增>基线3倍可能大促或异常重试

【追问链】(三层)

L1|“Lag 为 0 就安全吗?” → 不一定。可能消费失败进死信、或消息被错误丢弃、或根本没生产。要结合生产量、失败率、死信一起看。

L2|“如何区分消费慢和下游挂了?” → 先把第 5 条的两个指标拆开看:端到端延迟(生产到消费,含排队)与单条处理耗时。端到端涨但单条耗时正常=消息在排队,是消费能力不足;单条耗时本身涨,才往下游查。第二步看质量指标:消费失败率(第 3 条)与重试次数分布(第 7 条「大量消息频繁重试=下游不稳定」)一起飙升、错误类型集中在超时/连接拒绝,就是下游故障;只有耗时高、失败率平,是消费自己慢(线程池、GC、慢 SQL)。第三步排除中间件本身:ISR 收缩、磁盘水位、Rebalance 次数(第 6 条)异常时问题在 MQ。三者对应的处置完全不同——扩分区与消费者、熔断降级保护下游、修 Broker。

L3|“大盘之外还要什么?” → 链路追踪(msgId 贯穿)、审计(谁发的什么)、容量预测(磁盘增长、分区水位)、混沌演练验证告警有效性。监控是可观测性的一部分,不是全部。

【评分标准】

档位答案特征
60 分只提积压量
80 分Lag + 失败率 + Broker 资源
95 分进度/质量/资源三层;ISR 与死信;趋势解读;lag=0 的陷阱

【关联题】

  • 同簇: 第 83 题(积压)、第 88 题(死信)、第 90 题(重试)、第 95 题(流控)、第 80 题(不丢)
  • 运维闭环: 第 91 题(降级)

【自测】

  1. 为什么死信量要单独告警? 参考答案: 死信代表最终处理失败,可能造成业务数据丢失;仅看 lag 会漏掉“消费成功但进死信”的假阴性。
  2. ISR 收缩意味着什么? 参考答案: 同步副本减少,数据可靠性下降,leader 故障时可能丢消息或不可用,需立即处理。
  3. 判断对错:只要消费 TPS 高就说明系统健康。 参考答案: 错。可能伴随高失败率或大量重试;需结合质量指标。


持续学习,持续积累。