← 返回题库 / 大纲

消息队列架构

第 61 题:Kafka vs RocketMQ vs RabbitMQ 选型 选型

【真实企业业务场景】
公司有三类需求:①日志/埋点(百万 TPS、可丢);②订单交易(不能丢、要事务消息);③延迟任务(30 分钟后关单)。要求:三个场景分别选对 MQ 并说明取舍。
【面试官问题】
    1. 三者定位差异?2. 各适合什么场景?3. 吞吐量对比?4. 事务消息谁支持?5. 延迟消息?
【候选人的标准回答】

第一步:分析。选型看语义需求:高吞吐日志用 Kafka;金融级事务/顺序/延迟用 RocketMQ;复杂路由/低延迟用 RabbitMQ。

第二步:挑战。场景匹配、语义支持、运维。

第三步:架构。日志/埋点:Kafka(百万 TPS、分区并行、可丢);②订单交易:RocketMQ(事务消息、exactly-once 语义、顺序);③延迟关单:RocketMQ 延迟消息(18 级)或 RabbitMQ TTL+死信;④复杂路由:RabbitMQ(Exchange 绑定)。

第四步:选型。混合:Kafka(日志/流)+ RocketMQ(交易/延迟)+ RabbitMQ(内部复杂路由,可选)。

第五步:一致性。Kafka at-least-once + 消费幂等;RocketMQ 事务消息保证本地事务与发消息原子。

第六步:高可用。三者均多副本;Kafka ISR、RocketMQ 主从、RabbitMQ 镜像队列。

第七步:优化。按场景分集群,避免混用致相互影响。

【架构设计】
日志/埋点 → Kafka(高吞吐) 订单/交易 → RocketMQ(事务/顺序/延迟) 内部复杂路由 → RabbitMQ(可选) 按场景分集群, 不混用
【技术方案深度解析】

Kafka:分区+顺序写+零拷贝,吞吐极致,但事务/延迟支持弱,适合日志流。RocketMQ:阿里系,事务消息(半消息+回查)、定时/延迟消息、顺序消费成熟,适合交易。RabbitMQ:AMQP,Exchange 灵活路由,低延迟,但吞吐与堆积能力弱于前两者。一张对比表最直观:

维度KafkaRocketMQRabbitMQ
吞吐极高
事务消息强(半消息)无原生
延迟消息需外挂原生(18级)TTL+死信
顺序分区内队列内单队列
适用日志/流交易/金融内部路由
【关键技术点】
KafkaRocketMQRabbitMQ事务消息延迟消息场景选型
【Java 实现示例】
// RocketMQ 事务消息(订单)
Message m = new Message("order", JSON.toJSONString(order).getBytes());
SendResult r = txProducer.sendMessageInTransaction(m, order);
// Kafka 日志(可丢, acks=1, 批量)
props.put("acks", "1"); props.put("linger.ms", "5");
【面试官可能继续追问】
【常见错误回答】
  • ❌ 只用一种 MQ 包打天下——语义不匹配。正确:按场景选。
  • ❌ 交易用 Kafka 无事务——丢消息。正确:RocketMQ 事务。
  • ❌ 日志用 RabbitMQ——吞吐不够。正确:Kafka。
【架构师评分标准】
初级 0~40
只知一种。
中级 40~60
知三者,无场景论证。
高级 60~80
场景驱动+语义对比+分集群。
架构师 80~100
再加:①语义保证深度;②运维/成本权衡;③演进策略;④统一封装治理。

第 62 题:消息不丢失全链路保障 可靠性

【真实企业业务场景】
订单支付成功后发 MQ 通知仓储发货,曾因「Broker 宕机 + 生产者未确认」丢失 2000 单,仓库未发货客诉。要求:保障消息从生产到消费「零丢失」。
【面试官问题】
    1. 哪些环节会丢?2. 生产者怎么保?3. Broker 怎么保?4. 消费者怎么保?5. 极端宕机?
【候选人的标准回答】

第一步:分析。丢失在三处:生产(未确认)、Broker(未刷盘/副本)、消费(未处理就 ack)。需逐环节加固

第二步:挑战。性能 vs 可靠权衡、确认机制、刷盘。

