Skip to content

消息队列

消息队列面试的主线:为什么用(解耦/异步/削峰)→ 怎么保证不丢 → 怎么保证不重复 → 怎么保证有序。以 Kafka 和 RabbitMQ 为参照。

Q1: 为什么要用消息队列?有什么缺点? 「🟢 校招/初级」

考察点:技术选型的收益与成本分析。

参考答案

三大作用:

  • 解耦:下游系统增减不影响上游。
  • 异步:非核心逻辑(发短信、记日志)异步化,缩短主链路耗时。
  • 削峰:突发流量先进队列,消费端按能力处理。

代价(必须主动说):系统复杂度上升、可用性问题(MQ 挂了整条链路受影响)、一致性问题(消息丢失、重复、乱序都要额外处理)。

追问延伸

  • 什么场景不适合用消息队列?(强一致同步调用)
  • Kafka、RabbitMQ、RocketMQ 怎么选型?

Q2: 怎么保证消息不丢失? 「🟡 中级」

考察点:端到端可靠性设计,必问题。

参考答案

按三个环节逐一分析:

  • 生产端:开启确认机制(Kafka 的 acks=all、RabbitMQ 的 publisher confirm),发送失败重试;同步发送拿回调。
  • Broker 端:Kafka 设置副本数 ≥3 + min.insync.replicas ≥2,禁用不干净选举;RabbitMQ 队列和消息都持久化。
  • 消费端:关闭自动确认,业务处理成功后再手动提交(Kafka 关 enable.auto.commit;RabbitMQ 手动 ack)。

追问延伸

  • acks=all 有什么性能代价?
  • 先消费再提交和先提交再消费各有什么问题?

Q3: 怎么保证消息不重复(幂等)? 「🟡 中级」

考察点:重复不可避免,幂等才是正解。

参考答案

  • 重复的原因:生产端重试、消费端处理成功但 ack 失败、消费者重启重新消费。所以重复无法根除,只能让重复无害
  • 幂等手段:
    • 唯一键约束(订单号入库去重)。
    • 数据库乐观锁/状态机(UPDATE ... WHERE status=待支付)。
    • Redis SETNX 记录已处理的消息 ID。
    • 业务天然幂等(SET 覆盖写、删除操作)。

追问延伸

  • 幂等和"只消费一次"(exactly once)是什么关系?
  • Kafka 的事务消息能做到什么程度?

Q4: 怎么保证消息有序? 「🟡 中级」

考察点:分区模型的理解。

参考答案

  • MQ 只能保证分区/队列内有序,跨分区无序。
  • 方案:把需要有序的消息用同一个 key(如订单号)路由到同一分区,分区内单线程顺序消费。
  • 代价:热点分区问题;消费端出异常时不能跳过(否则破坏顺序),需要重试或暂停该分区。
  • RabbitMQ:单 queue + 单 consumer 才能保证全局有序,吞吐受限。

追问延伸

  • 订单的"创建→支付→发货"全局有序怎么设计?
  • 消费失败又必须保序时怎么办?(本地重试表 + 人工兜底)

Q5: 消息积压了怎么处理? 「🔴 高级」

考察点:线上事故的应急处理能力。

参考答案

  1. 止血:确认消费端是否有 bug(一直失败重试),有的话先修复或临时跳过。
  2. 扩容消费端:增加消费者实例,前提是有足够分区(Kafka 消费者数 ≤ 分区数)。
  3. 临时扩容方案:新建一个分区数翻倍的临时 topic,写一个临时消费程序只做转发,再用多倍消费者消费临时 topic。
  4. 预防:消费端压测基线、积压监控告警、大流量入口提前限流。

追问延伸

  • 为什么不直接给原 topic 加分区?(原数据仍在旧分区,且破坏顺序性)
  • 积压的消息过期丢了怎么办?(死信队列 + 补偿)

Q6: Kafka 的高性能是怎么做到的? 「🟡 中级」

考察点:对高性能中间件设计的理解。

参考答案

  • 顺序写磁盘:消息追加写入,顺序写速度接近内存。
  • 页缓存(PageCache):读写利用操作系统缓存,减少用户态-内核态拷贝。
  • 零拷贝:消费用 sendfile 直接从页缓存到网卡。
  • 分区并行:多分区分散到多 broker、多磁盘。
  • 批量与压缩:生产者批量发送、压缩(如 lz4),消费者批量拉取。

追问延伸

  • Kafka 怎么选举 Leader?ISR 是什么?
  • 消费者 rebalance 是怎么触发的?有什么影响?

Q7: 什么是死信队列和延迟队列? 「🟡 中级」

考察点:消息系统的进阶机制。

参考答案

  • 死信队列(DLQ):消费失败重试超限、消息过期被拒收后进入死信队列,供人工排查和补偿,避免问题消息卡住正常流程。
  • 延迟队列:消息发送后延迟一段时间才可被消费。实现:RabbitMQ 的 TTL + 死信、Redis ZSet(score 为到期时间轮询)、RocketMQ 原生支持延迟级别、Kafka 需要自建(时间轮/分层 topic)。
  • 应用:订单超时取消、定时提醒、重试退避。

追问延伸

  • 大量不同延迟时间的消息用 ZSet 实现要注意什么?
  • 时间轮的原理?(Netty 的 HashedWheelTimer)

Q8: Kafka 的架构是怎样的?Topic、Partition、Replica、ISR 分别是什么? 「🟡 中级」

考察点:Kafka 核心概念和架构理解。

参考答案

核心概念

  • Broker:Kafka 集群中的服务器节点,负责存储和转发消息。
  • Topic:消息的主题(分类),逻辑上的队列,生产者往 Topic 发,消费者从 Topic 读。
  • Partition:分区,Topic 的物理分片,一个 Topic 可以有多个 Partition
    • 每个 Partition 是一个有序的、不可变的消息队列(追加写)
    • 消息按 offset 顺序编号,offset 是消费者消费进度的标记
    • 分区是 Kafka 高吞吐的关键(并行读写、分散负载)
  • Replica:副本,每个 Partition 有多个副本(1 个 Leader + N 个 Follower)
    • Leader:处理该分区的所有读写请求
    • Follower:从 Leader 同步数据,做备份,不对外提供服务
    • 副本机制提高可用性(Leader 挂了可以选新 Leader)
  • ISR(In-Sync Replicas):同步中的副本列表
    • 和 Leader 保持同步的副本集合(落后不多)
    • 只有 ISR 中的副本才有资格被选为 Leader
    • 判定条件:replica.lag.time.max.ms 时间内同步过数据

数据分布特点

  • 每个 Partition 的多个副本分布在不同 Broker 上(容灾)
  • Leader 均匀分布在各个 Broker(负载均衡)
  • 生产者写入 Leader,消费者从 Leader 消费

追问延伸

  • 分区数是不是越多越好?(不是,越多文件句柄越多、元数据越多、rebalance 越慢、延迟越高)
  • ISR 收缩和扩张的条件是什么?(落后超过时间阈值收缩;追赶上了扩张)

Q9: Kafka 消费者组(Consumer Group)和 rebalance 机制? 「🔴 高级」

考察点:Kafka 消费模型和重平衡的深入理解。

参考答案

