第四章 消息队列(第79-96题)
79. 架构评审:引入 MQ 的理由与代价(MQ 价值与代价)
【考察内容】MQ 核心价值与引入代价(解耦/异步/削峰)
【题目】架构评审会上,你提出要引入消息队列,有人反对说“多一个组件多一份故障”。请讲清楚:引入 MQ 到底解决什么痛点(削峰/解耦/异步),不引入会怎样,以及引入后新增了哪些风险与成本?
【参考答案】三大核心作用:
- 解耦:A 系统产生订单事件,需要通知 B/C/D 系统——同步调用会“新增一个下游就改一次代码”;通过 MQ,A 只发消息,下游自行订阅,新增下游不动 A 的代码;
- 异步:下单链路同步调用“发短信、送积分、更新推荐”导致接口 RT 长——MQ 异步化后主链路只做核心操作,RT 从 2s 降到 200ms;
- 削峰:秒杀瞬间 10 万请求,直接打到 DB 会挂——MQ 缓冲,消费端按 DB 可承受的速率慢慢消费(如每秒 1000),保护下游;
- 代价:引入新组件(运维成本)、消息可靠性问题(丢失/重复/乱序需要治理)、最终一致性(业务要接受异步延迟)。
- 容量估算(步进):大促峰值 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。
- 失败与降级: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)——异步最终一致的理论基础
【自测】
- 判断对错:只要 QPS 高,就该上 MQ 削峰。 参考答案: 错。若消费端能力与生产端匹配、无突发峰值、且需要同步结果,上 MQ 只增加复杂度。削峰只解决“峰值 ≫ 稳态”问题。
- 架构评审时,反对者说“多一个组件多一份故障”,你怎么回应? 参考答案: 承认代价,再算收益:MQ 把故障域从“全链路同步串联”变成“主链路短 + 异步可重试”;同时给出可靠性三段防护与降级预案,证明引入后整体可用性可更高,而不是简单多一个单点。
- 为什么说解耦是“发布事实”而不是“调用接口”? 参考答案: 发布事实只声明事件与 schema,不绑定下游地址与存在性;新增订阅方无需改生产者。接口调用则把下游身份硬编码进上游,形成编译期耦合。
80. 业务方投诉“消息丢了”,订单状态没推进(消息不丢失)
【考察内容】消息不丢失是 MQ 最高频题
【题目】业务方投诉:异步任务经常“消息丢了”,比如订单支付成功但积分没到账。消息在“生产者发送→Broker 存储→消费者消费”三个阶段都可能丢失。请给出每个阶段不丢的配置与机制(ack、持久化、重试、确认机制)。
【参考答案】逐环节防护:
- 生产端:
- Kafka:acks=all(等所有 ISR 副本确认)+ 重试机制;同步发送并在失败时记录告警;
- RabbitMQ:开启 publisher-confirm(Broker 返回 ack/nack,nack 重发);
- RocketMQ:同步发送(syncSend),不丢才返回成功;
- Broker 端:
- 持久化:Kafka topic 副本数≥2 +
min.insync.replicas=2+ unclean.leader.election 关闭(缺 mir 时 acks=all 只需 leader 落盘即算成功,leader 崩仍丢已 ack 消息);RabbitMQ 队列持久化 + 消息 delivery_mode=2;RocketMQ 刷盘策略(同步刷盘更强); - 集群高可用(副本机制),防止单节点宕机丢数据;
- 持久化:Kafka topic 副本数≥2 +
- 消费端:
- 手动 ack(而不是自动):业务处理成功后才提交 offset/ack,处理失败不 ack,消息重新投递;
- 关闭“收到即确认”(RabbitMQ autoAck=false);
- 兜底:发送失败写本地重试表/MQ 自身重试;消费失败进死信队列人工处理;对账(业务侧核对消息与结果)。
- 容量估算(步进):日订单 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,但换来不丢。
- 失败与降级:生产端 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 acks | 0 / 1 / all | all 最安全;1 在 leader 故障时可能丢 |
| Kafka 副本数 | ≥2,金融 ≥3 | ISR 同步副本 |
| RabbitMQ delivery_mode | 2(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 要什么就得显式开什么”
【自测】
- Kafka 的
acks=1在什么情况下会丢消息? 参考答案: leader 写入成功并 ack,但 follower 尚未同步时 leader 宕机且发生非干净选举,已 ack 消息丢失。acks=all 可避免(等 ISR 全部写入)。 - 为什么消费端要“处理成功再 ack”? 参考答案: 若收到即 ack,处理中崩溃则消息被标记成功不再投递 → 丢失。处理后再 ack,崩溃时未 ack 会重投 → 最多重复,不丢。
- 判断对错:只要 Broker 做了副本,消息就一定不丢。 参考答案: 错。还需生产端确认、消费端正确 ack、禁止不干净选举,否则任一段仍可能丢。
81. 网络抖动导致同一条消息被消费两次(重复消费)
【考察内容】消息重复消费是 MQ 必考题
【题目】网络抖动或消费者宕机重平衡,导致同一条消息被投递了两次,业务被重复处理(比如重复发奖)。怎么保证“消息只被处理一次”?消费端的幂等设计怎么做?
【参考答案】
- 认知:重复消费在至少一次(at-least-once)投递语义下必然存在(ack 丢失、重试、Rebalance、进程崩溃),只能靠“消费幂等”解决;
- 幂等方案:
- 业务唯一键 + 数据库唯一索引:处理前先插入唯一键(如 order_id),重复插入报错则跳过——最可靠;
- 状态机幂等:消息处理推进状态机(如订单:待支付→已支付),重复消息发现状态已推进则忽略;
- Redis SETNX 去重:以消息 ID/业务 ID 为 key,SETNX 成功才处理(配合过期时间),Redis 挂了由 DB 唯一索引兜底;
- 消息表去重:本地记录已处理消息 ID 表,查重后处理;
- 关键:幂等要基于业务唯一键(订单号、支付流水号)而非消息本身 ID(同业务可能产生多条消息);
- 设计原则:消费逻辑本身要“可重入”(同参数执行多次结果一致)。
- 容量/一致性窗口估算:幂等存储(DB 唯一索引/Redis SETNX)要按峰值消费 TPS 设计——若峰值 5000 msg/s,Redis SETNX 需能扛同量级写(单实例通常 10 万+ ops,足够);DB 唯一索引插入在热点业务键上可能成瓶颈,建议按业务键分表或走 Redis 预检 + DB 最终落库。去重 key 的 TTL:至少覆盖最大重试窗口 + 时钟偏差 + 业务最长处理时间,常见 24h~7 天;过短会在延迟投递时失效导致重复。
- 失败与降级: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/s | Redis 轻松;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 题(发送端也可能重复)
【自测】
- 为什么说重复消费“无法避免,只能幂等”? 参考答案: 因为 ack 与业务副作用不在同一原子单元;at-least-once 协议下,崩溃/网络问题会导致“业务已做但未确认”从而重投。
- 幂等键为什么不用 msgId? 参考答案: 同一业务事件可能对应多条 msgId(生产端重发),用 msgId 会漏重;业务唯一键才能表达“只应发生一次”。
- 判断对错:Redis SETNX 成功就可以不落库唯一索引。 参考答案: 错。Redis 非持久、可能丢、与业务非原子,只能前置过滤;正确性必须靠 DB 唯一约束或状态机。
82. 订单的“创建→支付→发货”消息乱序了(消息顺序性)
【考察内容】消息顺序是 MQ 三大可靠性问题之一
【题目】订单系统的状态流转依赖消息顺序:必须先消费“创建”、再“支付”、再“发货”,一旦乱序订单状态就错。消息队列怎么保证这种顺序性?Kafka/RocketMQ 分别怎么实现?代价是什么?
【参考答案】
- 认知:全局有序成本极高(只能单分区单消费者,吞吐骤降),生产上只做局部有序(同一业务 key 的消息有序);
- 方案:
- Kafka:同一业务 key(如 order_id)作为消息 key,hash 到同一个分区(同一分区内消息有序);消费者组内一个分区只能被一个消费者线程消费,单线程按序处理;
- RocketMQ:顺序消息(MessageQueueSelector 按业务 id 选队列)+ 单消费者;
- RabbitMQ:同一条队列 + 单消费者(队列天然 FIFO,但要避免多消费者并发消费同一队列);
- 代价:同一 key 的消息串行处理,吞吐受限——拆 key 粒度(用户级/店铺级)平衡;热点 key 可在允许乱序的业务维度上拆(如按用户/店铺分队列);不能按「订单子流程」拆——那会把同一个 order_id 的创建/支付/发货打散到不同队列,直接破坏题干要求的同单有序;真要提速只能改成状态机版本号+条件更新来容忍乱序(另见原理溯源的旁路 topic 方案,同样必须配版本裁决);
- 失败处理:消息处理失败不能跳过继续(会乱序),需重试阻塞或进死信并暂停后续(业务补偿);
- 兜底:对账/重放机制纠正极端乱序。
- 容量与并行度估算:订单状态消息峰值假设 1 万 msg/s,若 order_id 哈希到 N 个分区,则有序消费最大并行度 = N。要保证同一订单串行:同 order_id 固定分区;不同订单可并行。若单消费者处理 500 msg/s,则需要至少 20 个分区 + 20 个消费者实例才能扛 1 万/s。代价:热点商户/热点订单单分区可能过热,需在业务键中加入合理分散因子(如 seller_id+order_id)——前提是同一订单的创建/支付/发货三条消息上该因子恒等(同一卖家、同一订单),否则同一 order_id 会因字段取值不同落到不同分区,顺序照样破;所以分散因子只能用下单时就固定、与事件类型无关的稳定字段(如订单号自带的分片号),绝不能用「消息类型」「当前状态」这类会变的属性。也不能破坏“必须有序的最小键”。
- 失败与降级:单分区内消息失败 → 该分区阻塞(保序必然牺牲局部可用);处理策略:① 可跳过的非关键乱序消息进死信;② 关键状态走状态机幂等,旧状态消息直接忽略;③ 极端时对该分区做“人工跳过 + 对账修复”。禁止为提速把同 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,吞吐无法支撑业务;实际只需同业务实体内有序。
- 同一订单的三条消息 hash 到不同分区会怎样? 参考答案: 跨分区无顺序保证,可能先消费支付再消费创建,状态错乱。必须保证同 order_id 同分区。
- 判断对错:顺序消息消费失败应该跳过,先处理后面的。 参考答案: 错。跳过会破坏顺序,应阻塞重试或进死信并暂停该 key 后续消息。
83. 一觉醒来 MQ 积压了几百万条消息(消息积压)
【考察内容】消息积压是生产事故高频题
【题目】凌晨上线的一个 bug 让消费者全挂了,早上发现 MQ 积压了几百万条消息,业务链路大面积延迟。先恢复消费还是先修 bug?积压的消息怎么快速消化(扩容、临时消费者、跳过)?会不会挤垮下游?
【参考答案】
- 先定位原因:
- 消费能力不足(消费者实例少、消费逻辑慢);
- 下游故障(DB/依赖服务挂了,消费一直失败重试);
- 生产端突发流量(大促)超过消费速率;
- 处理流程:
- 修复下游故障:如果是下游问题,先恢复下游(否则扩容也没用);
- 第 0 步:先回滚或热修(题干根因就是「凌晨上线的 bug 让消费者全挂」):扩容、临时 topic、转发都跑在同一份有 bug 的代码上,消费照样失败,重试还会把下游再压一遍;所以下列手段都以「消费代码已能正常处理」为前提;
- 临时扩容消费者:增加消费者实例(注意:Kafka 单分区只能一个消费者消费,分区数=最大并行度——可临时增加分区并重新分布,或用“多线程消费”);更快的办法:把积压消息转发到新建的临时 topic(分区数扩大),用更多消费者并行消费,消费完再处理结果;
- 降级策略:非核心消息(如日志、统计)在积压严重时可跳过/丢弃(业务可接受);核心消息不能丢;
- 限流生产端:暂停/减缓非关键生产,先清积压;
- 根治:优化消费逻辑(批量消费、异步化)、评估分区数与消费者数匹配、监控堆积量告警(堆积超过阈值自动告警/扩容);
- 关键防坑:不要盲目重置位点、或用新消费组从
auto.offset.reset=latest起(会把积压直接跳过,等效丢消息);单纯重启不会丢已提交 offset——位点存在 Broker 上;扩容要保证消费幂等(重复消费仍安全)。 - 容量估算(步进):积压 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。
- 失败与降级:下游未恢复前禁止盲目扩容;非核心 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:30 | Kafka 特殊点 | 分区=最大并行度;临时 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 题(突发流量应急)——同样是“先止血再根治”
【自测】
- 下游 DB 挂了导致积压,扩容消费者有用吗? 参考答案: 没用且有害。瓶颈在下游,应先恢复下游,再限流爬坡恢复消费。
- 为什么 Kafka 加消费者实例有时速度不变? 参考答案: 消费者数超过分区数后,多余实例分不到分区,空闲。并行度上限=分区数。
- 判断对错:清积压时应全速消费尽快追平。 参考答案: 错。全速可能二次打挂刚恢复的下游,应限流爬坡。
84. 选型会:Kafka/RabbitMQ/RocketMQ 用哪个(MQ 选型)
【考察内容】MQ 选型是高频综合题
【题目】公司要做消息中间件选型,候选 Kafka、RabbitMQ、RocketMQ。按吞吐量、可靠性、顺序性、事务消息、延迟、运维成本几个维度对比,并给出你的业务场景(比如日志收集 vs 交易订单)该怎么选?
【参考答案】
- Kafka(顺序性:key→分区内原生局部有序,跨分区无序;组内分区独占消费,无原生死信/延迟):优点——超高吞吐(顺序写+零拷贝+批量)、天然分布式分区、消息可回溯(按 offset 重放)、生态好(大数据/日志);缺点——功能简单(无死信/延迟等高级特性)、消息可能重复(at-least-once)、管理较复杂。适合:日志、埋点、大数据、削峰场景;
- RabbitMQ(顺序性:三家最弱——只有单队列+单消费者才有序,多消费者或调大 prefetch 即乱序):优点——功能丰富(延迟队列、死信、优先级、灵活路由)、Erlang 稳定可靠、社区成熟;缺点——吞吐中等(万级)、分布式扩展(镜像队列)不如 Kafka。适合:业务解耦、复杂路由、中小流量;
- RocketMQ(顺序性:官方顺序消息 API=MessageQueueSelector 队列选择器+失败阻塞重试,最省事):优点——阿里出品、吞吐高(十万级)、事务消息/延迟消息/死信内置、Java 生态友好;缺点——社区相对小、功能多但运维复杂。适合:电商交易、金融业务消息(需要事务消息和可靠性的场景);
- 选型结论:日志/大数据→Kafka;交易/事务/延迟消息→RocketMQ;轻量业务/复杂路由→RabbitMQ;技术栈统一优先考虑团队熟悉度。
- 容量锚点(选型用):日志/埋点类常见日增 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。消费端能力要 ≥ 生产峰值,否则必须削峰策略。选型时同时估:集群规模、运维人力、机房/跨城复制带宽。
- 失败与降级:任何 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 题(监控)
【自测】
- 日志收集与订单状态通知,分别选什么?为什么? 参考答案: 日志选 Kafka(高吞吐、回放、大数据生态);订单通知选 RocketMQ(可靠、事务/延迟/死信内置,Java 交易栈)。
- 为什么不能说“Kafka 有事务所以等价于 RocketMQ 事务消息”? 参考答案: Kafka 事务保证跨分区读写的原子性与幂等生产,不解决“本地数据库事务与发消息”的两阶段问题;后者需要 RocketMQ 半消息或本地消息表。
- 判断对错:功能越多的 MQ 越好,应优先选功能全的。 参考答案: 错。功能多常伴随吞吐/复杂度代价。应按场景选:日志不需要事务消息,交易不需要百万 TPS 日志管道。
85. 自研消息队列技术预研,核心能力有哪些(MQ 设计)
【考察内容】“设计一个 MQ”是考察系统设计深度的经典题
【题目】团队做技术预研:想自研一个消息队列。从零设计,需要哪些核心能力?存储模型(日志/队列)、生产消费模型、ack 与重试、顺序、堆积、集群高可用分别怎么设计?
【参考答案】按“最小可用→逐步补齐”的顺序:
- 基础模型:生产者 → Broker(队列/分区)→ 消费者,消息的发送、存储、消费基本流程;
- 持久化:消息落磁盘(顺序写),解决重启丢消息——顺序写文件+批量刷盘;
- 高可用:副本机制(leader/follower),主故障自动选举,解决单点;
- 可靠性:生产端 ack、Broker 持久化、消费端 offset/ack 机制(处理成功才提交);
- 顺序性:分区/队列模型 + 同一 key 路由同分区 + 单消费者消费;
- 负载均衡:消费者组内分区分配(一个分区一个消费者)、Rebalance;
- 高级能力:延迟消息(时间轮/分级)、重试设计(独立重试队列/Topic+指数或分级退避+最大重试次数+不可重试异常直投死信;重试按延迟等级分档存放,避免定时器轮询全量消息;顺序 topic 的重试必须阻塞该 key 而不是并发重投,否则重新引入乱序)、死信队列(重试耗尽转人工,DLQ 要可查询、可回放)、事务消息(半消息+确认)、消息回溯(按 offset 重放)、削峰(消费限流);
- 其他:监控(积压量、消费 lag)、权限、压缩、多租户。
- 容量估算(自研 MQ 预研目标可写进设计文档):目标单机写入 5万–10万 msg/s(1KB 消息、3 副本、异步刷盘);同步刷盘场景降到 1万–3万 msg/s 量级。分区数规划:生产峰值 / 单分区消费能力;例峰值 5 万、单消费者 500/s → 至少 100 分区(再留 50% 余量)。存储:日增消息条数×体×副本×保留天;堆外内存/页缓存要能覆盖热数据。延迟消息:时间轮精度 vs 扫描频率的 trade-off(秒级精度可接受则定时扫描更简单)。
- 失败与降级:副本同步失败/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:30 | P0/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–3 | ISR 同步 |
| 默认语义 | 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 题(选型思维)
【自测】
- 为什么 MQ 存储多用追加日志? 参考答案: 顺序写高吞吐、实现简单、支持按 offset 回放;出队即删的队列模型难以高吞吐持久化与多订阅。
- 消费端 offset 提交的正确时机? 参考答案: 业务处理成功之后。过早提交会丢,过晚会重复(可接受)。
- 判断对错:自研 MQ 应优先实现延迟消息和事务消息。 参考答案: 错。应先完成持久化、ack、分区、副本的最小闭环;高级特性建立在可靠核心之上。
86. 单机几十万 QPS 写入,Kafka 凭什么(Kafka 高吞吐原理)
【考察内容】Kafka 高吞吐原理是大数据/后端双高频题
【题目】压测显示 Kafka 单机写入能到几十万 QPS,而传统 MQ 只有几万。Kafka 靠什么做到?请从磁盘顺序写、页缓存、零拷贝、批量、分区并行几个角度解释。
【参考答案】多维度设计:
- 顺序写磁盘:消息 append-only 追加到分区文件末尾,顺序写接近内存速度(600MB/s 级),避开随机 IO;
- 零拷贝:消费时用 sendfile(内核态直接发到网卡),避免用户态拷贝;生产时页缓存(page cache)直接写;
- 批量与压缩:生产者批量发送(batch)、Broker 批量写、消费批量拉取;支持压缩(lz4/zstd)减少网络与磁盘;
- 分区并行:topic 分成多个分区,分区是并行单元,多磁盘/多线程并行处理;
- 页缓存而非应用缓存:利用 OS page cache,读写都命中缓存;数据已落盘,进程重启后 page cache 通常仍可命中;但 page cache 在机器重启/断电后必然丢失,此时需回读磁盘;
- 稀疏索引:每个 segment 只索引少量偏移,查找快,索引内存占用小;
- 常数级开销:O(1) 磁盘读写(append 与读 offset 附近数据),不受数据总量影响。
- 容量估算: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 倍)。分区数:生产峰值 / 单分区吞吐,并保证 ≥ 消费者实例数。
- 失败与降级:开启高吞吐配置(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 为什么快)——同属“单机为什么能扛”方法论
【自测】
- Kafka 零拷贝主要优化的是哪条路径? 参考答案: 消费路径。Broker 用 sendfile 将 page cache 数据直接送网卡,避免用户态拷贝。
- 为什么批量化能提高吞吐? 参考答案: 网络 RTT、系统调用、刷盘等固定成本被多条消息分摊,单条成本下降。
- 判断对错:用了页缓存就不用落盘了。 参考答案: 错。page cache 机器重启/断电会丢;仍需顺序写落盘保证持久性。
87. 组内加消费者想加速,结果速度没变(Kafka 消费者组)
【考察内容】Kafka 消费模型
【题目】Kafka 消费速度不够,你往消费者组里加了几台机器,结果速度几乎没变。为什么?请讲清消费者组与分区的关系:为什么同一个分区不能被组内多个消费者同时消费?并行度由什么决定?
【参考答案】
- 概念:分区是 Kafka 的并行存储单元;消费者组(group)内的消费者共同消费一个 topic,分区分配给组内消费者(一个分区分配给一个消费者,一个消费者可消费多个分区);
- 为什么同组不能多消费者共用一个分区:
- 顺序性:分区内消息有序,若多消费者并发消费同分区,处理顺序无法保证;
- offset 管理:分区进度(offset)以“消费者组+分区”为单位提交,多消费者并发会导致 offset 提交互相覆盖(A 处理到 10,B 处理到 5,提交错乱);
- 推论:分区数=组内最大并行度——消费者数超过分区数时,多余的消费者空闲;分区数决定吞吐上限;
- 不同组可以消费同一分区(各组独立 offset,互不影响)——用于多业务各自消费同一份数据;
- Rebalance:消费者增减/分区变化触发再平衡,分区重新分配(可能引起重复消费,需幂等)。
- 容量/并行度估算:消费者组内有效并行度 = min(消费者实例数, 分区数)。若 topic 有 8 个分区,拉起 20 个消费者也只有 8 个在干活,空转 12 个。目标吞吐反推:峰值 2 万 msg/s、单消费者稳定 2000/s → 需要 ≥10 个分区 + ≥10 个消费者。加机器前先看:分区数、单消费者处理 RT、是否有共享下游瓶颈(DB 连接池打满则扩消费者无效)。Rebalance 频率过高也会吃掉吞吐——扩缩容要平滑。
- 失败与降级:消费者实例故障触发 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 + 处理 N | N 需自管顺序 |
| 有效并行度 | 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 重复)
【自测】
- 12 个分区,加到 20 个消费者,有几个在干活? 参考答案: 最多 12 个分到分区,其余空闲。
- 不同业务组能否消费同一分区? 参考答案: 可以。各组独立 offset,互不影响。
- 判断对错:消费慢就加消费者实例总能提速。 参考答案: 错。受分区数上限约束,且若瓶颈在下游或单条逻辑,加实例无效。
88. 消息重试了 10 次还是失败,一直卡住队列(死信队列)
【考察内容】死信队列是 MQ 可靠性治理高频题
【题目】有一条消息消费一直失败(比如数据格式非法),重试了 10 次还失败,占着队列头,后面的消息全被堵住。怎么处理?死信队列怎么设计?死信消息如何告警、人工介入、回放?
【参考答案】
- 概念:消息消费失败超过重试次数上限(如 3 次)后,转入死信队列(DLQ)——专门存放“处理不了”的消息,与正常队列隔离;
- 触发场景:业务异常(数据不合法)、依赖故障持续超时、消息格式错误;
- 处理机制:
- RabbitMQ:死信交换机(DLX)+ 死信路由键,消息在队列中被拒绝/过期/超限后转发;
- RocketMQ:重试队列(%RETRY%)消费失败自动重试,超过次数进死信队列(%DLQ%);
- Kafka:无内置死信,需自建——消费失败写“失败 topic”或落库;
- 死信处理:
- 告警通知:死信产生即告警(钉钉/短信),人工介入;
- 死信消费端:读取死信,分析失败原因(日志/消息体),修复后重新投递(重建消息回原队列);
- 定时重放:对可恢复的死信定时重试(如依赖恢复后);
- 兜底:无法恢复的登记后人工/对账处理;
- 设计要点:死信消息要保留原始信息(原 topic、消息体、失败原因、重试次数)。
- 容量估算:若业务消息峰值 5000 msg/s,失败率 1%,则重试流量约 50 msg/s;若单消费者重试逻辑更慢(如人工接口 1s/次),死信与重试队列会快速堆积。重试队列容量按“峰值失败量 × 最大重试次数 × 消息体 × 保留天”估算。例:50/s×10 次×1KB,保留 7 天 ≈ 约 300GB 量级(50×10×1KB ≈ 500KB/s × 604800s ≈ 302GB,再乘副本)。生产上重试 topic 要与主 topic 隔离分区与消费者,避免重试流量拖垮主链路。
- 失败与降级:重试耗尽 → 进死信,主队列必须立刻继续消费,不能被单条毒消息卡死;死信可人工修复后重新投递,或对账补偿;对可丢弃的非核心消息,业务确认后可定期清理死信。监控死信堆积量与增长速率,超过阈值告警。禁止:无限原地重试、或为了“清死信”无限制并发重放把下游打挂。
【原理溯源】
- 为什么要死信,而不是一直重试? 无限重试有三宗罪:① 毒消息(格式错误)永远失败,空耗资源;② 占住队列头/分区,阻塞后续消息(顺序消费时尤其致命);③ 重试风暴打挂下游。死信把“处理不了”从主通道剥离,主通道保持流动。
- 为什么毒消息必须隔离? 数据非法类失败与“下游暂时超时”不同:重试一万次也一样。若不识别异常类型并快速进死信,会把瞬时故障策略误用到永久故障上。
- Kafka 为什么没有内置 DLQ? Kafka 定位日志流,消费失败的语义交给应用。常见自建:失败发到
xxx.DLTtopic、或写入失败表;由专门消费者处理。这也说明死信是消费端模式,不是存储端魔法。 - 死信闭环为什么必须“告警 + 原因 + 可回放”? 只堆积不告警 = 静默数据丢失(业务以为在处理)。必须保留原 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 题(不丢)
【自测】
- 为什么要设死信而不是无限重试? 参考答案: 防止毒消息永久失败、阻塞通道、造成重试风暴;把不可处理消息隔离出来人工治理。
- Kafka 的死信一般怎么实现? 参考答案: 无内置,消费失败写入失败 topic 或落库,由专门消费者告警与回放。
- 判断对错:死信队列里的消息可以忽略。 参考答案: 错。死信常代表数据丢失或业务未完成,必须告警、分析、回放或对账补偿。
89. 订单 30 分钟未支付自动关单,怎么定时触发(延迟消息)
【考察内容】延迟任务与 MQ 结合
【题目】订单 30 分钟未支付自动关闭。用 MQ 的延迟消息能力怎么做?RocketMQ 延迟消息怎么实现(延迟级别)、RabbitMQ 的 TTL+死信方案怎么搭?各自的延迟精度和限制是什么?
【参考答案】
- RocketMQ 延迟消息:内置延迟等级(1s/5s/10s/30s/1m/.../2h 等 18 级),发送时设 delayTimeLevel,到点投递——简单可靠,但只支持固定档位;
- RabbitMQ 延迟队列:TTL(消息过期时间)+ 死信交换机——消息进“延迟队列”设 TTL,过期后被死信路由转发到真正消费队列;TTL 可设任意时长(需注意队列级/消息级 TTL 语义差异);
- Kafka:无内置延迟,需自建:方案 A——发消息时带“目标执行时间”,消费端判断未到时间则重新投递/暂存(可用分层时间轮+重投);方案 B——时间轮(Netty HashedWheelTimer)在应用层管理;
- 对比:RocketMQ 延迟消息最省事;RabbitMQ 用 TTL+DLX 组合;Kafka 生态需自研或用时间轮;
- 生产建议:MQ 延迟消息为主 + 定时任务扫表兜底(延迟消息丢失/超时的订单由定时任务扫描“待支付超时”订单关闭),双保险。
- 容量估算:日订单 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 万条)。关单消费要幂等,且与支付回调竞争同一状态机。
- 失败与降级:延迟消息丢失/投递失败 → 依赖兜底定时扫描(扫描 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 题(积压)
- 业务: 订单状态机、库存回滚
【自测】
- RocketMQ 延迟消息的主要限制? 参考答案: 只支持固定延迟级别,任意时长需对齐或配合扫表。
- 为什么关单要扫表兜底? 参考答案: 延迟消息可能丢失或消费失败;关单影响库存与资金正确性,需要数据库扫描作为最终一致保障。
- 判断对错:RabbitMQ 中给每条消息设不同 TTL 就能实现任意延迟。 参考答案: 不完全对。存在头阻塞与队列/消息 TTL 交互问题,需谨慎设计。
90. 消费者一抛异常就无限重试,队列被堵死(消费重试策略)
【考察内容】消息消费的容错设计
【题目】消费者处理消息抛异常后,如果无限重试会把队列堵死;不重试又会丢消息。重试策略怎么设计才合理(重试次数、退避间隔、区分可重试/不可重试异常、超过次数进死信)?
【参考答案】
- 失败分类处理:
- 可重试错误(下游超时、临时故障):自动重试,采用指数退避(1s、2s、4s...)或固定间隔,设最大次数(3~5 次);
- 不可重试错误(消息格式错误、业务数据非法):重试无意义,直接进死信/记录告警(避免无限重试浪费);
- 实现方式:
- RocketMQ:默认重试 16 次(延迟等级递增),可配最大重试次数;
- RabbitMQ:消息 reject/nack + 重试队列(或延迟重试);
- Kafka:消费端手动控制——捕获异常后不提交 offset,等待重投;或用“重试 topic”(失败消息发到重试 topic,延迟后重新消费);
- 关键防坑:
- 重试风暴:大量消息同时失败重试,打爆下游——指数退避+熔断(下游故障时暂停消费,恢复后再继续);
- 重试消息乱序:顺序消息重试要阻塞后续(或单线程消费);
- 重试与幂等配合:重投的消息消费时要幂等(去重);
- 监控:失败率、重试次数分布、死信量告警。
- 容量/重试风暴估算:假设峰值消费 2000 msg/s,其中 5% 下游异常 → 100 msg/s 进入重试;若策略是“立即无限重试”,等效把失败流量以乘数放大,可能在秒级打满下游与线程池。合理预算:最大重试 3~5 次、指数退避(1s/2s/4s)、重试 topic 与主 topic 隔离。线程池:重试消费与主消费隔离,重试线程数 ≤ 下游可承受并发的一定比例(如 20%)。观察指标:失败率、重试队列 lag、下游 RT/错误率。
- 失败与降级:识别不可重试异常(参数错误、业务拒绝)→ 直接死信/失败回调,不再重试;可重试异常(超时、限流、短暂不可用)→ 退避重试;达到上限 → 死信 + 告警 + 补偿任务。业务侧开关:极端时可暂停该 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 题(流控背压)
【自测】
- 消息格式错误为什么要直接进死信? 参考答案: 永久性失败,重试无意义且浪费资源、可能阻塞通道。
- 指数退避解决什么问题? 参考答案: 避免在下游故障期高频重试造成重试风暴,给恢复时间。
- 判断对错:重试次数越多越可靠。 参考答案: 错。过多重试放大负载、延迟暴露问题;应分类 + 上限 + 死信 + 告警。
91. 核心链路依赖的 MQ 突然不可用(MQ 故障降级)
【考察内容】依赖故障的容灾设计
【题目】MQ 集群故障不可用,而核心下单链路依赖它做异步化。业务不能停,怎么办?有没有降级方案(开关切同步、本地暂存、DB 轮询兜底)?恢复后积压数据怎么处理?
【参考答案】
- 认知:MQ 是核心依赖,要先做高可用(集群、副本、监控)降低挂的概率;但必须准备“MQ 挂掉”的降级预案;
- 降级方案:
- 本地消息表兜底:业务与消息同库同事务写入本地消息表,MQ 恢复后定时任务把积压消息补发(消息表是“永不丢失”的持久层)——最稳妥;
- 同步调用降级:MQ 不可用时,异步链路降级为同步调用(如通知/积分直接同步执行,牺牲 RT 保正确性)——适合低频轻量操作;
- 服务降级:非核心功能(推荐、统计)直接关闭,核心链路(下单、支付)优先保障;
- 消息落本地文件/磁盘队列:轻量场景可写本地队列文件,恢复后重放(有丢失风险);
- 恢复流程:MQ 恢复后,先补发积压(本地消息表/重放),再恢复消费,监控积压量回落;
- 预案要求:降级方案要提前演练(混沌工程),不能事发时临场设计。
- 容量估算(降级通道):MQ 全挂时,若下单峰值 5000 TPS、每单平均 2 条异步事件,则本地消息表要承接 1 万行/秒 写入——这通常不现实,必须配合:① 只保留核心事件(订单状态),非核心关闭;② 本地表批量写入;③ 业务限流到 DB/本地表可承受范围(如 500–1000 TPS)。口径必须连着说清,否则自相矛盾:题干峰值 5000 TPS 限到 1000 即只放行 20%,所以【关键数字】里「降级期间 ≥99%」只能指被放行的那部分流量的成功率,其余是排队或明确拒绝;不能同时承诺「全量下单 99% 成功」和「限流 1000 TPS」。降级窗口内用户侧表现:主流程可下单,积分/通知延迟;恢复后补发速度按下游容量限流,避免二次打挂。
- 失败与降级:预案分三级——① 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)
【自测】
- MQ 完全不可用时,本地消息表如何保证消息不丢? 参考答案: 消息与业务同一本地事务落库;MQ 只是投递通道,恢复后定时任务扫“待发送”补发。
- 为什么不适合把所有异步都改成同步降级? 参考答案: 同步会把下游 RT/故障引入主链路,高并发下可能导致下单超时雪崩;应只保留轻量低频,其余暂存或关闭。
- 判断对错:降级方案上线前不用演练,真挂了再按文档操作。 参考答案: 错。未演练的预案不可信;需混沌工程验证开关、补发、回切。
92. “先发消息还是先写库”才能两全(RocketMQ 事务消息)
【考察内容】事务消息是 RocketMQ 标志性能力,分布式事务专题高频题
【题目】下单要“写订单库 + 发消息给下游”,先写库可能消息没发出去,先发消息可能库没写成。RocketMQ 事务消息怎么解决这个两难?半消息、本地事务、回查机制分别是什么?
【参考答案】
- 解决的问题:先发消息再提交事务(事务失败消息白发了)/先提交事务再发消息(消息发送失败业务成功但下游不知道)——本地操作与发消息无法原子;
- 半消息机制:
- 生产者发送“半消息”(half message)——Broker 收到但不投递给消费者;
- 执行本地事务(如订单落库);
- 根据本地事务结果向 Broker 提交 commit(消息可见可投递)或 rollback(删除消息);
- 事务状态未知兜底(关键):若生产者本地事务执行中挂了/超时,Broker 长期持有半消息——Broker 会回查(check)生产者(回调 checkLocalTransaction),生产者查询本地事务状态后回复 commit/rollback;
- 保证:半消息+提交/回滚+回查,使“本地事务成功⇔消息最终可见”达到最终一致;
- 与本地消息表对比:事务消息把“消息表”搬到 MQ 内部,业务少维护一张表,但回查逻辑仍需业务实现。
- 容量估算:事务消息比普通消息多一次半消息写入 + 状态回查,吞吐通常低于普通发送(经验值约为普通 sync 的 50%–80%,视回查比例)。若交易峰值需发送 5000 msg/s 事务消息,按 80% 效率折算为 5000÷0.8=6250、按 50% 折算为 10000,故集群容量要按 6250–10000 普通消息当量预估。回查压力:本地事务未决比例越高,回查 QPS 越高;状态表要有索引且查询足够快(按 transactionId/msId)。半消息在 Commit 前对消费者不可见,会增加端到端延迟约一次确认往返(通常毫秒~十毫秒级)。
- 失败与降级:半消息发送成功但本地事务失败 → 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)
【自测】
- 半消息与普通消息的区别? 参考答案: Broker 暂存且不对消费者投递,直到 commit;rollback 则删除。
- 回查机制解决什么故障? 参考答案: 生产者在本地事务后未能发送 commit/rollback 时,Broker 主动查询本地事务状态以完成提交或回滚。
- 判断对错:用了事务消息,消费端就不会重复消费。 参考答案: 错。事务消息只管生产端一致性;消费仍可能重复,需幂等。
93. 不用事务消息,怎么保证“业务和发消息”不丢(本地消息表)
【考察内容】本地消息表是分布式事务最经典落地方案
【题目】如果 MQ 不支持事务消息,怎么保证“本地业务操作”和“发消息”要么都成功要么都补偿?本地消息表方案怎么设计?定时任务怎么扫描投递?重复投递怎么幂等?
【参考答案】
- 核心思想:业务操作与消息记录在同一个本地事务中完成,消息不丢;投递失败可重试;
- 流程:
- 生产者:开启本地事务 → 执行业务操作(如订单表插入)→ 同时写消息表(msg 状态=待发送)→ 提交事务(原子保证);
- 发送服务:扫描消息表中“待发送”消息(定时任务)→ 投递 MQ → 收到 ack 后更新状态=已发送;
- 消费者:消费消息处理业务(幂等),处理成功返回 ack;
- 可靠性保证:
- 不丢:消息与业务同事务落库,投递失败/宕机后定时任务继续扫发(最终一定发出);
- 不重:消费者幂等(唯一键)+ 发送端状态机(已发送不重发,或重发但消费幂等);
- 缺点:每套业务要多维护消息表、扫表有延迟、消息表与业务表耦合;与 RocketMQ 事务消息是同一思路的两种实现(一个在业务侧、一个在 MQ 侧);
- 适用:无事务消息功能的 MQ(Kafka/RabbitMQ)场景,或已有消息表基础设施。
- 容量估算:本地消息表扫描线程按轮询间隔 T 与批量大小 B 设计。若峰值需投递 2000 msg/s,轮询间隔 500ms,则每批需投递约 1000 条——批量发送比逐条快,但单批不宜过大(避免长时间事务与内存)。表增长:日订单 50 万×平均 2 消息 = 100 万行/日;保留 7 天 ≈ 700 万行,要按时间分区/归档。扫描索引:
(status, next_retry_time)必须有,否则全表扫描会拖垮 DB。 - 失败与降级: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 题(幂等)
【自测】
- 本地消息表的原子性靠什么保证? 参考答案: 业务数据与消息记录在同一数据库本地事务中提交。
- 为什么消费端仍要幂等? 参考答案: 扫表可能重复投递,MQ 也可能重复,必须保证业务效果只发生一次。
- 判断对错:本地消息表可以完全替代对账。 参考答案: 错。消费失败、业务部分成功等仍需对账兜底。
94. 一条消息 ack 没配对,消费者重启后消息全丢了(ack 机制)
【考察内容】MQ 投递语义体系
【题目】有同事把消费者的 auto-ack 打开后,消费者一处理完(甚至还没处理完)消息就被确认,重启后消息丢失。请讲清生产端 ack 与消费端 ack 的机制,ack 时机配错分别会带来什么问题?
【参考答案】
- 生产端 ack(Broker 确认收到):
- Kafka:acks 参数——acks=0 不等确认(可能丢);acks=1 leader 落盘即确认(leader 挂丢数据);acks=all 等所有 ISR 副本确认(最可靠,延迟最高);
- RabbitMQ:publisher-confirm,Broker 落盘后回调 ack,未确认/nack 可重发;
- 意义:生产端确认“消息已持久化”,避免“发了就当成功”导致丢失;
- 消费端 ack(消费者确认处理完):
- 自动 ack:收到消息即 ack——若处理中崩溃,消息已确认=丢失(默认禁用安全场景);
- 手动 ack:处理成功才 ack;处理失败不 ack(消息重新投递);
- Kafka:offset 提交机制(enable.auto.commit=false,处理成功后手动提交);
- ack 丢失场景:
- 生产 ack 丢失:消息实际已到 Broker 但生产者没收到 ack → 重发 → 重复消息(消费端幂等解决);
- 消费 ack 丢失:处理成功但 ack 没送达 → 消息重投 → 重复消费(幂等解决);
- 结论:ack 机制解决“丢”与“不丢”,但引入了“重复”——所以幂等是 MQ 可靠性闭环的最后一块拼图。
- 容量/确认开销估算:手动 ack 会让消费 RT 从“拉取即可”变成“处理完才确认”。若单条业务处理 10ms、单线程,则上限 100 msg/s/线程;要扛 5000 msg/s 需要 50 线程或 50 分区并行。同步刷盘/acks=all 会增加生产与 Broker 开销,吞吐下降但更安全。缓冲:in-flight 消息过多会导致 Rebalance 时重复面变大;建议限制
max.poll.records与处理批量,使单次 poll 在会话超时内处理完。 - 失败与降级:处理失败不 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 acks | 0/1/all | all 最安全 |
| enable.auto.commit | 生产建议 false | 手动提交 |
| autoAck | 核心业务 false | RabbitMQ |
| 重复率 | 网络抖动时可见 | 必须幂等 |
| ISR | 副本同步集合 | acks=all 等待对象 |
| 单线程消费上限 | 1/单条RT | 10ms→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 题(生产端一致)
【自测】
- auto-ack 在什么时点确认? 参考答案: 收到消息时(或按配置极早确认),处理中崩溃即丢失。
- Kafka acks=1 与 all 的差别? 参考答案: 1 只等 leader;all 等全部 ISR 副本,更安全更慢。
- 判断对错:ack 越早系统越高效,应尽量早确认。 参考答案: 错。过早确认会丢消息。应在业务成功后确认,用幂等处理重复。
95. 消息每秒 1000 条,下游接口只扛得住 100(消费端流控)
【考察内容】消费端容量治理
【题目】消费端每处理一条消息要调一个下游接口,该接口每秒最多承受 100 次调用;但消息流入量每秒 1000 条。怎么设计消费端(限速、分批聚合、削峰填谷、背压)才能既不丢消息又不压垮下游?
【参考答案】
- 需求本质:消费速率必须与下游处理能力匹配,防止把下游打挂、防止自己 OOM;
- 方案:
- 限流消费:消费端加限流器(Semaphore/Guava RateLimiter),按下游阈值控制处理速率;
- 批量消费:一次拉取/确认多条(Kafka 批量拉取、RabbitMQ prefetch 设置),批量处理(如攒 100 条调一次下游批量接口);
- 控制并发:消费者并发数(线程数/实例数)与下游容量匹配,不盲目扩;
- 背压/暂停:下游 RT 升高或错误率上升时,熔断暂停消费(消息留在 MQ),恢复后继续——比“继续消费硬重试”安全;
- 预取限制:RabbitMQ 设 prefetch count(如 1~10),避免一次性拉太多到本地内存;
- 监控:消费 lag、处理耗时、下游错误率联动告警;
- 特殊场景:大促削峰——允许积压(MQ 天然缓冲),下游恢复后追平,但要设积压告警上限。
- 容量估算(步进):上游 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 被重试吃掉。
- 失败与降级:限流阈值动态调整(下游 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 prefetch | 1–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 题(监控)
【自测】
- 为什么 1000/s 进、100/s 出时不应该全速消费? 参考答案: 会打挂下游并引发重试风暴;应限流到下游容量,积压留在 MQ。
- 批量消费的前提是什么? 参考答案: 下游支持批量或单次调用固定开销可被摊薄;且业务允许微批延迟。
- 判断对错:消费端 prefetch 越大吞吐越高越好。 参考答案: 错。过大导致内存压力与故障恢复变慢,应按处理能力设置。
96. MQ 上线后要盯哪些指标,才能提前发现风险(MQ 监控)
【考察内容】MQ 可观测性
【题目】MQ 接入核心链路后,运维要建立监控大盘。要盯哪些指标才能提前发现风险(积压、消费停滞、broker 负载、延迟)?哪些指标出现什么趋势就说明要出事了?
【参考答案】
- 积压量(Lag):消费 offset 落后生产 offset 多少——积压持续上涨=消费能力不足或下游故障,是最核心指标;按 topic/分区维度监控;
- 消费速率与生产速率:对比判断是否失衡;
- 消费失败率:失败率高=业务异常/消息异常,联动重试与死信监控;
- 死信队列量:死信持续产生=有系统性不可恢复问题,告警人工介入;
- 吞吐与延迟:生产/消费 TPS、端到端延迟(生产到消费的耗时);
- Broker 健康:CPU、磁盘(分区文件增长)、内存、网络、副本同步状态(ISR 收缩=副本异常,数据安全风险)、消费组 Rebalance 次数与时长;
- 重试次数分布:大量消息频繁重试=下游不稳定;
- 告警设计:积压量超过阈值(如 10 万)告警;失败率突增告警;死信非零告警;磁盘水位 75% 预警、85% 严重(与第 9 条同口径)。
- 容量相关监控阈值示例:生产 TPS 突增 > 基线 3 倍告警;消费 lag > 10 万或持续上涨 5 分钟告警;单 broker 磁盘水位 > 75% 预警、> 85% 严重;副本同步延迟 > 5s 关注;生产失败率 > 0.1% 告警;死信增长速率突变告警。容量巡检:按周评估峰值 TPS vs 集群水位,提前扩容分区/broker。大盘应同时展示:流入、流出、lag、失败、磁盘、CPU——只看 TPS 会漏掉“在堆积”和“在失败”。
- 失败与降级:监控本身挂了要有备用(多采集通道/云监控);告警要分级(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 题(降级)
【自测】
- 为什么死信量要单独告警? 参考答案: 死信代表最终处理失败,可能造成业务数据丢失;仅看 lag 会漏掉“消费成功但进死信”的假阴性。
- ISR 收缩意味着什么? 参考答案: 同步副本减少,数据可靠性下降,leader 故障时可能丢消息或不可用,需立即处理。
- 判断对错:只要消费 TPS 高就说明系统健康。 参考答案: 错。可能伴随高失败率或大量重试;需结合质量指标。