第三步:架构。生产:同步发送 + 重试 + 回调确认(RocketMQ 同步 send;Kafka acks=all+重试);②Broker:同步刷盘 + 多副本(Raft/ISR);③消费先处理业务再 ack(非自动提交),失败不 ack 触发重试;④兜底:本地消息表/事务消息防生产丢。

第四步:选型。RocketMQ(同步双写+同步刷盘)/ Kafka(acks=all+min.insync.replicas)。

第五步:一致性。at-least-once + 消费幂等 = 不丢不重。

第六步:高可用。多副本防 Broker 丢;消费重试防处理丢。

第七步:优化。关键链路同步刷盘,非关键异步提吞吐。

【架构设计】
生产: 同步send+重试+确认 → Broker(同步刷盘+多副本) → 消费(处理完再ack) 兜底: 本地消息表/事务消息; 不丢=at-least-once+幂等
【技术方案深度解析】

为什么先处理后 ack?若先 ack 再处理,处理崩=丢消息;先处理后 ack,处理失败重投(at-least-once)。acks=all/min.insync.replicas:Leader 等多数副本写入才确认,防单点丢。代价:同步刷盘/多副本降吞吐,关键链路值得。

【关键技术点】
同步发送同步刷盘acks=allISR处理完再ack本地消息表
【Java 实现示例】
// Kafka: acks=all + 重试, 防生产/Broker丢
props.put("acks", "all");
props.put("min.insync.replicas", "2");
props.put("retries", "5");
// 消费: 手动提交, 处理完再 ack(关闭 auto-commit)
consumer.subscribe("order");
for (Record r : consumer.poll(1000)) {
  handle(r); consumer.commitSync(); // 处理成功才提交
}
【面试官可能继续追问】
【常见错误回答】
  • ❌ 发完就当成功——未确认丢。正确:同步+确认。
  • ❌ 先 ack 再处理——处理崩丢。正确:处理后 ack。
  • ❌ 异步刷盘单副本——Broker 丢。正确:同步+多副本。
【架构师评分标准】
初级 0~40
不知丢失环节。
中级 40~60
知生产确认,不论 Broker/消费。
高级 60~80
三环节+刷盘+多副本+处理后ack。
架构师 80~100
再加:①性能/可靠权衡;②本地消息表兜底;③极端恢复;④监控(丢失率≈0)。

第 63 题:消息重复与幂等消费 幂等

【真实企业业务场景】
积分服务消费「支付成功」消息,因 Broker 重投 + 消费 ack 超时,同一笔支付被加积分 3 次,用户白拿积分。要求:从架构上保证消费幂等。
【面试官问题】
    1. 为什么会有重复?2. 幂等怎么实现?3. 数据库幂等?4. Redis 幂等?5. 幂等键怎么选?
【候选人的标准回答】

第一步:分析。MQ 是 at-least-once(不丢),重复必然存在。应对=消费幂等(同消息处理多次=一次效果)。

第二步:挑战。重复来源、幂等键、并发、存储。

第三步:架构。幂等键:消息带业务唯一 ID(如支付流水号);②去重:a) 数据库唯一索引/冲突更新(INSERT ... ON DUPLICATE);b) Redis SETNX 幂等表(带 TTL);c) 状态机(已处理则跳过);③并发:唯一键保证原子;④兜底:对账修复。

第四步:选型。唯一索引(最稳)+ Redis 幂等表(高性能)+ 状态机。

第五步:一致性。幂等保证重复处理不重复副作用;与「不丢」互补。

第六步:高可用。幂等存储高可用;失败重试不破坏幂等。

第七步:优化。幂等表设 TTL 防无限增长;热点用 Redis。

【架构设计】
消息(业务唯一ID) → 消费 → 幂等判断 DB唯一索引 / Redis SETNX(幂等表,TTL) / 状态机 重复: 已处理→跳过, 保证一次效果
【技术方案深度解析】

为什么重复不可避免?at-least-once 语义下,网络/ack 超时必重投。幂等键:必须用稳定业务 ID(支付流水号),不能用 MQ 自带的 offset(会变)。唯一索引最稳:INSERT 冲突即重复,DB 层原子保证;Redis SETNX 高性能但有 TTL 窗口需对账兜底。