消费者组

  • 一组消费者共同消费一个 Topic,每个分区只分配给组内一个消费者
  • 组 ID(group.id)相同的消费者属于同一个组
  • 消费模式:
    • 点对点:一个组消费,消息只被消费一次
    • 发布订阅:多个组消费,每个组都能收到所有消息

Rebalance(重平衡)

消费者组内消费者数量变化、Topic 分区数变化时,重新分配分区和消费者的对应关系。

触发条件

  • 消费者加入组(新增实例)
  • 消费者离开组(实例挂了、主动离开、session 超时)
  • 订阅的 Topic 分区数变化(扩容分区)
  • 订阅的 Topic 数量变化(正则订阅匹配到新 Topic)

Rebalance 过程(以消费端 Coordinator 为例)

  1. Join Group:所有消费者申请加入组,选一个消费者作为 Leader(负责分配方案)
  2. Sync Group:Leader 制定分配方案,发给 Coordinator,Coordinator 同步给所有消费者
  3. 消费者拿到分配结果,开始消费

问题和影响

  • Rebalance 期间整个组暂停消费(STW,Stop The World)
  • 重复消费(还没提交 offset 就 rebalance 了,新消费者从上次提交的 offset 重新消费)
  • 频繁 rebalance 严重影响性能和消费延迟

追问延伸

  • 怎么减少 rebalance 的影响?(静态成员、增量协商、调大会话超时/心跳间隔、减少不必要的 rebalance)
  • 心跳和会话超时的关系?(session.timeout.ms 判定是否存活,heartbeat.interval.ms 心跳发送频率,心跳间隔 < 会话超时)

Q10: 消息队列选型对比?Kafka vs RabbitMQ vs RocketMQ vs Pulsar? 「🟡 中级」

考察点:技术选型能力,对各 MQ 特点的了解。

参考答案

特性KafkaRabbitMQRocketMQPulsar
开发语言Java / ScalaErlangJavaJava
模型发布订阅(分区模型)AMQP(Exchange + Queue)发布订阅 + 队列模型发布订阅(分区模型)
吞吐量极高(百万级)中(万级)高(十万级)极高
延迟毫秒级微秒级毫秒级毫秒级
可靠性高(副本机制)高(持久化 + 确认)
顺序消息分区内有序队列内有序支持(MessageQueue)分区内有序
延迟消息不支持(需自建)插件支持(TTL + DLX)原生支持(18 个级别)原生支持
事务消息支持不支持支持支持
死信队列不支持(需自建)支持支持支持
回溯消费支持(按 offset / 时间)不支持支持(按时间)支持
生态大数据生态极好路由灵活、插件多阿里系、金融场景好云原生、存算分离

选型建议:

  • 大数据 / 日志采集 / 高吞吐场景 → Kafka
  • 路由复杂、业务消息、需要灵活 Exchange → RabbitMQ
  • 金融 / 电商 / 事务消息 / 阿里云生态 → RocketMQ
  • 云原生 / 多租户 / 存算分离 / 大规模 → Pulsar

追问延伸

  • RabbitMQ 的 Exchange 类型有哪些?(direct、fanout、topic、headers)
  • Pulsar 的存算分离是什么意思?(Broker 无状态负责计算,BookKeeper 负责存储,可独立扩容)

Q11: RocketMQ 的事务消息是怎么实现的? 「🔴 高级」

考察点:事务消息的实现原理,分布式事务的一种方案。

参考答案

RocketMQ 事务消息:保证本地事务和消息发送的原子性(要么都成功,要么都失败)。

实现流程(两阶段提交 + 回查补偿)

第一阶段:发送半消息(Half Message)

  • 生产者发送半消息到 Broker
  • 半消息对消费者不可见(存在 RMQ_SYS_TRANS_HALF_TOPIC 中)
  • Broker 返回发送成功确认

第二阶段:执行本地事务 + 提交 / 回滚

  • 生产者执行本地事务(如更新数据库)
  • 本地事务成功 → 发送 Commit 消息给 Broker → 半消息变成普通消息,消费者可见
  • 本地事务失败 → 发送 Rollback 消息给 Broker → 删除半消息

回查机制(补偿)

  • 如果第二阶段的 Commit / Rollback 消息丢失或超时未收到
  • Broker 会定期扫描半消息,回查生产者事务状态
  • 生产者回查本地事务状态,返回 Commit / Rollback / Unknown
  • 多次回查还是 Unknown 则默认回滚(可配置)

应用场景

  • 订单创建 + 扣库存(分布式事务)
  • 本地操作 + 发消息需要原子性的场景
  • 替代本地消息表方案(不需要额外建表)

追问延伸

  • 事务消息和本地消息表有什么区别?(事务消息由 MQ 中间件保证原子性,本地消息表需要业务方自己维护)
  • 半消息为什么对消费者不可见?(事务还没提交,消费者看到了会出现脏读)

Q12: 消息的幂等怎么实现?有哪些工程实践? 「🟡 中级」

考察点:幂等设计的实战经验。

参考答案

幂等:同一个操作执行一次和执行多次结果一样。消息队列中重复消费不可避免(重试、网络超时等),所以消费端必须做幂等。

实现方案

1. 数据库唯一键约束

  • 利用主键或唯一索引,重复插入报错(或 IGNORE)
  • 最简单可靠,适合写入操作
  • 如:订单号作为唯一键,重复创建会失败

2. 数据库乐观锁

  • 加版本号字段:UPDATE ... SET status=新状态, version=version+1 WHERE id=xxx AND version=旧版本
  • 适合更新操作,并发不高的场景

3. 状态机幂等

  • 业务有明确状态流转(如:待支付 → 已支付 → 已发货 → 已完成)
  • UPDATE ... WHERE status=待支付(状态不对就更新不了,天然幂等)
  • 适合有状态流转的业务

4. Redis SETNX 去重

  • 用消息 ID / 业务唯一键作为 key,SETNX 设置成功才处理
  • 设置过期时间(防止内存泄漏)
  • 适合处理速度快、量大的场景

5. Token 机制

  • 操作前先申请 token(存 Redis),操作时带上 token,处理完删除 token
  • 适合前端防重复提交、接口幂等

6. 下游天然幂等

  • 如 SET 覆盖写、删除操作、select 查询等
  • 利用业务特性实现幂等,零成本

选型建议

  • 写入操作:数据库唯一键约束最可靠
  • 更新操作:状态机 / 乐观锁
  • 消费端处理:Redis 去重表 + 数据库唯一键兜底(双层保障)

追问延伸

  • Redis 去重表 key 过期了怎么办?(配合数据库唯一键兜底,Redis 是快速失败,数据库是最终保障)
  • Exactly Once 和幂等的关系?(Exactly Once 是传输语义,幂等是业务效果;At-Least-Once + 幂等 = 业务上的 Exactly Once)

Q13: 怎么保证消息的顺序性?全局有序和分区有序的权衡? 「🔴 高级」

考察点:消息有序性的深入理解和设计权衡。

参考答案

有序性的两个层次

1. 分区有序(局部有序)

  • 同一个分区内的消息有序
  • 实现:相同 key 的消息路由到同一个分区 + 单线程消费
  • 优点:吞吐高(多个分区并行消费)
  • 缺点:只能保证同一 key 内有序,跨 key 无序
  • 适用:大部分业务场景(如同一订单的状态变更有序即可)

