消息队列通用面试题
消息队列通用面试题
导语:本篇只讲与具体 MQ 无关的通用考点——投递语义、可靠投递、幂等去重、顺序性、积压治理、死信与延迟消息、事务消息,以及消费模型与工程实践。各 MQ 的差异化实现见《RabbitMQ》《RocketMQ》《Kafka》三篇。共 25 题。
一、MQ 的价值与代价
1. 为什么要引入消息队列?典型应用场景有哪些?
答: 引入 MQ 的核心价值是解耦、异步、削峰,对应三类典型场景:
- 解耦:多个下游系统都依赖同一份数据(如下单后通知库存、积分、物流),若直接同步调用,任何下游变更都会影响上游。引入 MQ 后,上游只管发消息,下游各自订阅,互不耦合。
- 异步:非核心链路(发短信、写操作日志、埋点)不需要同步返回,丢进 MQ 后台慢慢处理,缩短主链路 RT。
- 削峰填谷:秒杀、抢购等瞬时流量可达平时的数十倍,MQ 把请求先缓存起来,下游按自身处理能力匀速消费,保护数据库不被冲垮。
注意:引入 MQ 会显著提升系统复杂度(要考虑消息丢失、重复、顺序、积压),不是"免费的午餐",需在收益与代价间权衡。
2. 消息队列会带来哪些副作用?
答: 主要成本:
- 系统可用性降低:MQ 挂了,依赖它的链路全部受影响,需保证 MQ 自身高可用(集群 + 多副本)。
- 一致性复杂度上升:引入了「分布式数据不一致」风险(如本地事务已提交但消息未发出)。
- 复杂度提升:要处理消息丢失、重复消费、顺序错乱、消费积压、消息堆积导致的磁盘/内存压力等问题。
- 运维成本:多了需要监控、扩容、调优的组件(Lag 监控、分区规划、副本管理)。
一句话:MQ 把"同步调用的强耦合"换成了"异步的最终一致"——换来吞吐与解耦,付出的是"必须自己保证最终一致"的代价。
3. 消息队列的三种投递语义是什么?
答: 这是 MQ 最核心的基础概念,面试高频:
| 语义 | 含义 | 实现方式 | 性能 |
|---|---|---|---|
| 至多一次(at most once) | 可能丢,绝不重 | 先提交位移/ACK,再处理业务 | 最高 |
| 至少一次(at least once) | 不丢,但可能重 | 先处理业务成功,再提交位移/ACK | 中等 |
| 恰好一次(exactly once) | 不丢也不重 | 幂等生产者 + 事务(MQ 内部);端到端需消费端幂等 | 最低 |
关键认知(面试必答):
- 绝大多数 MQ 的默认语义是"至少一次"——「先业务、后 ACK」能保证不丢,但 ACK 失败/超时会导致重投,所以必然可能重复;
- "恰好一次"要分两层看:
- MQ 内部:如 Kafka 通过幂等生产者(PID + 序列号去重)与事务实现;
- 端到端:因为"业务处理"(写业务库)与"提交位移"(写 MQ)无法放在同一个原子操作里,所以端到端的恰好一次几乎不可能只靠 MQ 实现。
- 工程上的标准答案:"至少一次投递 + 消费端幂等" = 业务等效于恰好一次。这也是为什么"消费幂等"是必考题。
二、消息可靠投递
4. 一条消息从生产到消费,可能在哪些环节丢失?
答: 三个阶段、六种典型场景:
| 阶段 | 丢失场景 |
|---|---|
| ① 生产者 → Broker | ① 网络抖动,消息没到 Broker 却没被感知;② Broker 收到但未持久化就返回成功,随后宕机 |
| ② Broker 存储 | ③ 消息只在内存/Page Cache、未刷盘时宕机;④ 主从异步复制时主库宕机,未同步的副本数据丢失 |
| ③ Broker → 消费者 | ⑤ 拉到消息后先 ACK 再处理,处理前进程崩溃;⑥ 消费端自动提交位移,业务还没跑完就提交了 |
对应的三类防线:
- 生产端:同步发送 + 重试;RabbitMQ 用 publisher confirm、Kafka 用 acks=all、RocketMQ 用 同步刷盘 + 同步复制;
- Broker 端:消息持久化 + 多副本(Kafka 的
min.insync.replicas、RocketMQ 的SYNC_MASTER); - 消费端:关闭自动 ACK/自动提交位移,业务处理成功后再手动确认——"先业务、后提交"是铁律。
5. 消费端拿到消息但处理到一半宕机,会丢吗?
答: 取决于 ACK / 提交位移的时机:
- 先处理业务、成功后再提交(正确做法):宕机时该消息尚未 ACK/未提交位移,Broker 会把它重新投递给其他消费者(或本消费者重启后从上次位移继续),不丢——但会重复消费,所以需要幂等;
- 先提交位移、再处理业务(错误做法):位移已推进,重启后直接从下一条开始消费,宕机的这条永久丢失。
结论:务必「先业务、后提交」。这条规则同时说明了为什么"至少一次"必然带来"可能重复"。
延伸:若业务处理本身包含"写多个库/发其他消息",还要考虑局部成功的情况——所以消费逻辑要设计成幂等且可重试,而不是假设"要么全成功、要么全失败"。
6. 如何设计消息的 TTL 与重试策略?
答: 按业务重要度分级设计,核心原则是"核心消息绝不丢,非核心消息可降级":
| 消息等级 | 重试策略 | TTL | 兜底 |
|---|---|---|---|
| 核心(支付、订单) | 无限重试 + 死信 | 不设 TTL 或设很长 | 死信告警 + 人工/定时补偿 |
| 重要(通知、积分) | 有限重试(如 5~10 次) | 按业务合理设置 | 死信 + 低优先级补偿 |
| 非核心(埋点、日志) | 少量重试(如 3 次) | 设置较短 TTL | 允许丢弃 |
三个关键设计点:
- 重试间隔要用指数退避(如 1s、5s、30s、5min…),避免下游故障时被重试流量打垮(雪崩);
- 重试要有上限并落入死信,否则"无限重试"会变成"无限积压";
- 不要用 TTL 来"清理积压"——TTL 到期是直接丢弃(见第 14 题),会造成静默的数据丢失。
三、幂等与重复消费
7. 为什么消息会被重复消费?
答: 重复的根源是「至少一次(at-least-once)」投递语义——为了保证不丢,各环节都会用"重试/重投"兜底,而重试就会带来重复:
| 环节 | 重复场景 |
|---|---|
| 生产端 | 发送超时但实际已成功,生产者的重试导致 Broker 收到两条 |
| 消费端 | 业务处理成功但 ACK/提交位移失败或超时,Broker 认为未消费而重投 |
| Rebalance | 分区重新分配给新消费者,而旧消费者的位移尚未提交,从旧位移重新消费 |
| Broker 复制 | 主从切换后,新 Leader 的位移回退,导致部分消息被再次投递 |
核心认知:MQ 一般只保证「不丢」,不天然保证「不重」——幂等必须由消费端自己实现。这是 MQ 面试的第一原则。
8. 如何实现消费幂等?(含唯一索引的注意事项)
答: 核心思路是让"同样的消息处理多次"与"处理一次"的业务效果完全相同。四种主流方案:
| 方案 | 做法 | 适用 |
|---|---|---|
| 业务唯一键 + 去重表(推荐) | 用业务唯一键建唯一索引,消费前 INSERT ... ON DUPLICATE KEY UPDATE(或 INSERT IGNORE),冲突则视为已处理、直接返回成功 | 通用,最简单可靠 |
| 状态机 / 版本号 | 业务状态单向流转(待支付 → 已支付),重复消息因"状态已变更"而被判定为无效操作 | 有明确状态机的业务(订单、工单) |
| Redis 去重 | 以 msgId/业务键为 key 做 SETNX 并设过期时间,已存在则跳过 | 高并发、可接受极小的不一致窗口 |
| 乐观锁 | UPDATE ... SET ... WHERE id=? AND version=?,影响行数为 0 说明已处理 | 更新类操作 |
用唯一索引做幂等的三个注意事项:
- 唯一键要代表"业务动作",而不是"消息 ID"——否则同一条业务的不同消息会被误拦(例如"订单支付"和"订单发货"不能共用一个键);
- 区分"已处理"与"真的失败":唯一键冲突说明已成功处理过,应返回成功(ACK);而业务异常(如参数非法)应返回失败并重试;
- 去重记录要能清理:去重表会持续膨胀,需按时间归档(如只保留 7 天),但要确保归档时间 > 消息可能重投的最大时间窗口。
选型建议:优先"业务唯一键 + 数据库唯一索引"——它把幂等交给数据库的强一致约束,最可靠;Redis 方案更快但存在"Redis 写入成功、业务失败"的窗口,需要更复杂的兜底。
四、顺序消息
9. 什么场景需要保证消息顺序?
答: 只要业务有状态流转、且步骤之间存在前后依赖,就必须保序,否则状态会错乱:
| 场景 | 说明 |
|---|---|
| 订单状态流转 | 「创建 → 支付 → 发货 → 完成」乱序会出现"未支付先发货" |
| 增量数据同步 | MySQL binlog 同步到 ES/数仓,乱序会导致下游数据与源库不一致 |
| 账户/积分变动 | 「+100、-50」乱序会读到中间态甚至错误余额 |
| 库存扣减 | 先扣后加与先加后扣的结果不同(若含条件判断) |
不需要保序的场景:彼此独立的通知、埋点、日志——盲目保序会白白牺牲并发度。
关键取舍:顺序性与并发度天然冲突。所以工程上只对「同一业务键」保序(如同一
orderId的消息有序),不同 key 之间仍然并行——这才是既有顺序又有吞吐的正确做法。
10. 保证消息顺序的通用手段有哪些?
答: 三个层面配合,缺一不可:
① 生产端:把需要保序的消息路由到同一个分区/队列
| MQ | 做法 |
|---|---|
| Kafka | 发送时指定 key(如 orderId),Kafka 按 hash(key) % partition 固定分区 |
| RocketMQ | 用 MessageQueueSelector 按 key 选择队列 |
| RabbitMQ | 发送到同一个 Queue(可用一致性哈希交换机按 key 路由) |
② 消费端:同一分区/队列必须串行消费
- 用单线程消费该分区(Kafka 中一个分区只能被同组的一个消费者消费,天然如此);
- 若要多线程并发,需要按 key 分桶:把同一 key 的消息交给固定的内存队列/线程处理,保证 key 内串行、key 间并行。
③ 生产端:避免"重试导致的乱序"
- Kafka:开启幂等生产者(
enable.idempotence=true,Kafka 3.0+ 默认开启)或把max.in.flight.requests.per.connection设为 1。否则"第一个请求失败重试、第二个请求先成功"会造成分区内乱序; - RabbitMQ:单 Channel 中连续发送并配合 confirm,避免重发到其他 Channel。
一句话总结:"同一 key → 同一分区 → 同一线程",这是顺序消息的万能公式;其余消息并行处理,兼顾吞吐。
五、消息积压治理
11. 如何发现消息积压?
答: 核心是监控 消费滞后量(Lag)——即"已生产但尚未被消费的消息数":
| MQ | 监控方式 |
|---|---|
| Kafka | kafka-consumer-groups.sh --describe --group xxx 看 LAG 列;生产上接入 Prometheus + kafka_exporter 看 kafka_consumergroup_lag |
| RocketMQ | mqadmin consumerProgress 查看各消费组的堆积量与消费位点差;控制台也有积压图 |
| RabbitMQ | 管理界面看队列的 messages_ready(待消费)与 messages_unacknowledged(已投递未确认) |
告警设计要点:
- 不要只看绝对值(1000 条可能对某业务是异常、对另一业务是正常);
- 更好的指标是 Lag 的增速(持续增长说明消费能力不足)与 消费延时(消息从生产到被消费的时间差);
- 结合业务 SLA 设阈值(如"订单消息必须 1 分钟内消费完")。
12. 消息大量积压该如何处理?
答: 按"先止血 → 再消化 → 后根治"三步走:
第一步:先恢复消费能力
- 扩容消费者实例,提升并行度;
- ⚠️ Kafka 的关键限制:消费者并发上限 = 分区数。若消费者数已等于分区数,再加实例也没用(多出来的实例会空闲);
- 检查是否有个别消费者"卡住"(如某分区在慢处理),必要时重启。
第二步:紧急降级 / 限速,避免雪崩
- 对非核心消息临时丢弃,或转到"旁路队列"异步补偿;
- 生产端限流(若下游已过载,继续灌入只会更糟);
- 优先消费核心业务消息(可按 Topic 隔离,避免非核心消息把核心消费者占满)。
第三步:存量快速消化(关键技巧)
当"单个分区消费太慢、而扩容又受分区数限制"时,用中转 Topic 并行消费:
- 临时创建一个多分区的新 Topic(分区数远大于原 Topic);
- 写一个"搬运程序":不做业务逻辑,只把旧 Topic 的消息快速转发到新 Topic,并按业务 key 打散到更多分区;
- 起N 倍数量的消费者消费新 Topic 并执行业务;
- 积压消化完后,切回原链路,并删除临时 Topic。
注意:以上无论如何都要先确认"是不是消费逻辑变慢了"——如果是下游接口超时/慢查询导致的,扩容只是把压力更快地传导给下游,可能引发更大的故障。扩容前先看一眼消费耗时与下游健康状况。
13. 消费者处理太慢,能否通过增加分区解决?
答: 能提升并行上限,但不是万能药,有几个关键限制:
| 限制 | 说明 |
|---|---|
| 存量数据不在新分区 | 已积压的历史消息仍在旧分区里,增加分区只对新写入的消息生效。要消化旧积压,还得配合消费者扩容或中转 Topic |
| 消费者并发上限 = 分区数 | 增加分区后必须同步扩容消费者实例,否则新分区没有消费者,等于白加 |
| 分区不能减少 | Kafka 只支持增加分区,不支持减少;分区过多会让元数据膨胀、Controller 压力变大、Rebalance 变慢 |
| 可能加剧乱序与再均衡 | 增加分区会改变 hash(key) % partition 的映射,导致同一 key 的历史消息与新消息落入不同分区,破坏顺序性 |
| Rebalance 频率上升 | 分区数多、消费者进出时,Rebalance 的协调成本更高 |
正确做法:
- 先定位为什么慢(慢 SQL?下游 RT 高?单条消息处理含阻塞调用?)——优化消费逻辑通常比加分区更有效;
- 需要提高并行度时,优先用"按 key 分桶 + 内存队列 + 线程池"在一个消费者内并发(key 内保序、key 间并行),而不是简单地加分区;
- 加分区需在业务低峰期做,并评估对顺序性的影响。
14. 消息积压时如果消息过期或队列写满怎么办?
答:
1)消息过期(TTL)——最危险的情况
- RabbitMQ 的消息 TTL 到期会被 Broker 直接清理丢弃(不是重投!)。此时"积压"其实已经不存在了,但数据真的丢了;
- 补救方式:只能批量重导——从业务库/binlog 查出这段时间的数据,重新补发到 MQ(例如订单消息可按
create_time范围重新投递); - 预防:核心消息不要依赖 TTL,应设"无限重试 + 死信 + 告警",或用 RocketMQ/Kafka(默认按时间/大小整体删除旧段,不做单条 TTL)。
2)队列写满(内存/磁盘触发流控)
- RabbitMQ:达到内存/磁盘水位会阻塞生产者(流控),保护 Broker 不崩,但上游会被拖住;
- Kafka:消息按时间/大小保留,写满磁盘需靠及时扩容磁盘 + 降低保留时长解决;
- 应急手段:临时"消费一条丢一条"快速腾空间(牺牲这部分数据),再走夜间补数;
- 常规手段:先扩容消费者 + 中转 Topic 加速消化(见第 12 题)。
关键认知:TTL 过期是"丢消息",不是"积压"。所以真正要防的不是积压,而是"消息在没人消费的情况下被静默删除"——这必须靠 Lag 告警 + 核心消息不设 TTL 来规避。
15. 如何避免分区数据倾斜导致的消费不均?
答: 消费倾斜指某些分区数据量远大于其他分区,导致"个别消费者一直忙、其余空闲",本质是分区键分布不均。
成因:
- 分区 key 本身分布不均(如按"地区"分区但业务集中在少数省);
- key 存在"超级热点"(如某个大商户的订单量占 80%);
- key 过少(如只有 3 个商户却分了 16 个分区,必然有分区空转)。
解决方案:
| 方案 | 说明 |
|---|---|
| 优化分区键 | 选离散度高的键(如 orderId 而非 merchantId);保证 key 数量远大于分区数 |
| 加盐打散 | 对热点 key 加随机后缀(key + "_" + random(N))分散到多个分区,但会破坏该 key 的顺序性,需与顺序需求权衡 |
| 按 key 分桶内并发 | 分区内再用"hash(key) % 线程数"把消息分给不同线程,让慢 key 不阻塞其他 key |
| 热点数据特殊处理 | 大商户走独立 Topic/独立消费组 |
| 监控分区级指标 | 对比各分区的 Lag 与消费速率,及时发现倾斜 |
根本认知:哈希分区只能保证"键的数量分布均匀",无法保证"消息量分布均匀"。真实业务量服从幂律分布,所以倾斜是常态,需要持续监控而非一次性设计解决。
六、死信、延迟与重试
16. 什么是死信队列(DLQ)?什么消息会变成死信?
答: 死信队列(Dead Letter Queue) 用于存放无法被正常消费的消息,作用是"不直接丢弃、便于事后排查与补偿"。
各 MQ 产生死信的条件:
| MQ | 死信条件 |
|---|---|
| RabbitMQ | ① 消息被消费者 basic.reject / basic.nack 且 requeue=false;② 消息TTL 过期;③ 队列达到最大长度(x-max-length)被丢弃/挤出的消息。这些消息会被投递到绑定的死信交换机(DLX),再路由到死信队列 |
| RocketMQ | 消费重试达到最大次数(默认 16 次,且间隔按等级递增)仍失败,消息进入 %DLQ%<consumerGroup> 死信 Topic |
| Kafka | 原生没有死信队列,需要自己在消费端捕获异常后把消息投递到专用 Topic 来实现 |
注意:死信 ≠ 丢弃。RabbitMQ 的 TTL 过期消息进死信队列,可以保留下来排查;而没有配 DLX 时,TTL 过期消息会被直接丢弃(见第 14 题)。
17. 延迟消息有哪些实现方案?
答: 延迟消息(定时消息)指"消息发送后不立即投递,而是延迟一段时间才可被消费",典型场景是"订单 30 分钟未支付自动取消"。
| 方案 | 实现 | 优缺点 |
|---|---|---|
| MQ 原生支持 | RocketMQ 原生支持:4.x 提供 18 个固定延迟级别(1s~2h),5.0 起支持任意秒级精度;RabbitMQ 需靠 TTL+DLX 或插件;Kafka 原生不支持 | ✅ 最可靠、无需自研 ❌ 各 MQ 支持度差异大 |
| RabbitMQ TTL + DLX | 消息设 TTL,过期后进入死信交换机再被消费 | ⚠️ 有队头阻塞坑(见《RabbitMQ》第 10 题) |
| RabbitMQ 延迟插件 | rabbitmq_delayed_message_exchange,声明延迟交换机 | ✅ 官方推荐 ❌ 需装插件 |
| 定时任务扫表 | 业务表存"到期时间",定时任务扫描到期记录并投递 | ✅ 实现简单 ❌ 精度低(受扫描间隔限制)、数据量大时扫表压力大(需按时间分片 + 索引优化) |
| Redis ZSet | 以"到期时间戳"为 score 存入 ZSet,定时轮询取出到期的 | ✅ 精度较高、性能好 ❌ 依赖 Redis 可靠性 |
| 时间轮(Time Wheel) | 用环形数组 + 链表按时间槽组织定时任务,Netty/Quartz 内部均采用 | ✅ 精度高、内存占用小、O(1) 插入 ❌ 需自研 |
选型建议:能用 RocketMQ 原生延迟消息就别自研;其次是 RabbitMQ 延迟插件;"定时扫表"只适合小数据量与低精度场景,它会随数据量增长而退化。
18. 死信消息该如何处理?
答: 死信不是"处理完了",而是"转入人工/异步流程",核心是可感知 + 可恢复:
- 必须报警:死信率、死信数量要接入监控告警(放任不管等于静默丢数据);
- 保留完整上下文:死信消息要能查到"原始消息体 + 消费失败的异常栈 + 重试次数 + 发生时间",否则无法定位;
- 提供重投工具:RocketMQ 的死信消息没有自动重投,需要写工具读取 DLQ Topic 后重新投递(修复根因后再重投);
- 区分处理策略:
- 由代码 bug / 下游临时故障导致的 → 修好后批量重投;
- 由脏数据 / 业务不合法导致的 → 人工确认后归档或丢弃,并修复上游校验逻辑;
- 控制死信队列的容量与保留时间,避免死信本身变成新的积压。
七、事务消息
19. 如何保证「本地事务执行」与「消息发送」的一致性?
答: 这是分布式事务的经典问题——本地事务提交成功但消息没发出去(或反之),会导致数据不一致。三种主流方案:
方案一:本地消息表(最大努力通知,最通用)
- 在业务库中建一张"消息表";
- 把「业务操作」与「写入消息记录」放进同一个本地事务——保证"业务成功则消息记录一定存在";
- 另起定时任务/异步线程轮询未发送的消息记录,投递到 MQ,成功后更新状态;
- 投递失败则重试,实现最终一致。
- ✅ 不依赖 MQ 的高级特性,任何 MQ 都能用;
- ❌ 需要额外的表与轮询任务;消息表会成为热点(可用分片/归档优化)。
方案二:MQ 事务消息(RocketMQ 原生支持)
- 基于半消息 + 回查两阶段提交(详见《RocketMQ》第 5 题);
- ✅ 无需自建消息表,由 MQ 保证"本地事务成功则消息一定可见";
- ❌ 仅支持具备事务消息能力的 MQ。
方案三:MQ 事务(Kafka)
- Kafka 事务解决的是"读-处理-写"(消费-转换-生产)的原子性,即 EOS(详见《Kafka》第 9 题);
- 注意它解决的不是"本地业务库事务 + 发消息"的一致性问题,两者问题域不同。
方案四:最大努力通知(不追求强一致)
- 业务处理完后异步重试通知下游,下游提供查询接口做对账兜底;
- 适合"通知类"场景(如支付结果通知),不保证不丢,靠对账补偿。
面试答题要点:先点出问题本质(本地事务与消息发送无法原子化),再给方案并说明适用场景。首选"本地消息表"(通用),若用 RocketMQ 则可用"事务消息"(更省事)。
八、消费模型与工程实践
20. 集群消费与广播消费的区别?
答:
| 模式 | 含义 | 适用场景 |
|---|---|---|
| 集群消费(Clustering) | 同一消费组内,一条消息只被一个消费者处理 | 绝大多数业务(并行消费、提升吞吐),RocketMQ/Kafka 默认 |
| 广播消费(Broadcasting) | 一条消息被组内所有消费者各处理一遍 | 刷新本地缓存、配置下发、多实例同步状态 |
各 MQ 的实现差异:
- RocketMQ:原生支持两种模式(
MessageModel.CLUSTERING/BROADCASTING); - Kafka:没有广播概念,等效做法是让每个实例使用独立的消费组(每个 group 都会收到全量消息);
- RabbitMQ:本质是队列模型,要实现"广播"需每个消费者绑定各自独立的 Queue 到同一个 Exchange(如 Fanout Exchange)。
注意广播模式的风险:消息量会随实例数放大(N 个实例 = N 倍消费),且无法用消费组位移统一管理,容易造成积压与重复。
21. pull 与 push 消费模型有何差异?
答:
| 维度 | Push(推) | Pull(拉) |
|---|---|---|
| 主导方 | Broker 主动推给消费者 | 消费者主动拉取 |
| 实时性 | 好(有消息立即推) | 取决于拉取间隔,可能延迟 |
| 消费端压力 | 易被压垮,必须做流控 | 由消费端控制速率,天然背压 |
| 批处理 | 不友好 | 友好(可一次拉一批,提升吞吐) |
| 空轮询 | 无 | 无消息时会空轮询,浪费资源 |
| 典型实现 | RabbitMQ(basic.consume) | Kafka |
关键真相(面试加分):
"RocketMQ 的 PushConsumer 底层其实是长轮询(Long Polling)的 Pull"——消费者发起拉取请求,若暂无消息,Broker hold 住请求(默认 15s)直到有新消息或超时才返回。这样既避免了空轮询,又获得了接近 Push 的实时性。
Kafka 也有类似机制:fetch.max.wait.ms + fetch.min.bytes 组合实现"攒够数据或超时才返回"的批量拉取。
一句话:Pull 是本质,Push 是"长轮询包装出来的体验";现代 MQ 都在向"Pull 的可控性 + Push 的实时性"靠拢。
22. 消息体过大有什么问题?如何优化?
答: 问题:大消息会同时压垮三个地方——网络带宽(传输慢)、Broker 内存/磁盘(Page Cache 与日志段膨胀、GC 压力)、消费端内存(反序列化开销大)。典型表现是吞吐骤降、Broker 频繁 Full GC。
优化手段:
| 手段 | 说明 |
|---|---|
| 消息体只放 ID / 必要字段(最推荐) | 完整数据存 DB / 对象存储 / 缓存,消费者拿到 ID 后再查——把"消息传递"变成"传递引用" |
| 压缩 | Kafka 支持 gzip/snappy/lz4/zstd;RocketMQ 支持压缩消息体。注意压缩是 CPU 换带宽 |
| 拆分大消息 | 把一个 10MB 的消息拆成多个小消息 + 一个"结束标记",消费端聚合 |
| 限制单条消息大小 | Kafka 默认 max.request.size=1MB、max.message.bytes=1MB;RocketMQ 默认最大 4MB。主动设限并在生产端拦截,避免大消息进入系统 |
| 走大文件通道 | 超大文件用对象存储(OSS/S3)+ 消息只传 URL |
反面案例:把整个商品详情、Excel 报表内容塞进消息体——正确做法是消息里只放
productId或fileId。
23. 如何做消息链路追踪?
答: 目标是在全链路上定位"消息卡在哪一段"。三个层次:
1)透传 traceId(必做)
- 生产端生成唯一的
traceId(或复用上游链路的traceId),写入消息头/属性(而非消息体,避免污染业务数据); - 消费端从消息头取出
traceId,绑定到当前线程的日志上下文(如 SLF4J MDC),保证消费日志也能关联到同一次请求; - 下游继续透传:若消费逻辑还会发新消息,把
traceId继续带下去。
2)打印关键节点的耗时
- 生产:发送开始/收到 ACK 的时间;
- Broker:入队时间
- 消费:拉取时间、业务开始/结束时间;
- 由此可算出"生产→入队""入队→被消费""被消费→业务完成"三段耗时,快速定位瓶颈。
3)借助工具
| 工具 | 作用 |
|---|---|
| APM(SkyWalking / Pinpoint / ARMS) | 自动串联跨 MQ 的调用链,可视化展示耗时 |
| MQ 自带消息轨迹 | RocketMQ 提供消息轨迹功能,可查询一条消息的产生、存储、消费全过程;RabbitMQ 可用 Firehose Tracer(性能开销大,仅调试用) |
| 消费端打点上报 | 自定义指标(消费耗时分布、失败率)上报到 Prometheus,配合 Lag 监控形成完整可观测性 |
实践建议:"traceId 透传 + 三段耗时打点 + Lag 监控" 是投入产出比最高的组合;APM 是锦上添花。
24. 消息过滤是如何实现的?
答: 消息过滤的目标是"让消费者只收到自己关心的消息,减少无效传输与无效处理"。按"在哪过滤"分为两类:
| 类型 | 说明 | 代表 |
|---|---|---|
| Broker 端过滤(推荐) | 过滤在服务端完成,不匹配的消息根本不会传给消费者,节省带宽与消费端 CPU | RocketMQ 的 Tag 过滤(Producer 打 Tag,Consumer 订阅指定 Tag)、RocketMQ 的 SQL92 属性过滤(按消息属性写 SQL 表达式)、RabbitMQ 的 Exchange 路由规则(routing key / Headers) |
| 消费端过滤 | 消费者拉到全部消息后自行判断丢弃 | Kafka 没有 Broker 端过滤,只能在消费端过滤(或按 Topic/分区拆分) |
注意事项:
- RocketMQ 的 Tag 过滤在 Broker 端是"哈希比对",效率极高;而 SQL92 过滤需要解析表达式,会消耗 Broker 的 CPU,属性不宜过多、表达式不宜复杂;
- 不要用"消费端过滤"来替代"Topic 拆分":如果两类消息的消费逻辑与吞吐差异很大,应该用不同的 Topic,而不是发到一个 Topic 再过滤——因为过滤不掉的是带宽和存储;
- Kafka 中若确实需要服务端过滤能力,常见做法是用不同 Topic 或不同分区键区分,或引入 Pulsar/支持过滤的中间件。
25. 什么是消息回溯(消息重放)?有哪些应用场景?
答: 消息回溯指重新消费已经消费过的历史消息——本质是把消费者的位移(offset)重置到更早的位置。
应用场景:
| 场景 | 说明 |
|---|---|
| 数据修复 | 下游消费逻辑有 bug,把错误数据写入后,修好代码再重放这段时间的消息修正数据 |
| 新增下游 | 新上线一个消费方(如新加一套数仓),需要把历史全量数据补给它 |
| 故障恢复 | 消费端长时间宕机,或误提交了位移导致消息被跳过,需要回退位移重新消费 |
| 压测/演练 | 用历史消息回放来验证下游的容量与稳定性 |
各 MQ 的实现方式:
| MQ | 方式 |
|---|---|
| Kafka | kafka-consumer-groups.sh --reset-offsets --to-datetime/--to-offset/--shift-by;或对新消费组设置 auto.offset.reset=earliest(新组会从最早位置开始读) |
| RocketMQ | 按时间回溯(ConsumeFromTimestamp,控制台按时间点重置消费位点);5.0 也支持按 offset 重置 |
| RabbitMQ | 不支持回溯——消息被 ACK 后即从队列删除,只能靠"重新生产"(从业务库/binlog 补发) |
两个关键前提(面试常追问):
- 消息还没被删除:回溯能力受消息保留策略限制——Kafka 默认保留 7 天(
retention.ms)、RocketMQ 默认 3 天。超出保留期的消息无法回溯,所以"数据修复"必须在保留窗口内做; - 消费端必须幂等:重放会把消息再投递一次,若下游不幂等会造成数据重复。这也是"幂等"如此重要的另一个原因。
一句话对比:Kafka/RocketMQ 靠"位移可重置 + 消息保留"实现回溯;RabbitMQ 靠"ACK 即删"的设计,天生不支持回溯——这是选型时要考虑的一个重要差异。