【关键技术点】
幂等业务唯一ID唯一索引SETNX状态机at-least-once
【Java 实现示例】
// DB 唯一索引幂等(支付流水号唯一)
try { pointMapper.add(txId, uid, score); }      // 唯一键冲突=重复
catch (DuplicateKeyException e) { return; } // 跳过
// Redis 幂等表(高性能)
if (!redis.setnx("dedup:"+txId, "1", 24h)) return; // 已处理跳过
handle(msg);
// 状态机: if(order.status==PAID) return;
【面试官可能继续追问】
【常见错误回答】
  • ❌ 不处理重复——重复副作用。正确:幂等。
  • ❌ 用 offset 做幂等键——会变。正确:业务 ID。
  • ❌ 仅 Redis 无兜底——丢后重复。正确:DB 唯一索引。
【架构师评分标准】
初级 0~40
不知重复必然。
中级 40~60
知幂等,键选错。
高级 60~80
业务ID+唯一索引/SETNX+状态机+兜底。
架构师 80~100
再加:①并发原子保证;②TTL 与对账;③存储可靠性分层;④监控(重复率)。

第 64 题:消息顺序性保障 顺序

【真实企业业务场景】
订单状态流「创建→支付→发货」,消费者并发处理致「发货」先于「支付」被处理,状态机异常。要求:保证同一订单消息严格有序。
【面试官问题】
    1. 为什么乱序?2. 怎么保证分区内有序?3. 全局有序?4. 并发消费与顺序矛盾?5. 乱序怎么修复?
【候选人的标准回答】

第一步:分析。MQ 只保证分区/队列内有序。全局有序代价极高(单分区/单消费者)。按业务键路由到同一分区即可。

第二步:挑战。路由、并发、全局序。

第三步:架构。路由:同一订单号 hash 到同分区(Kafka partitionKey=orderId;RocketMQ 同 queue);②单分区单消费线程:保证顺序处理;③全局序:仅单分区单消费者(牺牲并行,仅极端需要);④修复:状态机拒绝非法跃迁 + 重排序缓冲。

第四步:选型。Kafka(partitionKey)/ RocketMQ(MessageQueue 顺序消费)。

第五步:一致性。顺序消息保证状态流转正确;失败重试保持顺序(不跳)。

第六步:高可用。顺序消费失败阻塞该分区(或特殊处理),不影响其他分区。

第七步:优化。仅对需顺序的 topic 用顺序,其余并行提吞吐。

【架构设计】
orderId → hash → 同分区/队列 → 单线程顺序消费 状态机: 创建→支付→发货(非法跃迁拒绝) 并发: 仅同key有序, 不同key并行
【技术方案深度解析】

分区内有序:Kafka 同 partition 内 offset 递增,单消费者线程顺序读即有序。路由关键:orderId 作 key 保证同订单同分区。全局序代价:单分区单消费者=无并行,吞吐极低,几乎不用。状态机防乱序:即使极少乱序,状态机拒绝「发货前无支付」,可缓冲重排。

【关键技术点】
分区有序partitionKey顺序消费状态机单线程路由
【Java 实现示例】
// Kafka: 同 orderId 同分区
producer.send(new ProducerRecord<>("order", orderId, msg));
// RocketMQ: 同 key 同 queue + 顺序消费监听器
SendResult r = producer.send(msg, (q,m)-> q.get( Math.abs(orderId.hashCode())%q.size()), msg);
// 消费: 状态机校验
if (!order.canTransit(toStatus)) { buffer.reorder(msg); return; }
【面试官可能继续追问】
【常见错误回答】
  • ❌ 多分区并发要全局序——乱。正确:key 路由同分区。
  • ❌ 顺序消费不用单线程——乱。正确:单分区单线程。
  • ❌ 无状态机——非法跃迁。正确:状态机兜底。
【架构师评分标准】
初级 0~40
不知分区有序。
中级 40~60
知路由,不懂单线程/状态机。
高级 60~80
key路由+单线程+状态机+失败处理。
架构师 80~100
再加:①全局序代价论证;②扩容兼容;③重排缓冲;④吞吐与顺序权衡。

第 65 题:消息积压治理 积压

【真实企业业务场景】
大促消费者 bug 致处理慢,Kafka 积压 5000 万条,消费滞后 2 小时,下游库存更新全延迟。要求:快速消化积压且不引发二次故障。
【面试官问题】
    1. 积压怎么发现?2. 怎么快速消费?3. 消费者扩容?4. 紧急丢弃/转存?5. 预防?