2. 全局有序

  • 整个 Topic 内所有消息严格有序
  • 实现:只有 1 个分区 + 1 个消费者
  • 优点:全局严格有序
  • 缺点:吞吐极低(单分区单消费者,完全没有并行度)
  • 适用:极少场景(如数据库 binlog 同步的严格有序、配置变更广播)

保证有序的注意事项

  • 发送端:相同 key 的消息必须按顺序发送(不能并发发同一个 key 的消息)
  • Broker 端:同一个分区内有序(Kafka / RocketMQ 天然保证)
  • 消费端:必须单线程顺序消费(不能多线程消费同一个分区)
  • 消费失败不能跳过:失败了要重试,否则顺序就断了
    • 处理方式:同步重试 → 还是失败就进死信队列 → 人工介入补偿

常见坑

  • 发送端重试导致乱序(Kafka 不会,因为重试是追加的;但如果开启了 max.in.flight.requests.per.connection > 1 且 acks < all 可能乱序)
  • 消费端多线程消费同一个分区导致乱序
  • 分区扩容导致 key 路由变化(Hash 取模变化,同一个 key 可能到不同分区)
  • 消息重试机制可能破坏顺序(重试的消息到了队尾,后面的消息先被消费了)

追问延伸

  • 订单的"创建 → 支付 → 发货 → 完成"怎么保证有序?(订单号作为 key,路由到同一分区)
  • 消费失败但必须保序,重试也不成功怎么办?(死信队列 + 人工补偿,不能跳过,跳过就破坏顺序了)

Q14: Kafka 的 Leader 选举机制是怎样的? 「🔴 高级」

考察点:Kafka 高可用核心机制的深入理解。

参考答案

Kafka 的 Leader 选举涉及两个层面:

  1. Controller 选举(集群级):

    • ZooKeeper 临时节点 /zookeeper/controller,先创建的 Broker 成为 Controller
    • Controller 负责管理分区和副本状态、分区 Leader 选举、Broker 上下线
    • Controller 挂了,ZooKeeper 临时节点消失,其他 Broker 竞争成为新 Controller
  2. Partition Leader 选举(分区级):

    • 由 Controller 负责
    • Leader 所在 Broker 挂了 → Controller 从 ISR 中选第一个作为新 Leader
    • ISR(In-Sync Replicas):和 Leader 保持同步的副本
    • 如果 ISR 为空:
      • unclean.leader.election.enable=true(默认 false):从非 ISR 中选 Leader → 可能丢数据
      • unclean.leader.election.enable=false:分区不可用,等 ISR 中的副本恢复

选举流程:

  1. Broker 心跳超时 → ZooKeeper 通知 Controller
  2. Controller 遍历该 Broker 上的所有 Leader 分区
  3. 为每个分区从 ISR 中选新 Leader
  4. 更新元数据,通知相关 Broker
  5. 客户端感知 Leader 变化(通过元数据刷新)

追问延伸

  • 为什么默认关闭 unclean leader 选举?
  • Controller 单点问题怎么解决?(KIP-500 移除 ZooKeeper 依赖)

Q15: Kafka 的日志压缩(Log Compaction)是什么?和日志清理有什么区别? 「🟡 中级」

考察点:Kafka 存储机制的深入理解。

参考答案

两种日志清理策略:

  1. 日志删除(Log Deletion)(默认):

    • 按时间或大小删除旧消息
    • retention.ms(默认7天)/ retention.bytes
    • 适合:日志、事件流等只需最近数据的场景
  2. 日志压缩(Log Compaction)

    • 按 key 保留最新值,旧版本被清理
    • 同一个 key 只保留最后一条消息
    • 类似 LSM 数据库的 compaction
    • 适合:状态存储、配置变更、数据库 CDC(Change Data Capture)

Log Compaction 原理:

  • 日志分为 Head(最新未压缩)和 Tail(已压缩)两部分
  • 后台线程定期扫描 Tail 段,对同一 key 保留最新值
  • 压缩后生成 .clean 文件,替换原始段

Segment 结构:

  • .log:消息数据
  • .index:偏移量索引(稀疏索引)
  • .timeindex:时间戳索引
  • .snapshot:清理快照

应用场景:

  • Kafka Streams 的状态存储(changelog topic)
  • __consumer_offsets topic(存消费进度)
  • 配置变更通知(key=配置名,value=配置值)

追问延伸

  • Log Compaction 和数据库的 MVCC 有什么相似之处?
  • 一个 topic 能同时用删除和压缩策略吗?

Q16: Kafka 消费者的分区分配策略有哪些? 「🔴 高级」

考察点:消费者组内分区分配算法的理解。

参考答案

三种分配策略:

  1. Range(默认)

    • 按主题分配,每个主题的分区按范围分配给消费者
    • 如:主题有 P0-P9,3 个消费者 → C0: P0-P3, C1: P4-P6, C2: P7-P9
    • 缺点:第一个消费者总是分到更多(不公平)
  2. RoundRobin

    • 跨主题轮流分配
    • 把所有主题的所有分区排成列表,轮流分给消费者
    • 更均匀
    • 要求:消费者实例的订阅主题必须相同
  3. StickyAssignor(粘性分配,推荐):

    • 尽量保持原有分配不变(减少 rebalance 时的分区迁移)
    • rebalance 时只重新分配变更的部分
    • 更均匀 + rebalance 开销小
  4. CooperativeStickyAssignor(增量协作式,2.4+):

    • 基于 StickyAssignor,但 rebalance 是增量的(不停止所有消费)
    • 先撤销需要变更的分区,其他消费者继续消费
    • 大幅减少 rebalance 停顿时间

配置:partition.assignment.strategy 参数

追问延伸

  • Range 为什么不公平?
  • 增量协作式 rebalance 比全量好在哪里?

Q17: Kafka 的 offset 提交机制?手动提交 vs 自动提交? 「🟡 中级」

考察点:消费进度管理的实战理解。

参考答案

offset 存储:消费者组的 offset 存在 __consumer_offsets topic 中(key: 组ID+主题+分区, value: offset)。

两种提交方式:

  1. 自动提交

    • enable.auto.commit=true(默认)
    • 定时自动提交(auto.commit.interval.ms,默认 5 秒)
    • 优点:简单
    • 缺点:
      • 可能重复消费(处理完消息还没提交就挂了)
      • 可能丢失消费(提交了但还没处理完就挂了)
      • 不精确
  2. 手动提交

    • enable.auto.commit=false
    • commitSync():同步提交(阻塞直到成功,更可靠但慢)
    • commitAsync():异步提交(不阻塞,但可能失败)
    • 可以提交到指定 offset(更精确控制)

最佳实践:

  • 重要业务:关闭自动提交,处理完成后同步提交
  • 性能要求高:异步提交 + 定期同步提交(兜底)
  • 精确控制:处理完一批后提交这批的最大 offset
  • 结合幂等设计:offset 提交和业务操作不原子,可能重复消费
java
// 手动同步提交
consumer.subscribe(Arrays.asList("topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        // 处理消息
    }
    consumer.commitSync(); // 处理完一批后提交
}

追问延伸

  • offset 提交失败会怎样?(下次从这里重新消费,可能重复)
  • 怎么保证业务处理和 offset 提交的原子性?(事务消息 / 本地消息表)