【候选人的标准回答】

第一步:分析。积压=生产 > 消费。治理=「提消费能力 + 临时扩容 + 降级非核心 + 修复根因」。

第二步:挑战。消费速度、扩容、下游扛不住、根因。

第三步:架构。监控:lag 告警(Kafka consumer lag);②提能力:临时增加分区+消费者实例(Kafka 分区数=最大并行度);③降级:非核心消息跳过/转冷存后补;④修复:定位慢查询/锁;⑤紧急:抽样转存(如转对象存储)避免 broker 爆。

第四步:选型。Kafka(增分区+消费组扩容)+ 监控(Burrow/Prometheus)+ 降级开关。

第五步:一致性。积压消息仍按顺序(同分区);降级需保证可补。

第六步:高可用。消费扩容不冲击下游(限流);转存防 broker 满。

第七步:优化。批量消费+并行;修复根因防再积。

【架构设计】
lag监控 → 增分区+扩消费者 → 批量/并行消费 非核心: 降级转冷存(后补); 根因: 慢SQL/锁修复 紧急: 抽样转存防broker爆
【技术方案深度解析】

为什么增分区?Kafka 并行度受分区数限制,消费者数 ≤ 分区数才有效,故先增分区再扩实例。下游扛不住:消费快了但下游(DB)被打,需限流或批量。根因优先:不修 bug 只扩容量,积压会再来;先止血(扩容)再修根因。

【关键技术点】
consumer lag增分区消费组扩容批量消费降级转存根因修复
【Java 实现示例】
// 批量消费提吞吐
@KafkaListener(topics="order", concurrency=32) // 并发=分区数
public void on(List<ConsumerRecord> batch){
  db.batchInsert(batch); // 批量写
  consumer.commitSync();
}
// 监控: kafka-consumer-groups.sh --describe --lag
【面试官可能继续追问】
【常见错误回答】
  • ❌ 只扩消费者不增分区——无效。正确:先增分区。
  • ❌ 不修根因——再积。正确:止血+修因。
  • ❌ 消费猛冲打垮下游——二次故障。正确:限流。
【架构师评分标准】
初级 0~40
不知 lag。
中级 40~60
知扩容,不增分区。
高级 60~80
监控+增分区+批量+降级+根因。
架构师 80~100
再加:①下游保护(限流);②紧急转存;③容量与预防;④可观测与告警。

第 66 题:死信队列与重试机制 死信

【真实企业业务场景】
通知服务消费消息,第三方短信网关偶发 5xx,消息一直重试阻塞队列,正常通知也被拖慢。要求:设计重试+死信,失败消息隔离不阻塞。
【面试官问题】
    1. 重试策略?2. 死信队列是什么?3. 重试阻塞怎么办?4. 死信怎么处理?5. 指数退避?
【候选人的标准回答】

第一步:分析。失败消息需有限重试 + 隔离,避免阻塞正常消费。死信队列(DLQ)承接最终失败消息。

第二步:挑战。重试风暴、阻塞、死信处理。

第三步:架构。重试:指数退避 + 上限(如 3~5 次),非业务错(4xx)不重试;②隔离:重试用独立重试队列/延迟重试,不占主消费;③死信:超上限进 DLQ,人工/定时补偿,不阻塞主队列;④告警:DLQ 有消息即告警。

第四步:选型。RocketMQ(重试队列+死信自动)/ RabbitMQ(DLX+ TTL)/ Kafka(自建重试 topic)。

第五步:一致性。死信保留消息可补;重试保证最终处理。

第六步:高可用。死信隔离保主队列流畅;补偿可重投。

第七步:优化。分类重试(瞬时 vs 永久错);死信人工台。

【架构设计】
主队列 → 消费失败 → 重试队列(指数退避,上限) 超上限 → 死信队列(DLQ) → 告警+人工/补偿重投 正常消息不受失败消息阻塞
【技术方案深度解析】

为什么死信?持续失败的消息若一直在主队列重试,占资源+阻塞同队列其他消息。进 DLQ 隔离,主队列继续。指数退避:1s→2s→4s 避免立即重试打爆下游。业务错不重试:参数错(4xx)重试无用,直接死信。

【关键技术点】
死信队列重试指数退避DLX隔离补偿
【Java 实现示例】
// RocketMQ: 重试(默认16次退避) + 死信(%DLQ%topic)
@RocketMQMessageListener(topic="notify", consumerGroup="g")
public void on(Msg m){
  try { sms.send(m); }
  catch (RetryableException e) { throw e; } // 触发重试
  catch (BizException e) { // 不重试
    deadLetter.save(m); return;
  }
}
// 死信消费组: 告警+人工补偿
【面试官可能继续追问】
【常见错误回答】
  • ❌ 失败无限重试——阻塞。正确:上限+死信。
  • ❌ 业务错也重试——无用。正确:直死信。
  • ❌ 死信不处理——丢消息。正确:告警+补偿。
【架构师评分标准】
初级 0~40
无重试/死信。
中级 40~60
知重试,无隔离。
高级 60~80
退避+死信+隔离+补偿+告警。
架构师 80~100
再加:①错误分类(瞬时/永久);②重试幂等;③自建 DLQ(Kafka);④可观测(死信率)。

第 67 题:延迟消息与定时任务 延迟

【真实企业业务场景】
订单 30 分钟未支付自动关单,每天 200 万未支付单。用定时任务扫全表太慢,用 RocketMQ 延迟消息精度不够(仅 18 级)。要求:高精度任意延迟方案。
【面试官问题】
    1. 延迟消息原理?2. 18 级不够怎么办?3. 任意精度方案?4. 关单与支付并发?5. 大量延迟消息性能?
【候选人的标准回答】

第一步:分析。延迟消息=「到时间才消费」。精度需求决定方案:固定档用 MQ 延迟,任意精度用时间轮/DB 扫描/RocketMQ 定时。

第二步:挑战。精度、并发、性能、可靠。

第三步:架构。固定档:RocketMQ 延迟消息(如 30 分钟档);②任意精度:自研时间轮(Netty HashedWheelTimer)+ DB 持久化,或 RocketMQ 5.0 定时消息(任意);③关单并发:消费时查订单状态(已付则跳过),幂等;④性能:延迟消息走独立 topic,批量。

第四步:选型。RocketMQ 定时消息(5.0+任意)/ 时间轮 + DB / 或 DB 扫描(低精度)。

第五步:一致性。关单前校验状态,已付不关;失败可重投。

第六步:高可用。延迟消息持久化;时间轮故障可重建(DB 兜底)。

第七步:优化。时间轮 O(1) 触发;DB 扫描低峰批量。

【架构设计】
下单 → 发延迟消息(30min) → 到点消费 → 查状态(未付→关单, 已付→跳过) 任意精度: 时间轮+DB持久化 / RocketMQ5定时消息
【技术方案深度解析】

MQ 延迟原理:消息暂存内部 Schedule 主题,到点才转正常队列。18 级局限:仅固定档(1s/5s/.../2h),「37 分钟」无法表达。RocketMQ 5.0 支持任意定时。时间轮:Hash 分桶,O(1) 入/出,适合海量定时;需 DB 持久化防丢。状态校验:延迟消息到时订单可能已付,必须查状态再关。

【关键技术点】
延迟消息时间轮定时消息状态校验任意精度持久化
【Java 实现示例】
// RocketMQ 5.0 任意定时消息
Message m = new Message("order-delay", body);
m.setDelayTimeSec(30*60); // 任意秒
producer.send(m);
// 消费: 状态校验(幂等关单)
public void on(OrderMsg m){
  Order o = orderMapper.get(m.orderId);
  if (o.status == UNPAID) orderMapper.close(m.orderId); // 已付跳过
}
【面试官可能继续追问】
【常见错误回答】
  • ❌ 全表定时扫——慢且锁。正确:延迟消息/时间轮。
  • ❌ 关单不校验状态——误关已付。正确:状态校验。
  • ❌ 18 级当任意精度——不准。正确:时间轮/5.0定时。
【架构师评分标准】
初级 0~40
全表扫描。
中级 40~60
知延迟消息,精度不管。
高级 60~80
延迟/时间轮+状态校验+持久+并发。
架构师 80~100
再加:①任意精度方案(时间轮/5.0);②可靠性(持久化);③状态机防并发;④性能与分片。

第 68 题:Kafka 高性能调优 调优