Q18: RabbitMQ 的 Exchange 类型有哪些?分别适合什么场景? 「🟡 中级」

考察点:RabbitMQ 路由模型的理解。

参考答案

四种 Exchange 类型:

  1. direct(直连):

    • Routing Key 完全匹配 Binding Key
    • 一条消息发给一个队列
    • 场景:点对点、任务分发
  2. fanout(扇出):

    • 忽略 Routing Key,广播给所有绑定的队列
    • 一条消息发给所有队列
    • 场景:广播通知、多系统同步
  3. topic(主题):

    • Routing Key 模糊匹配 Binding Key(支持通配符)
    • * 匹配一个单词,# 匹配零个或多个单词
    • 场景:按规则订阅、日志分级
    • 如:order.*.created 匹配 order.normal.createdorder.vip.created
  4. headers(头部):

    • 根据消息头(headers)匹配,忽略 Routing Key
    • 支持 x-match: all(全部匹配)或 x-match: any(任一匹配)
    • 性能最差,少用

RabbitMQ 消息流转模型:

  • Producer → Exchange → (Binding) → Queue → Consumer
  • Exchange 和 Queue 之间通过 Binding 关联
  • 消息有 Routing Key,Exchange 根据类型和 Binding 规则路由到 Queue

追问延伸

  • direct 和 topic 的区别?
  • fanout 模式下消息会被复制到每个队列吗?

Q19: RocketMQ 的架构和 NameServer vs ZooKeeper 的区别? 「🔴 高级」

考察点:RocketMQ 架构设计和注册中心选型的理解。

参考答案

RocketMQ 架构:

  • NameServer:注册中心,管理 Broker 路由信息
  • Broker:消息存储和转发,分 Master/Slave
  • Producer:生产者,从 NameServer 获取路由后直连 Broker
  • Consumer:消费者,从 NameServer 获取路由后直连 Broker

NameServer vs ZooKeeper(为什么 RocketMQ 不用 ZK):

  • NameServer:无状态、轻量、可集群部署但互不通信
    • 优点:简单、部署方便、性能好
    • 数据一致性:各节点独立接收 Broker 注册,不保证一致性
  • ZooKeeper:有状态、强一致(ZAB 协议)
    • 缺点:复杂、运维成本高、性能瓶颈
    • 早期 RocketMQ 用过 ZK,后来换成自研 NameServer

Broker 部署模式:

  • 单 Master:单点,不推荐
  • 多 Master:无 Slave,Master 挂了消息可能丢
  • 多 Master 多 Slave(同步刷盘):Master 和 Slave 同步写,数据安全但性能稍低
  • 多 Master 多 Slave(异步刷盘):Master 写完返回,Slave 异步同步,性能高但有丢数据风险

Producer 发消息流程:

  1. 从 NameServer 获取 Topic 的路由信息(哪些 Broker 有该 Topic 的队列)
  2. 选择一个队列(负载均衡)
  3. 直连 Broker 发消息
  4. Broker 写入 CommitLog → 异步构建 ConsumeQueue(索引)
  5. 返回发送结果

追问延伸

  • RocketMQ 的存储结构(CommitLog + ConsumeQueue)是什么?
  • NameServer 互相不通信怎么保证 Broker 信息一致?(Broker 向所有 NameServer 注册)

Q20: Kafka Producer 的发送流程是怎样的? 「🟡 中级」

考察点:生产者端完整链路的理解。

参考答案