【真实企业业务场景】
Kafka 集群单机吞吐仅 5 万 msg/s(预期 20 万),生产者吞吐上不去,消费 lag 涨。要求:定位瓶颈并调优到百万级。
【面试官问题】
    1. 吞吐瓶颈在哪?2. 生产者怎么调?3. 服务端怎么调?4. 消费怎么调?5. 零拷贝?
【候选人的标准回答】

第一步:分析。Kafka 高性能=顺序写 + 零拷贝 + 批量 + 压缩。瓶颈常在未批量/未压缩/单分区/磁盘

第二步:挑战。批量、压缩、分区、磁盘 IO。

第三步:架构。生产:批量(batch.size=64k, linger.ms=5)、压缩(snappy/lz4)、acks=1(非关键);②服务端:多分区、顺序写(SSD)、页缓存;③消费:批量拉(max.poll.records)、并行;④零拷贝:sendfile 消费传输。

第四步:选型。Kafka 参数调优 + 分区规划 + SSD。

第五步:一致性。acks=1 降可靠换吞吐;关键用 all。

第六步:高可用。多副本;调优不破坏副本。

第七步:优化。压测定位瓶颈(生产/网络/磁盘/消费)。

【架构设计】
生产: 批量+压缩+linger → Broker(顺序写+页缓存+多分区) → 消费(批量拉+并行) 零拷贝: sendfile 消费传输
【技术方案深度解析】

批量+linger:攒批发降低网络往返与开销,linger.ms 等一小会凑批。压缩:snappy/lz4 降网络与磁盘。顺序写:Kafka 追加写日志,磁带式顺序 IO 极快;页缓存由 OS 管理。零拷贝:消费时 sendfile 直接内核→socket,免用户态拷贝。

【关键技术点】
批量压缩多分区顺序写零拷贝页缓存
【Java 实现示例】
// 生产者调优
props.put("batch.size", "65536");
props.put("linger.ms", "5");
props.put("compression.type", "lz4");
props.put("acks", "1");
// 消费: 批量拉
props.put("max.poll.records", "500");
// 服务端: num.partitions 多 + log.segment 调优 + SSD
【面试官可能继续追问】
【常见错误回答】
  • ❌ 单条发——吞吐低。正确:批量+linger。
  • ❌ 不压缩——网络满。正确:lz4。
  • ❌ 分区太少——并行受限。正确:多分区。
【架构师评分标准】
初级 0~40
不知批量。
中级 40~60
知参数,不论证原理。
高级 60~80
批量+压缩+分区+零拷贝+压测。
架构师 80~100
再加:①顺序写/页缓存原理;②可靠与吞吐权衡;③分区规划;④瓶颈定位方法论。

第 69 题:事务消息与本地事务一致 事务消息

【真实企业业务场景】
下单成功后需发「订单创建」事件给 5 个下游(库存/积分/推荐...)。曾因「库里订单写成功但消息没发出」,下游全不一致。要求:保证本地事务与发消息原子。
【面试官问题】
    1. 为什么普通发消息不原子?2. 事务消息原理?3. 半消息与回查?4. 回查失败?5. 替代方案?
【候选人的标准回答】

第一步:分析。本地 DB 事务与发 MQ 是两个操作,无法原子。事务消息(半消息+回查)保证「本地成则消息发,本地败则消息不发」。

第二步:挑战。原子性、回查、可靠性。

第三步:架构。半消息:先发半消息(对消费者不可见);②本地事务:执行 DB 写;③提交/回滚:成则 commit(半消息变可见),败则 rollback;④回查:Broker 未收到 commit/rollback 时回调生产者查本地状态;⑤替代:本地消息表(见第二部分 Outbox)。

第四步:选型。RocketMQ 事务消息(半消息+回查)/ 本地消息表(Kafka)。

第五步:一致性。事务消息保证最终一致;回查兜底防悬挂。

第六步:高可用。回查保证中间状态可决断;消费者幂等。

第七步:优化。回查限频;幂等消费。

【架构设计】
发半消息 → 执行本地事务 → commit/rollback Broker未决 → 回查生产者(查DB状态) → 决断 下游: 幂等消费(防重复)
【技术方案深度解析】

为什么半消息?若先发消息再写库,写库失败但消息已出→不一致。半消息先占坑,本地事务成再「转正」。回查关键:网络丢 commit/rollback 时,Broker 不知状态,回调生产者查 DB(订单在=提交,不在=回滚),防消息悬挂。替代:本地消息表(同库写消息,CDC 发),Kafka 无原生事务消息时用。

【关键技术点】
事务消息半消息回查本地消息表最终一致幂等
【Java 实现示例】
// RocketMQ 事务消息
TransactionSendResult r = txProducer.sendMessageInTransaction(
    new Message("order", body), orderId);
// 本地事务
public LocalTransactionState executeLocal(Message m, Object o){
  return orderMapper.insert((Order)o) ? COMMIT : ROLLBACK;
}
// 回查: 查DB订单是否存在
public LocalTransactionState checkLocal(Message m, Object o){
  return orderMapper.exists(orderId) ? COMMIT : ROLLBACK;
}
【面试官可能继续追问】
【常见错误回答】
  • ❌ 先发消息再写库——不一致。正确:半消息+回查。
  • ❌ 不回查——悬挂。正确:回查决断。
  • ❌ 消费者不幂等——重复。正确:幂等。
【架构师评分标准】
初级 0~40
不知原子问题。
中级 40~60
知半消息,不论回查。
高级 60~80
半消息+回查+本地消息表+幂等。
架构师 80~100
再加:①回查可靠性与限频;②与 Saga/Outbox 对比;③异常决断;④多下游解耦。

第 70 题:消息轨迹与可观测性 可观测

【真实企业业务场景】
一笔支付消息「丢了还是消费慢」排查 2 小时,无链路可见性。要求:建立 MQ 全链路可观测,能定位「生产→Broker→消费」每一环。
【面试官问题】
    1. 观测看什么?2. 消息轨迹怎么实现?3. TraceID 透传?4. 积压/延迟监控?5. 告警?
【候选人的标准回答】

第一步:分析。MQ 可观测=「生产成功?Broker 收到?消费耗时?积压多少?」全链路指标+轨迹。

第二步:挑战。轨迹存储、Trace 透传、指标。

第三步:架构。指标:生产 TPS、ack 率、Broker lag、消费 TPS/耗时、重试/死信数(Prometheus);②轨迹:消息带 msgId+TraceID,生产/消费各记一条(RocketMQ 消息轨迹/自建);③Trace 透传:TraceID 随消息头传下游,串联链路;④告警:lag>阈值、消费失败率、死信。

第四步:选型。RocketMQ 轨迹 / Kafka+Prometheus / OpenTelemetry + Jaeger。

第五步:一致性。轨迹与业务解耦,异步记录不阻塞。

第六步:高可用。轨迹存储独立;不影响主链路。

第七步:优化。采样轨迹降成本;关键消息全量。

【架构设计】
生产(记msgId+TraceID) → Broker(落盘指标) → 消费(记耗时+TraceID) 指标: TPS/lag/失败率/死信 → 大盘+告警 轨迹: msgId 串联生产→消费全链路
【技术方案深度解析】

为什么轨迹?消息异步,出问题难定位「丢在哪一环」。轨迹记录每环时间+状态,秒级定位。TraceID 透传:把请求 TraceID 放进消息头,消费时取出继续透传,MQ 成为链路一环(OTel 支持)。lag 是关键指标:直接反映消费健康。

【关键技术点】
消息轨迹TraceID透传consumer lagOpenTelemetry死信监控告警
【Java 实现示例】
// 发消息带 TraceID
Message m = new Message(topic, body);
m.putUserProperty("traceId", Tracing.current());
// 消费取出继续透传
public void on(Msg m){
  Tracing.start(m.getUserProperty("traceId")); // 串联链路
  handle(m);
}
// 指标: Prometheus 暴露 lag/耗时, Grafana 告警
【面试官可能继续追问】
【常见错误回答】
  • ❌ 无轨迹——排查盲。正确:全链路轨迹。
  • ❌ Trace 不进 MQ——链路断。正确:头透传。
  • ❌ 不监控 lag——积压无感。正确:lag 告警。
【架构师评分标准】
初级 0~40
无观测。
中级 40~60
知 lag,无轨迹。
高级 60~80
指标+轨迹+Trace透传+告警。
架构师 80~100
再加:①采样降本;②跨 MQ Trace;③死信可观测;④排查 SOP。
第六部分 · 消息队列架构(第 61~70 题) · 返回大纲