Producer 发送流程:

  1. 序列化:Serializer 把 key/value 序列化为字节数组
  2. 分区选择:Partitioner 决定消息发往哪个分区
    • key 不为空:对 key 哈希取模(Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions
    • key 为空:轮询(Sticky Partitioner,黏性分区,4.0前默认轮询所有分区)
    • 自定义:实现 Partitioner 接口
  3. 追加到批次:消息加入 RecordAccumulator(消息累加器)
    • 每个分区对应一个 Deque<RecordBatch>
    • batch.size(默认16KB):批次满则发送
    • linger.ms(默认0):等待时间,即使批次不满也发送
  4. Sender 线程:后台线程从累加器取出就绪批次
    • buffer.memory(默认32MB):缓冲区大小
    • max.in.flight.requests.per.connection(默认5):每个连接未确认请求数
  5. 网络发送:通过 NIO 发送给 Leader Broker
  6. 重试机制:发送失败可重试
    • retries(默认Integer.MAX_VALUE):重试次数
    • retry.backoff.ms(默认100ms):重试间隔
    • delivery.timeout.ms(默认120s):总投递超时
  7. ACK 确认:等待 Broker 的 ACK
    • acks=0:不等确认
    • acks=1:Leader 确认即可
    • acks=all/-1:ISR 全部确认(最安全)

幂等 Producer(0.11+):

  • enable.idempotence=true:Producer 分配 PID + SequenceNumber,Broker 去重
  • 保证单分区内不重复

事务 Producer(0.11+):

  • transactional.id:跨分区和会话的事务
  • initTransactions → beginTransaction → send → commitTransaction

追问延伸

  • linger.ms=0 和 linger.ms=10 的区别?
  • max.in.flight.requests.per.connection > 1 和幂等的关系?

Q21: Kafka 的 Segment 文件结构是怎样的? 「🟡 中级」

考察点:Kafka 底层存储结构的深入理解。

参考答案

Kafka 存储结构层级:

  • Topic → Partition → Segment → File

Segment 文件组成(每个 Segment 一组文件):

  1. .log 文件:实际消息数据
    • 消息按顺序追加写入(不可变)
    • 每条消息有 offset、key、value、timestamp
  2. .index 文件:偏移量索引(稀疏索引)
    • 记录「逻辑 offset → 物理位置」的映射
    • 稀疏索引:不是每条消息都有索引,间隔约 4KB 一个索引项
    • 查找时用二分查找定位最近的索引,再顺序扫描 .log 文件
  3. .timeindex 文件:时间戳索引
    • 记录「timestamp → offset」的映射
    • 用于按时间查找消息(如 consumer.seek by timestamp)
  4. .snapshot 文件:快照(事务相关)

Segment 切分条件:

  • segment.bytes(默认1GB):文件大小达到阈值切新 Segment
  • segment.ms(默认7天):时间达到阈值切新 Segment

查找消息流程(如查找 offset=100):

  1. 二分查找确定在哪个 Segment(按文件名 offset 范围)
  2. 在 .index 文件中二分查找 ≤ 100 的最大索引项
  3. 从 .log 文件的该物理位置开始顺序扫描
  4. 找到 offset=100 的消息

清理策略:

  • 日志删除(Delete):按时间/大小删除整个 Segment
  • 日志压缩(Compact):按 key 保留最新值(在 Q15 已覆盖)

追问延伸

  • 为什么 .index 是稀疏索引而不是每条消息都有索引?
  • 怎么查看 .log 文件的内容?

Q22: Kafka 副本同步机制(ISR)的详细原理? 「🔴 高级」

考察点:Kafka 高可用核心机制的深度理解。

参考答案

副本同步核心概念:

  • LEO(Log End Offset):每个副本的日志末端 offset(下一条要写入的位置)
  • HW(High Watermark):所有 ISR 副本中最小的 LEO,消费者只能看到 HW 之前的消息
  • ISR(In-Sync Replicas):和 Leader 保持同步的副本集合

同步流程(Follower 从 Leader 拉取数据):

  1. Follower 发送 FetchRequest,携带自己的 LEO
  2. Leader 根据请求返回 LEO 之后的数据
  3. Follower 收到数据,追加到本地日志,更新自己的 LEO
  4. Leader 收到 Follower 的 FetchRequest 后,更新远端 LEO(RemoteLEO)
  5. Leader 更新 HW = min(所有 ISR 的 LEO)
  6. Leader 在下次 FetchResponse 中把 HW 返回给 Follower
  7. Follower 更新自己的 HW

ISR 收缩(移出 ISR):

  • 条件:Follower 超过 replica.lag.time.max.ms(默认30秒)未同步
    • 即 Follower 的 LEO 一直落后于 Leader 的 LEO,且时间超过阈值
  • 移出 ISR 后,Follower 不再有资格被选为 Leader

ISR 扩张(加入 ISR):

  • 条件:Follower 追上 Leader(LEO >= Leader HW)
  • Follower 持续拉取数据直到追上

关键配置:

  • min.insync.replicas:ISR 最少副本数,不满足则 Producer 写入失败(和 acks=all 配合)
  • replica.lag.time.max.ms:Follower 同步超时时间

追问延伸

  • HW 和 LEO 的区别?
  • 如果 ISR 为空,acks=all 会怎样?(写入失败,保证不丢数据)

Q23: RocketMQ 的存储结构 CommitLog + ConsumeQueue 是什么? 「🔴 高级」

考察点:RocketMQ 存储设计的深入理解。

参考答案

RocketMQ 存储核心:

  • 所有 Topic 的消息都写入同一个 CommitLog(一个大文件,顺序写)
  • ConsumeQueue 是逻辑索引,每个 Topic+QueueID 对应一个 ConsumeQueue

CommitLog:

  • 文件:固定大小 1GB,写满新建文件
  • 所有消息统一追加写入 CommitLog(顺序写,性能极高)
  • 消息格式:总长度、MagicCode、CRC、QueueID、QueueOffset、物理偏移量、消息体等

ConsumeQueue:

  • 每个 Topic 的每个 Queue 有一个 ConsumeQueue 文件
  • 固定大小:每条索引 20 字节(8字节偏移量 + 4字节大小 + 8字节 tag hashcode)
  • 存储的是 CommitLog 中的物理偏移量
  • 消费者通过 ConsumeQueue 找到消息在 CommitLog 的位置

存储流程:

  1. Producer 发消息 → Broker
  2. Broker 写入 CommitLog(追加写)
  3. 异步构建 ConsumeQueue 索引(ReputMessageService 线程)
  4. 消费者从 ConsumeQueue 获取消息偏移量 → 从 CommitLog 读取消息

为什么这么设计:

  • 顺序写性能远高于随机写
  • CommitLog 统一写入,避免多 Topic 随机写
  • ConsumeQueue 索引较小,可以放在 OS Page Cache 中
  • 消费者按 ConsumeQueue 顺序消费,实际读取 CommitLog 可能随机,但有 Page Cache 缓存

IndexFile(可选):

  • 按 key/时间建索引,支持按 key 和时间范围查询消息
  • 哈希索引结构

刷盘策略:

  • 同步刷盘:写入内存 + 刷盘 → 返回(安全但慢)
  • 异步刷盘:写入内存 → 返回 → 后台刷盘(快但有丢数据风险)

追问延伸

  • ConsumeQueue 大小怎么估算?(每条索引20字节,百万消息约20MB)
  • 为什么不每个 Topic 一个文件?(随机写性能差)

Q24: RabbitMQ 死信队列(DLX)怎么实现?哪些场景用? 「🟡 中级」

考察点:RabbitMQ 高级特性的理解。

参考答案

死信(Dead Letter):被拒绝、过期或队列满的消息。

成为死信的三种情况:

  1. 消息被消费者拒绝(basic.reject / basic.nack)且 requeue=false
  2. 消息 TTL 过期
  3. 队列达到最大长度(x-max-length)

死信队列(DLX):死信被转发到的交换机。

配置方式:

  • 队列声明时设置 x-dead-letter-exchange:指定死信交换机
  • x-dead-letter-routing-key:指定死信路由键(默认用原消息的 routing key)

应用场景:

  1. 延迟队列:消息 TTL 过期 → 进入死信队列 → 消费者从死信队列消费(实现延迟效果)
  2. 消费失败重试:消费失败 → 拒绝且不重新入队 → 进入死信队列 → 延迟重试
  3. 异常消息归档:多次消费失败的消息存入死信队列,人工处理

延迟队列实现(TTL + DLX):

  • 正常队列设置 TTL 和 DLX
  • 消息过期后进入死信队列
  • 消费者监听死信队列,实现延迟消费
  • 缺点:TTL 是队列级别的,不能精确到每条消息(除非用插件 rabbitmq_delayed_message_exchange)
python
# 声明死信交换机和队列
channel.exchange_declare(exchange='dlx_exchange', exchange_type='direct')
channel.queue_declare('delay_queue', arguments={
    'x-message-ttl': 60000,  # 60秒过期
    'x-dead-letter-exchange': 'dlx_exchange',
    'x-dead-letter-routing-key': 'delay_key'
})

追问延伸

  • 死信队列的消息有次数限制吗?(默认没有,需要自己实现)
  • TTL 是消息级别还是队列级别?(都可,消息级 TTL 更精确但性能差)

Q25: Kafka 消费者端的重试机制怎么设计? 「🟡 中级」

考察点:消费失败处理的工程实践。

参考答案

消费者重试的挑战:

  • 直接在消费循环中同步重试 → 阻塞后续消息消费
  • 重试次数过多 → 消息积压
  • 重试不设上限 → 无限循环

重试方案:

  1. 同步重试 + 限制次数

    • 消费失败后同步重试 N 次
    • N 次后仍失败 → 跳过(记日志)或 进死信
    • 缺点:阻塞消费,重试期间不能消费新消息
  2. 异步重试 + 延迟队列

    • 消费失败 → 发到重试 Topic(如 topic_retry_30s
    • 重试 Topic 的消费者延迟 30 秒后消费
    • 再失败 → 发到 topic_retry_1min → 再失败 → topic_retry_5min
    • 最终失败 → 死信 Topic(人工处理)
    • 优点:不阻塞正常消费、分层延迟、可控
    • 缺点:需要维护多个重试 Topic
  3. Kafka 原生重试(Spring Kafka RetryTemplate)

    • Spring Kafka 内置 RetryTemplate
    • 支持配置重试次数、退避策略(指数退避)
    • 超过次数后走 DLT(Dead Letter Topic)

重试设计最佳实践:

  • 指数退避:间隔逐渐增加(1s → 2s → 4s → 8s),避免风暴
  • 最大重试次数:设上限(如 3-5 次),超过进死信
  • 死信队列:最终失败的消息进死信 Topic,人工介入
  • 区分异常类型
    • 临时异常(网络抖动)→ 重试
    • 永久异常(数据格式错)→ 不重试,直接进死信
  • 幂等保证:重试可能重复消费,消费端必须幂等
  • 监控告警:重试次数、死信数量需要监控
java
// Spring Kafka 重试配置
@Bean
public RetryTemplate retryTemplate() {
    RetryTemplate template = new RetryTemplate();
    template.setRetryPolicy(new SimpleRetryPolicy(3)); // 最多3次
    ExponentialBackOffPolicy backOff = new ExponentialBackOffPolicy();
    backOff.setInitialInterval(1000);
    backOff.setMultiplier(2.0);
    backOff.setMaxInterval(10000);
    template.setBackOffPolicy(backOff);
    return template;
}

追问延伸

  • 指数退避为什么比固定间隔好?
  • 消费失败但不重试,直接跳过有什么风险?(数据丢失)

Q26: 消息队列是参考哪种设计模式? 「🟡 中级」

考察点:设计模式在中间件中的应用。

参考答案

消息队列参考了 观察者模式发布-订阅模式,两者思路相似但有重要区别:

模式核心区别适用场景例子
观察者模式主题直接通知观察者,双方有耦合公司内部事件公司行政直接给员工发月饼
发布-订阅模式发布者和订阅者通过中间件完全解耦分布式系统通信公司通过快递公司给外部人发快递

观察者模式

主题(Subject) ──直接通知──> 观察者A
                  ──直接通知──> 观察者B
                  ──直接通知──> 观察者C
  • 观察者先订阅主题,主题维护一个观察者列表
  • 主题状态变化时,循环遍历列表逐一通知
  • 主题和观察者在同一个进程中,有直接引用关系

发布-订阅模式

发布者 ──> 发布订阅中心(Broker/MQ) ──> 订阅者A
                                     ──> 订阅者B
                                     ──> 订阅者C
  • 发布者不知道消息会被谁消费
  • 订阅者不知道消息来自谁
  • 中间件负责消息的存储、路由、分发
  • 发布者和订阅者完全解耦

MQ 为什么选择发布-订阅模式

如果公司自己去管理快递配送,就会变成一个快递公司,业务繁杂难以管理。同理,如果服务 A 直接维护对服务 B、C、D 的调用关系,耦合度极高。使用 MQ 作为中间件,发布者只需发消息,订阅者按需消费,互不影响。

追问延伸

  • Kafka 的 Consumer Group 是观察者还是发布-订阅?
  • 主题模式(Topic Exchange)和设计模式中的 Topic 有什么关系?

Q27: 让你设计一个消息队列,该如何进行架构设计? 「🔴 高级」

考察点:系统设计能力、对 MQ 原理的深入理解。

参考答案

这道题考察三个方面:MQ 架构原理理解、个人设计能力、工程思想(高可用、可扩展、幂等)。

从以下角度展开设计:

1. 整体流程

Producer ──RPC──> Broker(存储) ──RPC──> Consumer ──ACK──> Broker

核心链路:生产者发消息 → Broker 存储 → Broker 推/拉给消费者 → 消费确认

2. RPC 通信层

Producer ←→ Broker ←→ Consumer
  • 序列化协议:Protobuf / JSON
  • 通信方式:Netty(NIO)或 gRPC
  • 服务发现:ZooKeeper / NameServer / Config Server

3. 存储模型

方案A(Kafka 模式):每个 Partition 独立文件,顺序写
  → 适合 Topic 少、吞吐量大的场景
  → Topic 多时文件句柄暴涨,随机写退化

方案B(RocketMQ 模式):统一 CommitLog + ConsumeQueue
  → 所有 Topic 消息追加写一个文件
  → 后台线程构建索引(ConsumeQueue)
  → 适合 Topic 多的微服务场景

4. 消费关系管理

  • 点对点(Queue):一条消息只被一个消费者消费
  • 发布-订阅(Topic):一条消息被所有订阅者消费
  • 消费组内分区分配:轮询 / 范围 / 粘性 / 自定义

5. 可靠性保证

环节机制
生产者同步发送 + ACK 确认 + 重试
Broker持久化(刷盘策略:同步/异步) + 副本(同步/异步复制)
消费者手动 ACK + 幂等处理 + 重试 + 死信队列

6. 高可用设计

多副本 → Leader & Follower → Leader 挂了 → 选举新 Leader
  - 副本同步:ISR(In-Sync Replicas)
  - 选举算法:Raft / ZAB
  - 故障感知:心跳 + 超时

7. 事务消息

1. 发送半消息(Half Message)→ Broker 存储但不可消费
2. 执行本地事务
3. 本地事务成功 → Commit(消息可消费)
   本地事务失败 → Rollback(删除消息)
4. 超时未确认 → 回查生产者本地事务状态

8. 伸缩性与扩展性

Broker → Topic → Partition
  - 每个 Partition 放一台机器,存一部分数据
  - 资源不够 → 增加 Partition → 数据迁移 → 加机器
  - 水平扩展,线性提升吞吐量

追问延伸

  • 如何保证消息的全局有序?(单分区,牺牲并行度)
  • 如何设计消息回溯机制?(按 offset / 时间戳)

Q28: Kafka 为什么一个分区只能由消费者组的一个消费者消费?消费线程数和分区数的关系? 「🔴 高级」

考察点:Kafka 消费者模型的底层逻辑。

参考答案

为什么一个分区只能被组内一个消费者消费

核心原因是 避免消费混乱和重复处理

如果 C1 和 C2 同时消费 Partition-0:
  C1 读到 offset=2,正在处理...
  C2 读到 offset=1,还没处理完...
  C1 又读到 offset=3...
  
  → 消息处理顺序无法保证
  → 可能重复处理(C1 和 C2 读了同一条)
  → offset 管理混乱(C1 提交了 3,C2 才提交 1,回退了)

如果两个消费者负责同一个分区,由于消费者自己控制 offset,就会造成多线程读取同一个消息,重复处理且无法保证顺序。

设计意义

优点说明
顺序保证单分区单消费者,消费顺序 = 写入顺序
无锁消费不需要加锁协调多个消费者
简化 offset 管理一个分区只有一个 offset 指针

消费线程数与分区数的关系

Topic 有 10 个 Partition,Consumer Group 有 N 个消费者线程:

N < 10:某些消费者会消费多个分区(正常)
N = 10:每个消费者消费一个分区(最优吞吐)
N > 10:多出的消费者空闲(浪费资源)

结论:分区数决定了同组消费者个数的上限。

分区数 = N → 最多 N 个消费者同时消费
  → 最好线程数也保持为 N → 达到最大吞吐量
  → 超过 N 的线程不会被分配到任何分区

多线程消费同一个分区的方案

如果单分区消费速度不够,想多线程处理:

java
// 方案:消费者单线程拉取,业务多线程处理
// 注意:多线程处理后需要保证 offset 提交正确
ExecutorService executor = Executors.newFixedThreadPool(10);

consumer.subscribe(Collections.singletonList("topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    // 提交到线程池异步处理
    executor.submit(() -> {
        for (ConsumerRecord<String, String> record : records) {
            processMessage(record);
        }
    });
    // 注意:多线程处理后再提交 offset
    consumer.commitSync();
}

追问延伸

  • 不同 Consumer Group 能同时消费同一个分区吗?(可以,互不影响)
  • 如何提高单分区的消费速度?(业务多线程处理,但 offset 管理变复杂)

Q29: 消息中间件如何做到高可用? 「🟡 中级」

考察点:MQ 高可用架构的通用理解。

参考答案

无论 Kafka、RocketMQ 还是 RabbitMQ,高可用的核心就是保障两件事:机器挂了服务不能停,机器挂了数据不能丢。

三套关键机制:

1. 集群与多副本机制(打破单点)

单机 → 单点故障
集群 + 副本 → 数据有分身

Producer → Master → [异步/同步复制] → Slave
复制方式特点适用场景
异步复制Master 收到消息即返回成功,后台同步给 Slave追求性能,容忍少量丢失
同步复制Master 等 Slave 也写完磁盘才返回成功追求数据安全

2. 自动选主与故障转移(故障切换)

Master 宕机
  → 心跳检测超时
  → 剩余节点发起选举(Raft / ZAB)
  → 选出数据最全的 Slave 提升为新 Master
  → 秒级/毫秒级完成
  → 业务层几乎无感知

各 MQ 的选主方式:

MQ选主依赖机制
KafkaZooKeeper / KRaftPartition 级别 Leader 选举
RocketMQDLedger(Raft)/ NameServerBroker 级别 Master-Slave 切换
RabbitMQ集群 + 镜像队列队列级别主从同步

3. 动态路由感知(客户端感知)

Broker 主从切换
  → NameServer / ZooKeeper 感知变化
  → 下发最新路由表给 Producer / Consumer
  → 客户端自动连接新 Master
  → 继续收发消息
Kafka: 依赖 ZK/KRaft 控制器 → 通知客户端新 Leader
RocketMQ: NameServer → Broker 主动上报 → 客户端定期拉取最新路由

总结

高可用 = 多节点副本机制(数据不丢)
       + 自动选举与故障转移(服务不停)
       + 动态路由协调组件(客户端无缝切换)

追问延伸

  • Kafka 的 ISR 机制是什么?
  • RabbitMQ 镜像队列和 Quorum Queue 有什么区别?

Q30: Kafka 和 RocketMQ 消息确认机制有什么不同? 「🔴 高级」

考察点:Kafka 与 RocketMQ 的深层对比。

参考答案

一句话概括:Kafka 是「进度条模式」,RocketMQ 是「单据签收模式」。

生产者端确认对比

维度KafkaRocketMQ
核心机制基于副本的 acks基于刷盘和复制
选项acks=0(不等确认)/ acks=1(Leader 确认)/ acks=all(ISR 全确认)同步刷盘 / 异步刷盘 + 同步复制 / 异步复制
设计思路围绕分布式副本同步关注物理落盘 + 主从架构

消费者端确认对比(核心差异)

Kafka 消费者确认 = 提交偏移量(Offset Commit)
  → "我已经看到第 100 页了"
  → 批量消费,第 99 条失败、第 100 条成功
  → 提交了 100,第 99 条就被跳过(丢失)
  → 原生没有单条消息级别的失败重试机制
  → 需要业务代码自己捕获异常处理
RocketMQ 消费者确认 = 单条消息状态返回
  → CONSUME_SUCCESS(消费成功)
  → RECONSUME_LATER(稍后重试)
  → 失败后自动进入「重试队列」
  → 阶梯式延迟重试(10s → 30s → 1min → ...)
  → 重试 16 次后进入「死信队列(DLQ)」
  → 人工排查兜底
维度KafkaRocketMQ
确认粒度批量 Offset单条消息状态
失败重试无原生重试,业务自实现内置阶梯式重试
死信队列无原生支持,需自建 DLT内置 DLQ
适用场景大数据日志,吞吐优先电商业务,可靠性优先

设计哲学差异

Kafka: 为海量日志处理而生 → 追求极致吞吐量 → 确认机制粗粒度
RocketMQ: 为双 11 电商打造 → 追求严苛的业务可靠性 → 确认机制精细到单条

追问延伸

  • Kafka 消费者如何实现类似 RocketMQ 的重试机制?(重试 Topic + 延迟消费)
  • 为什么 Kafka 不原生支持单条消息重试?(影响吞吐量,大数据场景容忍少量丢失)

Q31: Kafka 和 RocketMQ 的 broker 架构有什么区别? 「🔴 高级」

考察点:MQ 底层存储架构的设计差异。

参考答案

Kafka 为大数据日志流设计,RocketMQ 为电商业务场景打造。Broker 架构最大区别在三个方面:

1. 文件存储模型(最核心区别)

Kafka:独立分区文件
  每个 Partition → 独立文件
  Topic 少 → 纯顺序写,极快
  Topic 多(上万个)→ 大量文件同时写入 → 退化为随机写 → 性能暴跌

RocketMQ:统一混合存储(CommitLog + ConsumeQueue)
  所有 Topic 的消息 → 全部追加写到一个 CommitLog 文件
  → 无论 Topic 多少,永远是顺序写
  后台线程 → 构建每个 Topic 的索引(ConsumeQueue)
  消费者 → 顺着 ConsumeQueue 去 CommitLog 找消息
维度KafkaRocketMQ
存储模型每 Partition 独立文件统一 CommitLog + ConsumeQueue
Topic 多时性能严重下滑性能稳定
写入方式Partition 级顺序写全局顺序写
读取方式直接读 Partition 文件ConsumeQueue 索引 → CommitLog
适用场景少 Topic 大吞吐多 Topic 微服务

2. 协调机制(大脑不同)

Kafka: 依赖外部组件
  → 传统架构依赖 ZooKeeper
  → 新版换成内置 KRaft 控制器
  → Broker 是纯打工人,元数据管理交给 ZK/KRaft
  → 节点多时心跳开销大

RocketMQ: 轻量级 NameServer
  → 极简设计,NameServer 节点间互不通信
  → Broker 主动向所有 NameServer 上报
  → 挂掉一个 NameServer 不影响大局
  → 架构简单稳定

3. 高可用粒度

Kafka: Partition 级别高可用
  → 一个 Broker 宕机 → 成百上千个 Partition 重新选举 Leader
  → 选举开销与 Partition 数成正比

RocketMQ: Broker 级别高可用
  → Master-Slave 模式
  → Master 挂了 → 消费者无缝切换读 Slave
  → 管理层级更高,更像数据库主从
维度KafkaRocketMQ
高可用粒度Partition 级Broker 级
故障恢复逐 Partition 选举Master-Slave 切换
协调组件ZooKeeper / KRaftNameServer

总结

Kafka Broker = 分布式文件系统
  → 每个 Partition 自治
  → 适合大数据场景,少 Topic 大吞吐

RocketMQ Broker = 数据库引擎思维
  → CommitLog 统一存储
  → NameServer 极致轻量
  → 适合企业级微服务,多 Topic 业务

追问延伸

  • Kafka 为什么在 Topic 多的时候性能会下降?
  • RocketMQ 的 ConsumeQueue 是如何保证实时性的?

Q32: RabbitMQ 的核心组件有哪些?AMQP 是什么关系? 「🟡 中级」

考察点:RabbitMQ 基础概念理解。

参考答案

AMQP 与 RabbitMQ 的关系

AMQP = 高级消息队列协议(标准/规范)
RabbitMQ = AMQP 协议的具体实现(软件产品)

关系类比:Interface(AMQP) ↔ Implementation(RabbitMQ)
  • AMQP 是图纸:定义了消息队列系统应该长什么样,比如消息不能直接进队列,中间要有 Exchange
  • RabbitMQ 是真房子:用 Erlang 语言照着 AMQP 协议标准,写出了可以运行的软件产品
  • AMQP 赋予 RabbitMQ 最大王牌:强大的路由能力(Exchange + RoutingKey + Binding)

RabbitMQ 核心组件

Producer → Exchange → (RoutingKey + Binding) → Queue → Consumer

                                    Connection / Channel
组件作用说明
Producer(生产者)发送消息不直接发队列,发给 Exchange
Exchange(交换机)消息分发器根据规则把消息路由到队列
RoutingKey(路由键)路由标识生产者发送时携带
Binding(绑定)绑定关系Exchange 和 Queue 之间的关联规则
Queue(队列)存储消息FIFO,真正存消息的地方
Consumer(消费者)消费消息从 Queue 取消息执行业务
ConnectionTCP 长连接客户端与 RabbitMQ 的物理连接
Channel(信道)轻量级通道在 Connection 上创建,避免频繁创建 TCP

为什么需要 Channel

没有 Channel:
  每个生产者/消费者 → 一个 TCP 连接
  → 大量 TCP 连接 → 创建销毁开销巨大

有 Channel:
  一个 TCP Connection → 多个 Channel(虚拟连接)
  → 复用 TCP 连接 → 大幅节省资源

四种交换机类型

类型路由规则适用场景
Direct路由键完全匹配精准投递、一对一
Fanout不看路由键,广播群发通知、多服务同步
Topic路由键模糊匹配(* / #按业务主题分类路由
Headers消息头键值对匹配几乎不用,性能差

追问延伸

  • 为什么 RabbitMQ 用 Erlang 语言开发?(Erlang 天生适合并发和分布式)
  • AMQP 协议和 MQTT、STOMP 有什么区别?

Q33: RabbitMQ 的可靠性保障怎么做? 「🟡 中级」

考察点:RabbitMQ 端到端可靠性设计。

参考答案

RabbitMQ 可靠性保障核心是「消息不丢失、不重复、不积压」,从三个环节层层把关:

Producer → Exchange → Queue → Consumer
  ①          ②         ③

1. 生产者环节(保证消息投递到 Exchange)

机制说明
Publisher Confirm开启确认模式,消息到达 Exchange 后返回 ACK
Return 回调消息无法路由到队列时触发(路由失败通知)
事务机制channel.txSelect()txCommit(),但影响吞吐量
重试机制设置合理重试次数和间隔,处理网络异常
java
// Publisher Confirm 模式
channel.confirmSelect();
channel.basicPublish(EXCHANGE, ROUTING_KEY, props, body);
if (channel.waitForConfirms(5000)) {
    // 消息成功到达 Exchange
} else {
    // 重试
}

2. RabbitMQ 服务器环节(防止服务端丢失)

机制说明
队列持久化durable=true,重启后队列不消失
消息持久化deliveryMode=2,消息写入磁盘
交换机持久化durable=true,重启后绑定关系不丢失
镜像队列队列内容复制到多个节点,防单点故障
死信交换机无法路由的消息转发到死信队列,人工处理
java
// 声明持久化队列
channel.queueDeclare("durable_queue", true, false, false, null);

// 发送持久化消息
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
    .deliveryMode(2)  // 持久化
    .build();
channel.basicPublish("", "durable_queue", props, body);

3. 消费者环节(保证消息正确处理)

机制说明
手动 ACK处理成功后手动发送 ack,失败不发送(消息重新入队)
预取限制basicQos(N),每次只取 N 条,处理完再取,防积压
幂等设计唯一消息 ID + 去重表 / Redis
消费失败处理重试 → 死信队列 → 人工处理
java
// 手动 ACK + 预取
channel.basicQos(10);  // 每次最多 10 条未确认
channel.basicConsume("queue", false, (consumerTag, message) -> {
    try {
        processMessage(message);
        channel.basicAck(message.getEnvelope().getDeliveryTag(), false);
    } catch (Exception e) {
        // 处理失败,不 ACK,消息重新入队
        channel.basicNack(message.getEnvelope().getDeliveryTag(), false, true);
    }
});

完整可靠性链路

生产者 → Confirm + 重试
  → Exchange → 持久化 + 路由失败回调
    → Queue → 持久化 + 镜像
      → Consumer → 手动 ACK + 幂等 + 预取控制

追问延伸

  • Publisher Confirm 和事务机制有什么区别?为什么不推荐用事务?
  • 镜像队列和 Quorum Queue 有什么区别?

Q34: 使用消息队列还应该注意哪些问题? 「🟡 中级」

考察点:MQ 使用中的工程实践和踩坑经验。

参考答案

除了消息不丢失、不重复、有序性这三大经典问题,实际使用 MQ 还需要注意以下问题:

1. 消息大小控制

问题:单条消息过大 → 网络传输慢、Broker 内存压力、消费端 OOM
建议:单条消息不超过 256KB
方案:大消息拆分 / 引用传递(消息只存 ID,数据走文件/DB)

2. 消费速率与生产速率匹配

问题:生产远大于消费 → 消息积压 → 磁盘满 → Broker 宕机
监控:积压量告警(如 > 10 万条)
方案:
  - 动态扩容消费者(注意分区数限制)
  - 降级:非核心消息丢弃或延迟消费
  - 限流:生产端限速

3. 消息过期策略

问题:消息积压后,旧消息已无意义(如库存已超时释放)
方案:
  - Kafka:基于时间/大小的日志清理(Log Retention)
  - RabbitMQ:TTL + 死信队列
  - RocketMQ:延迟级别 + 消息过期

4. 消息回溯能力

问题:线上 bug 修复后,如何重新消费历史消息?
方案:
  - Kafka:重置 offset 到历史位置(按时间/offset)
  - RocketMQ:按时间回溯消息
  - 通用:消息表持久化 + 重新投递

5. 集群容量规划

需要考虑:
  - 单 Broker 磁盘容量
  - 单 Topic 分区数 vs 消费者数
  - 峰值吞吐量(生产 + 消费)
  - 消息保留时间 × 消息速率 = 存储需求

6. 监控告警体系

监控指标告警阈值说明
消息积压量> 10 万条消费能力不足
消费延迟> 5 分钟消费速度跟不上
生产失败率> 1%Broker 异常
消费失败率> 5%业务异常
磁盘使用率> 80%需要扩容或清理
Broker 存活任意节点宕机高可用告警

7. 灰度发布与消费组隔离

问题:新版本消费者上线,如果逻辑有 bug 影响全量消息
方案:
  - 灰度消费:新建消费组,先消费部分流量
  - 消费组隔离:不同环境(dev/staging/prod)用不同消费组
  - 版本兼容:新旧消费者并行期间,消息格式需向后兼容

8. 消息追踪与可观测性

方案:
  - Trace ID 贯穿全链路(生产 → Broker → 消费)
  - RocketMQ:消息轨迹(Message Trace)
  - Kafka:Interceptor + 日志
  - 通用:ELK / Jaeger / SkyWalking

追问延伸

  • 如何评估一个 MQ 集群的容量瓶颈?
  • 消费者灰度发布时如何保证不丢消息?