MQ 面试
MQ 面试
MQ 简介
【简单】MQ 是什么?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:MQ / 基础概念
💎 关键结论
MQ 是一种「存储并转发消息」的异步通信中间件,把生产者与消费者解耦。理由:把同步等待变成异步缓冲,双方按各自节奏处理,系统吞吐、弹性与可维护性同时提升。
⚡记忆卡片
口诀:存储转发,异步解耦,缓冲背压
关键词:生产者/消费者/消息/队列/Broker/背压
链路:生产者发送消息 → Broker 接收并存储 → 消费者按需拉取处理 → 消费慢时背压反制生产者降速
📖 核心知识
- 定义:MQ(Message Queue,消息队列) 是一种异步通信机制,用于在不同服务、应用或系统组件之间可靠地传递消息。它的核心思想是解耦生产者和消费者,通过缓冲消息来提高系统的可靠性、扩展性和可维护性。
- 核心概念:
- 生产者(Producer):发送消息的应用或服务。
- 消费者(Consumer):接收并处理消息的应用或服务。
- 消息(Message):传输的数据单位,可以是文本、JSON、二进制等格式。
- 队列(Queue):存储消息的缓冲区,遵循 FIFO(先进先出) 或优先级策略。
- Broker(消息代理):负责接收、存储和转发消息的中间件(如 RabbitMQ、Kafka)。
- 背压:消息传递机制中的流控策略。当消费者处理速度跟不上生产者发送速度时,通过反向施加压力,迫使生产者降速或停止发送,以防止系统被压垮。
- 主流产品:目前主流的 MQ 有 Kafka、RabbitMQ、RocketMQ、ActiveMQ。
🔬 扩展知识
【L3】MQ 与 RPC 的分工
详情
- RPC 面向同步调用,要求双方同时在线、即时响应;MQ 面向异步消息,双方无需同时在线、无需相互感知。
- 架构设计时的分界线:强时效、强依赖的调用走 RPC;可容忍延迟、需要缓冲或广播的交互走 MQ。
【L4】MQ 形态的演进
详情
- 新一代消息系统(Kafka KRaft、RocketMQ 5.0、Pulsar)普遍走向存算分离,Broker 层可弹性伸缩,存储层独立扩展。
- 消息队列正从「传输通道」演进为「事件流底座」:不仅投递消息,还支持消息重放、回溯与流式计算(如 Kafka + Flink)。
🔀 发散问题
MQ 为什么能提高系统吞吐量?
因为同步等待被异步化:生产者写入 Broker 后即可返回,无需等待下游处理完成;消费者按自身能力批量、并行消费。叠加削峰缓冲,整体并发吞吐显著提升,典型场景见本文档『MQ 有哪些应用场景?』。
引入了 MQ,系统就一定更可靠吗?
不一定。MQ 本身成为新的依赖点,若 MQ 宕机则通信中断;必须为 MQ 自身做高可用设计(集群、副本),并承担重复消费、消息丢失、顺序性等新增复杂度,详见本文档『MQ 存在哪些挑战?』。
【简单】MQ 有哪些应用场景?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 应用场景
💎 关键结论
MQ 的典型场景可归纳为六个词:异步、解耦、削峰、通信、缓冲、一致性。理由:它们都是「存储转发 + 异步化」这两个基本能力在不同业务问题上的投影。
⚡记忆卡片
口诀:异步解耦削峰,通信缓冲一致
关键词:异步处理/系统解耦/流量削峰/系统间通信/传输缓冲/最终一致性
链路:生产者发消息即返回 → MQ 缓冲堆积 → 消费者按能力消费 → 上下游解耦、峰值被削平、最终一致可达成
📖 核心知识
MQ 的典型应用场景:异步处理、系统解耦、流量削峰、系统间通信、传输缓冲、最终一致性。
异步处理
MQ 可以将系统间的处理流程异步化,减少等待响应的时间,从而提高整体并发吞吐量。一般,MQ 异步处理应用于非核心流程,例如:短信/邮件通知、数据推送、上报数据到监控中心、日志中心等。
假设这样一个场景,用户向系统 A 发起请求,系统 A 处理计算只需要 10ms,然后通知系统 BCD 写库,系统 BCD 写库耗时分别为:100ms、200ms、300ms。最终总耗时为:10ms+100ms+200ms+300ms=610ms。此外,加上请求和响应的网络传输时间,从用户角度看,可能要等待将近 1s 才能得到结果。

如果使用 MQ,系统 A 接到请求后,耗时 10ms 处理计算,然后向系统 BCD 连续发送消息,假设耗时 5ms。那么这一过程的总耗时为 10ms + 5ms = 15ms,这相比于 610ms,大大缩短了响应时间。至于系统 BCD 的写库操作,只要自行消费 MQ 后处理即可,用户无需关注。

系统解耦
通过 MQ,可以消除系统间的强耦合。它的好处在于:
- 消息的消费者系统可以随意增加,无需修改生产者系统的代码。
- 生产者系统、消费者系统彼此不会影响对方的流程。
- 如果生产者系统宕机,消费者系统收不到消息,就不会有下一步的动作。
- 如果消费者系统宕机,生产者系统仍然可以正常发送消息,不影响流程。
不同系统如果要建立通信,传统的做法是:调用接口。如果需要和新的系统建立通信或删除已建立的通信,都需要修改代码,这种方案显然耦合度很高。

如果使用 MQ,系统间的通信只需要通过发布/订阅(Pub/Sub)模型即可,彼此没有直接联系,也就不需要相互感知,从而达到 解耦。

流量削峰
当 上下游系统 处理能力存在差距的时候,利用 MQ 做一个 “漏斗” 模型,进行 流控。把 MQ 当成可靠的 消息缓冲池,进行一定程度的 消息堆积;在下游有能力处理的时候,再发送消息。
MQ 的流量削峰常用于高并发场景(例如:秒杀、团抢等业务场景),它是缓解瞬时暴增流量的核心手段之一。
如果没有 MQ,两个系统之间通过 协商、滑动窗口、限流/降级/熔断 等复杂的方案也能实现 流控。但 系统复杂性 指数级增长,势必在上游或者下游做存储,并且要处理 定时、拥塞 等一系列问题。而且每当有 处理能力有差距 的时候,都需要 单独 开发一套逻辑来维护这套逻辑。
假设某个系统读写数据库的稳定性能为每秒处理 1000 条数据。平常情况下,远远达不到这么大的处理量。假设,因为做活动,系统的瞬时请求量剧增,达到每秒 10000 个并发请求,数据库根本承受不了,可能直接就把数据库给整崩溃了,这样系统服务就不可用了。

如果使用 MQ,每秒写入 10000 条请求,但是系统 A 每秒只从 MQ 中消费 1000 条请求,然后写入数据库。这样,就不会超过数据库的承受能力,而是把请求积压在 MQ 中。只要高峰期一过,系统 A 就会很快把积压的消息给处理掉。

- 系统间通信:消息队列一般都内置了 高效的通信机制,因此也可以用于单纯的 消息通讯,比如实现 点对点消息队列 或者 聊天室 等。生产者/消费者 模式,只需要关心消息是否 送达队列,至于谁希望订阅和需要消费,是 下游 的事情,无疑极大地减少了开发和联调的工作量。
传输缓冲
(1)MQ 常被用于做海量数据的传输缓冲。
例如,Kafka 常被用于做为各种日志数据、采集数据的数据中转。然后,Kafka 将数据转发给 Logstash、Elasticsearch 中,然后基于 Elasticsearch 来做日志中心,提供检索、聚合、分析日志的能力。开发者可以通过 Kibana 集成 Elasticsearch 数据进行可视化展示,或自行进行定制化开发。

(2)MQ 也可以被用于流式处理。
例如,Kafka 几乎已经是流计算的数据采集端的标准组件。而流计算通过实时数据处理能力,提供了更为快捷的聚合计算能力,被大量应用于链路监控、实时监控、实时数仓、实时大屏、风控、推荐等应用领域。
- 最终一致性:最终一致性 不是 消息队列 的必备特性,但确实可以依靠 消息队列 来做 最终一致性 的事情。
- 先写消息再操作,确保操作完成后再修改消息状态。定时任务补偿机制 实现消息 可靠发送接收、业务操作的可靠执行,要注意 消息重复 与 幂等设计。
- 所有不保证
100%不丢消息 的消息队列,理论上无法实现 最终一致性。像 Kafka 一类的设计,在设计层面上就有 丢消息 的可能(比如 定时刷盘,如果掉电就会丢消息)。哪怕只丢千分之一的消息,业务也必须用其他的手段来保证结果正确。
🔬 扩展知识
【L3】场景与基本能力的对应关系
详情
- 异步处理、传输缓冲 → 利用的是 MQ 的「异步 + 堆积」能力。
- 系统解耦、系统间通信 → 利用的是 MQ 的「发布/订阅、双方互不感知」能力。
- 流量削峰 → 利用的是 MQ 的「缓冲池 + 按能力消费」能力。
- 最终一致性 → 利用的是 MQ 的「可靠投递 + 补偿重试」能力,但必须以不丢消息为前提。
【L4】MQ 在数据架构中的角色
详情
- 在大数据体系中,MQ(尤其 Kafka)承担「数据总线」角色:上游各业务系统写入,下游数仓、搜索、流计算各自订阅,一份数据多处消费。
- 这使得 MQ 的容量规划从「消息条数」变为「数据流量 × 保留时长」,存储成本成为核心约束。
📚 延伸阅读:Kafka 官方文档
🔀 发散问题
削峰会不会变成「削不掉」?
会。削峰的前提是「峰值可堆积、峰谷能消化」,若生产速率长期大于消费速率,积压只会无限增长,最终打爆 MQ 存储。此时必须对生产端限流,参见本文档『如何处理 MQ 消息积压?』。
用 MQ 做异步化后,用户体验如何保障?
核心流程改为「提交即受理」:用户请求落 MQ 后立即返回受理凭证,处理结果通过轮询、推送或回调通知用户;同时配合消息堆积监控,避免用户长时间等不到结果。
【中等】MQ 存在哪些挑战?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:MQ / 架构权衡
💎 关键结论
引入 MQ 等于引入新的依赖:可用性降低、复杂度提高、一致性问题浮现。理由:通信链路全部经过 MQ,MQ 自身的任何故障与特性缺陷都会放大为系统级问题。
⚡记忆卡片
口诀:可用性降,复杂度升,一致性难
关键词:可用性/复杂度/重复消费/消息丢失/顺序性/消息积压/一致性
链路:引入 MQ → 通信依赖 MQ → MQ 宕机则全链路不可用 → 还需额外解决丢失、重复、乱序、积压、一致性
📖 核心知识
任何技术都会有利有弊,MQ 给整体系统架构带来很多好处,但也会付出一定的代价。MQ 主要引入了以下问题:
- 系统可用性降低:引入了 MQ 后,通信需要基于 MQ 完成,如果 MQ 宕机,则服务不可用。因此,MQ 要保证是高可用的。
- 系统复杂度提高:使用 MQ,需要关注一些新的问题:
- 如何保证消息没有 重复消费?
- 如何处理 消息丢失 的问题?
- 如何保证传递 消息的顺序性?
- 如何处理大量 消息积压 的问题?
- 一致性问题:假设系统 A 处理完直接返回成功的结果给用户,用户认为请求成功。但如果此时,系统 BCD 中只要有任意一个写库失败,那么数据就不一致了。这种情况如何处理?
MQ 存储
【中等】Kafka、RocketMQ、RabbitMQ 如何存储数据?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 存储模型
💎 关键结论
三者存储模型各异:Kafka 是「分区日志流」,RocketMQ 是「CommitLog 混存 + ConsumeQueue 索引」,RabbitMQ 是「队列文件 + 确认即删」。理由:吞吐优先、业务功能优先、路由灵活优先三种设计目标,决定了不同的存储形态。
⚡记忆卡片
口诀:Kafka 日志流,Rocket 读写分,Rabbit 队列删
关键词:Partition/CommitLog/ConsumeQueue/Exchange/Queue/Offset
链路:消息写入分片/日志文件 → 建立偏移量或逻辑索引 → 消费者按 Offset 拉取或 Broker 推送 → 按保留策略清理或确认后删除
📖 核心知识
逻辑与物理存储模型对比
| 特性 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 核心逻辑模型 | 发布/订阅的日志流 | 发布/订阅的消息队列 | 队列与交换器的路由代理 |
| 逻辑存储 | 主题(Topic) 分区(Partition) 记录(Record) | 主题(Topic) 队列(MessageQueue) 消息(Message) | 交换机(Exchange) 队列(Queue) |
| 物理存储 | Log(对应 Partition) LogSegment | 提交日志(CommitLog) 消费队列(ConsumeQueue) | Queue 对应的消息存储文件 |
| 存储文件 | .log (数据文件).index (偏移量索引).timeindex (时间索引) | CommitLog (提交日志)ConsumeQueue (消费队列)IndexFile (哈希索引) | .rdq (消息体文件).idx (消息索引与元数据文件) |
| 持久化方式 | 追加写入 + 零拷贝 | 追加写入 + 内存映射 | 可变长度写入,支持内存/磁盘持久化 |
| 刷盘方式 | 数据先写入页缓存,再在合适时机写入磁盘 | 数据先写入页缓存,再在合适时机写入磁盘 | |
| 数据清理 | 基于时间或空间的保留策略 | 基于时间或空间的保留策略 | 消费确认后立即删除 |
| 消费组模型 | 消费者组(Consumer Group) | 消费者组(Consumer Group) | 无明确组概念,客户端竞争消费同一队列 |
| 消息获取模式 | 消费者按 Offset 主动拉取 | 消费者主要按 Offset 主动拉取 | Broker 主动推送 给消费者(主流),也支持拉取 |
🔬 扩展知识
【L3】RocketMQ 为什么把所有 Topic 混写进一个 CommitLog
详情
- 所有 Topic 的消息顺序写入同一个 CommitLog,保证任何 Topic 的写入都是顺序 I/O,避免多 Topic 各自独立文件导致的随机写。
- 代价是读取需要经 ConsumeQueue 二次寻址,RocketMQ 用 mmap + PageCache 缓解读放大;Kafka 则按 Partition 建文件,读路径更直接,但 Partition 过多会产生随机写。
【L4】清理策略对架构的影响
详情
- Kafka/RocketMQ 按保留策略清理,天然支持消息重放与回溯;RabbitMQ 确认即删,重放能力弱。
- 这决定了二者定位差异:前者可做数据管道与事件溯源,后者专注任务分发与即时投递。
【L3】Kafka 定制深挖:LogSegment 切分细节
详情


- 逻辑三级 Topic → Partition → Record,物理上每个 Partition 对应一个
Log对象、一个目录(如test-0/test-1);Partition 是副本复制与 Offset 编号的最小单元,无法再细分到多磁盘或多 Broker。 - Log 切成 LogSegment(默认每段 ≤1GB 且只含 7 天数据),段文件名用「上一段最后一条消息的 offset」,定位目标段只需比较文件名二分,检索成本从 O(消息数) 降到 O(段数)。
- 正在写入的活跃片段(active segment)永不删除;每段占用
.log+.index+.timeindex(用事务时还有.txnindex,记录已终止事务偏移量供read_committed过滤)文件句柄,分区多、段数多时需调大系统ulimit -n。 - 想要更高并行度,正确姿势是增加分区数并配置多个
log.dirs让分区分布在不同磁盘,而不是拆分单个分区。
【L3】RocketMQ 定制深挖:CommitLog 目录结构与读写流程
详情

- 存储根目录由
storePathRootDir决定:commitLog 文件夹存消息物理文件,consumeQueue 文件夹存逻辑队列索引。 - 写入流程:所有 Topic 消息混写顺序写入 CommitLog 后立即返回;索引由单线程
ReputMessageService异步分发构建,正常负载下与写入延迟在毫秒级;该线程积压会出现“写入成功但暂时消费不到”,dispatchBehindBytes是核心监控指标。 - 读取流程:消费者先按 offset 查 ConsumeQueue 索引项(每条固定 20 字节:8 字节 CommitLog 偏移 + 4 字节消息大小 + 8 字节 Tag 哈希),再按物理偏移回查 CommitLog 取消息体——二次寻址是混存设计的读取代价。
📚 延伸阅读:RocketMQ 官方文档
🔀 发散问题
为什么 RabbitMQ 不适合海量消息堆积?
因为它按「确认即删」设计,堆积的消息会持续占用磁盘且索引结构对大量未确认消息不友好,性能随堆积量显著下降。海量堆积场景应选择 Kafka 或 RocketMQ,详见本文档『Kafka、RocketMQ、RabbitMQ 有什么区别?如何选型?』。
Kafka 的 Offset 存在哪里?
Offset 是 Partition 内每条记录的连续编号,Broker 端持久化在 __consumer_offsets 主题(新版本),客户端消费后定期提交已消费位点;这也是重复消费问题的根源,参见本文档『如何保证 MQ 消息不重复?』。
【中等】Kafka、RocketMQ、RabbitMQ 如何持久化?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 持久化
💎 关键结论
持久化设计由定位决定:Kafka「日志即存储」追求极致顺序写,RocketMQ「读写分离」兼顾低延迟读与事务,RabbitMQ「队列即存储」换取路由灵活与确认即删。理由:写入路径、刷盘策略、清理策略三者的组合,本质是吞吐、延迟、可靠性之间的取舍。
⚡记忆卡片
口诀:Kafka 拼吞吐,Rocket 顾均衡,Rabbit 重路由
关键词:顺序追加/页缓存/同步刷盘/异步刷盘/保留策略/稀疏索引
链路:消息顺序追加写入 → 先落页缓存再择机刷盘 → 按时间/空间保留策略清理(或确认即删)→ 索引支撑检索
📖 核心知识
核心要点对比
| 机制 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 核心设计 | 日志即存储 | 读写分离 | 队列即存储 |
| 写入方式 | 顺序追加写入 | 顺序追加写入 (CommitLog) | 可变长度写入 |
| 刷盘机制 | 异步(高性能) 同步(高可靠) | 异步 (高性能) 同步 (高可靠) | 异步 确认后删除 |
| 数据清理 | 基于保留策略 时间 大小 压缩 (Key 去重) | 基于保留策略 时间 大小 | 消费后删除 确认后立即删除 队列长度限制 |
| 索引机制 | 稀疏索引.index(偏移量).timeindex(时间戳) | 多级索引ConsumeQueue(核心逻辑索引)IndexFile(关键字查询) | 队列索引 消息存储索引 内存状态记录(性能瓶颈) |
总结
- Kafka:为超高吞吐和海量数据堆积设计,采用最极致的顺序 I/O,像一个只增不删的日志系统。
- RocketMQ:在 Kafka 的基础上做了优化,通过读写分离的设计,在保证高吞吐写入的同时,兼顾了低延迟读取和事务消息等金融级需求。
- RabbitMQ:核心在于灵活的路由和可靠交付,其存储设计不适合海量消息堆积,但能实现复杂的消息路由和“确认即删除”的精准控制。
🔬 扩展知识
【L3】同步刷盘与异步刷盘的代价
详情
- 异步刷盘:消息写入页缓存即返回,性能最高,但掉电可能丢失尚未刷盘的消息。
- 同步刷盘:每条消息落盘后才返回确认,可靠性高,但单次写延迟从亚毫秒级升至毫秒级,吞吐显著下降,通常只对资金类 Topic 开启。
- 工程上常用「异步刷盘 + 多副本」替代「同步刷盘」,用副本冗余对冲单机掉电风险。
【L4】Kafka 稀疏索引的设计意图
详情
- Kafka 不索引每条消息,仅每隔一定字节数记录一个偏移量索引项,索引文件小、可常驻内存。
- 检索时先二分定位索引项,再在段文件内顺序扫描,用少量内存换取顺序读的高吞吐。
【L3】Kafka 定制深挖:不强制 fsync,持久性押注副本
详情
- Kafka 把持久性押在副本而非单机刷盘上:
acks=all+min.insync.replicas≥2时,同一消息同时存在于多台机器的页缓存,多机同时损毁概率极低。 - 若强制
log.flush.interval.messages=1每条 fsync,吞吐会从百万级跌到万级;单机强持久只能接受这个代价。 - 0.11+ 的 Leader Epoch 机制修复了旧版基于 HW 截断可能丢已提交消息的问题;配合跨机架部署(
broker.rack)可兜底机柜级故障。
【L3】RocketMQ 定制深挖:flushDiskType 与 TransientStorePool
详情
- 两档刷盘由
flushDiskType配置:ASYNC_FLUSH性能高但有极低概率丢消息;SYNC_FLUSH可靠性高但性能有损耗。 - 清理策略:磁盘使用率超阈值(默认 75%)强制清理过期文件;超过保留时长(默认 72 小时)且不再被消费者需要的文件定期清理。
transientStorePoolEnable=true启用堆外内存池(TransientStorePool),写入先进堆外内存再异步刷盘,减少 JVM GC 影响并提升写入吞吐,代价是 Broker 宕机丢失窗口变大(仅异步刷盘 + 主从架构下可用)。- RocketMQ 依赖 OS Page Cache 加速读写,部署时应预留充足空闲内存并调低 swappiness,避免页缓存被换出。
【L3】RabbitMQ 定制深挖:durable 三件套与刷盘时机
详情
- 三件套缺一不可:队列
durable=true、消息deliveryMode=2、交换机durable=true;仅消息持久化而队列非持久化时,重启后队列不存在,消息照样丢。 - 持久化消息以批量/短窗口方式写盘而非每条同步 fsync,“持久化”承诺的是重启后可恢复而非写入瞬间绝对安全——单节点在写盘窗口内崩溃仍可能丢少量消息;要崩溃场景也不丢,需仲裁队列的多数派落盘确认。
- 仲裁队列(3.8+)消息天然全部持久化并 Raft 多副本落盘,三件套更多是经典队列时代的遗留知识;且持久化(单节点重启不丢)与高可用(节点故障不丢)正交,完整可靠性方案是持久化 + 副本 + 确认机制的组合。
🔀 发散问题
为什么三者都先写页缓存而不是直接写磁盘?
页缓存把随机、细碎的写操作合并为 OS 调度的批量刷盘,且读路径可命中缓存,是顺序 I/O 高吞吐的关键;直接写盘(O_DIRECT)反而失去 OS 缓存优化。
消息保留期设置过短会有什么后果?
消费者一旦故障时间超过保留期,消息被清理后连重放的机会都没有,只能靠外部归档或数据库反查补数,参见本文档『MQ 如何实现消息重放(重新消费)?』。
MQ 生产消费
【中等】MQ 数据传输有哪些模式?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 生产消费
💎 关键结论
MQ 数据传输分 Push 与 Pull 两种,长轮询是二者的折中:低延迟选 Push,高吞吐选 Pull,平衡选长轮询。理由:推送实时但易过载,拉取可控但有轮询开销,工程上靠背压与长轮询各补短板。
⚡记忆卡片
口诀:推快拉稳,长轮折中
关键词:Push/Pull/长轮询/背压/prefetch_count
链路:Broker 有消息 → Push 主动送达(实时但需背压防过载)/ Pull 按需拉取(可控但有空轮询)→ 长轮询挂起请求直到有消息或超时
📖 核心知识
MQ 数据传输的模式主要分为 Push(推) 和 Pull(拉) 两种,不同消息队列中间件采用不同的策略,部分系统还支持 混合模式。

消息队列消费模式对比:
| 维度 | Push | Pull | 长轮询 | 混合模式 |
|---|---|---|---|---|
| 特点 | Broker 主动推送给消费者,消费者被动接收;实时性高,减少消费者轮询开销 | 消费者主动从 Broker 拉取消息,按需获取;消费者控制消费速率 | Push 和 Pull 的折中方案;请求被保持,直到有消息或超时 | 关键消息用 Push,批量数据用 Pull;消费者可动态切换模式 |
| 优点 | 低延迟,适合实时性要求高的场景(如即时通讯) | 避免消息堆积冲击消费者,适合高吞吐场景(如日志处理) | 减少无效轮询,平衡实时性与服务端压力 | 灵活性强,能兼顾实时性和吞吐量,适应复杂场景 |
| 缺点 | 可能造成消费者过载(需背压机制控制流速) | 存在空轮询开销(可通过长轮询优化) | 实现复杂度高于纯 Pull,需维护挂起的请求连接 | 系统设计和实现的复杂度最高 |
| 典型实现 | RabbitMQ、ActiveMQ、RocketMQ(默认长轮询模拟 Push) | Kafka、Pulsar(原生 Pull)、RocketMQ(支持显式 Pull) | RocketMQ(默认模式)、HTTP 长轮询(如 WebSocket) | Pulsar(支持多模式)、部分自研 MQ 系统 |
选择建议
- 需要低延迟 → Push(如 RabbitMQ)
- 需要高吞吐 → Pull(如 Kafka)
- 平衡场景 → 长轮询(如 RocketMQ)
Push 模式的深度分析
Push 模式的核心问题:消费者过载
Push 模式下,Broker 主动推送消息,若消费者处理速度跟不上,会导致:
- 消息在消费者端积压,内存溢出
- 消费者因过载而崩溃
解决方案:背压机制
不同 MQ 实现背压的方式:
| MQ | 背压实现方式 |
|---|---|
| RabbitMQ | prefetch_count 限制未确认消息数,达到上限停止推送 |
| RocketMQ | 长轮询 + 消费者拉取频率控制 |
| Kafka | 消费者主动 Pull,天然具备背压能力 |
Pull 模式的深度分析
Pull 模式的核心问题:实时性差
Pull 模式下,消费者主动拉取,若队列无消息,会产生空轮询,浪费资源。
解决方案:
- 长轮询:Broker 在无消息时挂起请求,有消息或超时后才返回(如 RocketMQ 默认 15s 长轮询)。
- 退避策略:空轮询后增加下次拉取间隔,避免频繁无效请求。
- Kafka 的优化:消费者维护
Fetch Request,Broker 在有数据时才响应。
实践建议
- 实时性要求高的业务(如订单通知):选择 Push 模式或短轮询间隔。
- 吞吐量优先的业务(如日志处理):选择 Pull 模式,配合批量拉取。
- 混合场景:根据消息优先级动态调整策略。
🔬 扩展知识
【L3】RocketMQ Push 消费的本质
详情
- RocketMQ 的 PushConsumer 并非 Broker 真推送,而是客户端内部封装的「长轮询拉取」:Pull 请求在无消息时被 Broker 挂起(默认最长 15s),一有消息立即返回,对用户呈现为推送语义。
- 这种设计同时获得 Push 的低延迟与 Pull 的流速可控。
【L4】Pull 模式的消费位点语义
详情
- Pull 模式下消费进度(Offset)由客户端掌控并提交,天然支持按位点回溯与重放;Push 模式的确认语义(ACK/NACK)则更贴近「单条消息状态机」。
- 因此重放、回刷类能力在 Pull 体系(Kafka/RocketMQ)中是一等公民。
🔀 发散问题
为什么 Kafka 坚持 Pull 而不是 Push?
消费者速率差异巨大,Push 容易把慢消费者打垮;Pull 让消费者按自身能力取数,配合批量拉取天然具备背压能力,更适合高吞吐与重放场景。
长轮询会不会占用大量连接资源?
会挂起请求占用连接,但相比短轮询大幅减少了无效请求;Broker 端以超时兜底(如 RocketMQ 15s),连接数与消费者数同阶,工程上可接受。
【中等】MQ 有哪些通信模型?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:MQ / 通信模型
💎 关键结论
MQ 有两种核心通信模型:点对点(一条消息一个消费者)与发布/订阅(一条消息多个订阅者);现代 MQ 用消费者组把二者融合。理由:组间广播实现订阅,组内独占消费实现点对点,一套机制覆盖两类需求。
⚡记忆卡片
口诀:点对点独占,订阅者共享,组来融合
关键词:点对点/发布订阅/消费者组/任务分发/事件广播
链路:消息入队 → 点对点模型下仅一个消费者竞争成功 → 发布订阅模型下所有订阅者各收一份 → 消费者组实现组间订阅、组内独占
📖 核心知识
消息队列(MQ)的两种核心通信模型:
- 点对点模型
- 核心特点:消息存储在队列中。一条消息只能被一个消费者成功消费和处理,存在消费竞争关系。
- 应用场景:适用于任务分发和工作队列,例如视频转码、报表生成等需要并行处理任务的场景。
- 发布/订阅模型
- 核心特点:消息被发布到主题中。一条消息会被广播给所有订阅了该主题的消费者,消息被所有订阅者共享,没有竞争关系。
- 应用场景:适用于事件广播,例如一个“订单创建”事件需要被库存、营销等多个系统同时感知和处理。
现代 MQ(如 Kafka)的模型融合:
- 通过消费者组的概念融合了两种模型。
- 不同消费者组之间是发布/订阅模式(每个组都能收到全量消息)。
- 同一消费者组内部是点对点模式(组内只有一个消费者能处理某条消息)。
【中等】MQ 如何实现延迟消息?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:12 min | 🏷 标签:MQ / 延迟消息
💎 关键结论
延迟消息的本质是「到期才投递」,主流有三种实现范式:专用调度存储 + 定时扫描重投(RocketMQ 4.x)、TTL + 死信中转(RabbitMQ 原生)、延迟交换机 / 时间轮(RabbitMQ 插件、RocketMQ 5.0)。选型核心是延迟时间是否需要逐条动态设置、以及对精度和存量的要求。Kafka 无原生延迟消息,只能靠外部调度(如时间轮服务 + 二次投递)模拟。
⚡记忆卡片
- 口诀:定级别扫描重投、TTL 过期转死信、时间轮任意定
- 关键词:delayTimeLevel / SCHEDULE_TOPIC_XXXX / x-message-ttl / x-delayed-message / 时间轮
- 链路:发送时携带延迟参数 → Broker 转存到调度存储(专用 Topic/中转队列/插件存储)→ 到期检查(扫描/过期/时间轮)→ 重新投递回目标 Topic/队列 → 消费者正常消费
📖 核心知识
三种实现范式
| 范式 | 代表实现 | 延迟粒度 | 精度 | 特点 |
|---|---|---|---|---|
| 调度 Topic + 定时扫描重投 | RocketMQ 4.x | 固定 18 个级别 | 秒级 | 按级别分队列避免全量遍历,吞吐高 |
| TTL + 死信交换机(DLX)中转 | RabbitMQ 原生 | 队列级固定 | 秒级 | 无需插件,但存在队首阻塞 |
| 延迟交换机 / 时间轮 | RabbitMQ 插件、RocketMQ 5.0 | 消息级任意时间 | 毫秒级 | 逐条动态延迟,O(1) 槽位定位到期检查 |
通用流程
- 发送阶段:Producer 发送时携带延迟参数(级别/时间戳/毫秒数)。
- 暂存阶段:Broker 不把消息直接放入目标 Topic,而是转存到调度存储(内部 Topic、中转队列或插件存储)。
- 到期检查:后台定时任务扫描、TTL 过期触发或时间轮槽位到期。
- 重新投递:到期消息被移回目标 Topic/队列,消费者正常消费;此后与普通消息无区别,失败同样走重试与死信流程。
Kafka 的缺口
Kafka 原生不支持延迟消息(消息写入分区即可见)。工程上通常用外部时间轮/调度服务暂存到期时刻,到期后再写入 Kafka;或把消息发到专用 Topic 由消费者轮询判断到期(堆积大、精度差)。
典型业务形态
- 订单超时关闭:下单 30 分钟未支付自动取消(到期消费时检查支付状态,幂等处理)。
- 重试退避:失败任务延迟一段时间再触发。
- 定时提醒:预约通知、日程提醒。
🔬 扩展知识
【L3】RocketMQ 定制深挖:4.x 的 18 级延迟
详情
- Producer 通过
setDelayTimeLevel(3)指定级别(1s/5s/10s/30s/1m...2h 共 18 级),Broker 将消息转存到SCHEDULE_TOPIC_XXXX按级别分的队列(级别 3 的消息进 Queue3),定时任务每秒扫描到期消息重投回原 Topic。 - RocketMQ 5.0+ 支持任意时刻定时消息:
setDeliverTimeMs直接指定时间戳,底层基于时间轮 + 专用存储引擎,精度秒级;时间轮把到期检查从遍历变为 O(1) 槽位定位,海量定时任务下开销显著低于 4.x 逐队列轮询。 - 选型建议:延迟时间可枚举为少数几档时用 4.x 级别即可,无需升级 5.0。
【L3】RabbitMQ 定制深挖:TTL+DLX 与延迟插件
详情
- TTL+DLX:中转队列设
x-message-ttl+x-dead-letter-exchange,消息过期后由 DLX 路由到目标队列;缺陷是消息级 TTL 只在队首检查过期,队首长 TTL 消息会阻塞后面短 TTL 消息(队首阻塞),缓解办法是按延迟时长拆多个队列。 - 延迟插件:
rabbitmq_delayed_message_exchange提供x-delayed-message交换机,消息带x-delay(毫秒)逐条动态延迟,到期由插件内部调度投递;代价是插件存储(Mnesia/磁盘)在百万级长期延迟场景内存与磁盘开销可观,且吞吐低于原生交换机。 - 关键业务同样要开启持久化并监控插件交换机积压;海量长期延迟场景可评估时间轮调度 + 普通 MQ 的自研方案。
# 插件方案:声明延迟交换机并发送延迟 5 秒的消息
channel.exchange_declare(
exchange='delayed_exchange',
exchange_type='x-delayed-message',
arguments={'x-delayed-type': 'direct'}
)
channel.basic_publish(
exchange='delayed_exchange',
routing_key='order_queue',
body=message,
properties=pika.BasicProperties(headers={'x-delay': 5000})
)【L4】延迟消息到期后就是普通消息:消费失败走正常重试与死信流程;因此延迟消息的可靠性设计要与消费重试、死信治理一起规划,见本文档『什么是死信队列?如何设计死信队列的处理机制?』。
🔀 发散问题
Q:订单 30 分钟未支付自动关闭怎么实现?
A:下单时发一条延迟 30 分钟的消息,到期后消费检查支付状态:已支付则忽略,未支付则关单并释放库存;消息体带订单号保证幂等。
Q:三种方案怎么选型?
A:延迟可枚举为少数几档 → RocketMQ 4.x 级别或 RabbitMQ 按档位拆队列;逐条动态延迟 → RabbitMQ 延迟插件或 RocketMQ 5.0 定时消息;Kafka 技术栈则需外部时间轮服务二次投递。
MQ 集群
【困难】Kafka、RocketMQ、RabbitMQ 生产消费有什么相同和不同之处?⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:15 min | 🏷 标签:MQ / 生产消费
💎 关键结论
三者都支持同步/异步发送、压缩、多序列化与分片并发;核心差异在推拉模型、批量能力与扩展方式。理由:Kafka/RocketMQ 以批处理与拉模式拼吞吐,RabbitMQ 以推模式与灵活路由拼交付确定性。
⚡记忆卡片
口诀:K/R 拉+批量拼吞吐,Rabbit 推+路由拼交付
关键词:同步异步/批量/推拉模型/序列化/压缩/分片并发
链路:生产端同步/异步发送 → 批量+压缩写入分片 → 消费端拉取(Kafka/RocketMQ)或推送(RabbitMQ)→ 分片内独占消费保序
📖 核心知识
Kafka、RocketMQ、RabbitMQ 生产消费的相同点:
- 生产消息都支持同步、异步发送。
- 都支持消息压缩以节省带宽/存储。
- 都支持多种序列化方式。
- 都通过分区机制来支持并发与扩展。
Kafka、RocketMQ、RabbitMQ 生产消费的不同点:
| 特性维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 生产/消费模型 | 基于日志的发布/订阅 | 基于队列的发布/订阅 | 灵活的路由 (Exchange -> Queue) |
| 同步/异步发送 | 同步:send().get()异步: send(callback)单向: send() 无回调(不关心结果) | 同步:send() 阻塞异步: sendAsync(callback)单向: sendOneway() | 同步:使用 Channel 的确认机制异步:使用 Confirm Mode 和回调 |
| 推拉模型 | 消费者主动拉取 | 消费者主动拉取 | Broker 主动推送 (默认) (也支持拉取 Get,但不常用) |
| 是否支持批量生产 | 是,核心功能 生产者会累积消息,批量发送到同一 Partition,极大提升吞吐 | 是,核心功能 支持将多条消息打包成一个批次发送 | 是,但较弱 通过 Publisher Confirms 和 Tx 模拟,或使用第三方插件。原生对批量不友好 |
| 是否支持批量消费 | 是,核心功能 消费者一次拉取一批消息进行处理。 | 是,核心功能 消费者可以设置一次拉取的消息条数 | 是,通过 Qos 通过设置 prefetchCount 来控制未确认消息的批量大小,实现“准批量” |
| 序列化方式 | 高度自由,与协议解耦 生产者/消费者自行配置序列化器(如 String, JSON, Avro, Protobuf) | 高度自由,与协议解耦 同上,客户端可灵活配置 | 与协议绑定 依赖于 AMQP 等协议规定的格式,灵活性较低。消息体为二进制,具体格式由应用层决定 |
| 是否支持压缩 | 是,端到端压缩 支持 Snappy, Gzip, LZ4, Zstandard。生产者压缩,消费者解压,节省带宽和存储。 | 是,端到端压缩 支持 Gzip, Snappy, LZ4, Zstandard。机制与 Kafka 类似。 | 是,传输层压缩 通常指在传输协议层面(如 TLS)的压缩,而非消息内容的端到端压缩。 |
| 如何支持并发 | 基于 Partition 一个 Consumer Group 内的多个消费者并行消费同一个 Topic 的不同 Partition 单 Partition 内消息有序,Partition 间无序 | 基于 Queue 一个 Consumer Group 内的多个消费者并行消费同一个 Topic 的不同 Queue 单 Queue 内消息有序,Queue 间无序 | 基于 Queue 多个消费者绑定到同一个 Queue 时,消息以竞争方式被消费(如 Round-Robin) 通过多个 Queue 来实现并发 |
| 如何扩展 | 水平扩展:增加 Partition 通过增加 Topic 的 Partition 数量来提升并发度 消费者数量最好 ≤ Partition 数量 | 水平扩展:增加 Queue 通过增加 Topic 的 Queue 数量来提升并发度 消费者数量最好 ≤ Queue 数量 | 水平扩展:增加 Queue/Node 纵向:为一个 Exchange 绑定多个 Queue 横向:搭建集群,使用 镜像队列 实现高可用 |
核心差异总结
- 推拉模型差异:
- Kafka & RocketMQ:拉模式。消费者按需拉取消息,可以方便地实现批处理和流量控制。
- RabbitMQ:推模式。Broker 主动将消息推送给消费者,延迟更低,消费者需要有背压机制。
- 并发与有序性:
- 三者都通过 分区机制(Partition/Queue) 来实现并发。
- 它们都通过在 分区内部保证顺序 来兼顾并发与局部有序。若需要全局有序,就必须只有一个分片和一个消费者,这会牺牲并发性。
- 批处理:
- Kafka & RocketMQ:批处理是核心设计理念,从生产到消费都深度优化,是追求极致吞吐量的关键。
- RabbitMQ:设计核心是复杂的路由和可靠的单条消息传递,批量处理是其短板,通常通过 prefetch 来模拟。
- 扩展性:
- Kafka & RocketMQ:扩展性直接与主题的分片数挂钩,规划初期需要合理设置。
- RabbitMQ:扩展性更灵活,可以通过增加队列或搭建集群来实现,但其队列镜像对性能有一定影响。
🔬 扩展知识
【L3】批量机制的实现差异
详情
- Kafka 生产端按「Partition + 时间/大小」双条件积累批次,一次网络请求发送一批;消费端一次 Fetch 返回一批记录。
- RocketMQ 支持批量发送但对单批大小有限制(默认不超过 4MB,需自行分批)。
- RabbitMQ 无原生批量协议,靠
prefetchCount控制未确认消息窗口实现“准批量”。
【L4】序列化与 Schema 治理
详情
- Kafka 生态常用 Schema Registry(如 Confluent Schema Registry)管理 Avro/Protobuf 模式演进;RocketMQ 需业务自行约定消息格式。
- RabbitMQ 消息体对 Broker 透明,格式完全由应用层决定,跨语言集成时需额外约定。
🏭 实战场景
详情
(推演案例)某日志平台峰值写入约 50 万条/秒,单条消息平均 500 字节:若选 Kafka,生产端批量积累(如 linger.ms=10ms、batch.size=1MB)可将网络请求数降低一到两个数量级,单集群即可承载;若同一链路改用 RabbitMQ 推模式,需按单条消息估算,同等吞吐下 Broker 内存与网络开销显著升高,且堆积能力受限。反之,若业务是“订单创建后 100ms 内通知下游”,消息量仅千级 TPS,RabbitMQ 推模式的低延迟与灵活路由更合适。结论:批量吞吐型选 Kafka/RocketMQ,低延迟交付型选 RabbitMQ。
⚠️ 常见误区
详情
常见误区:
- ❌ “RabbitMQ 不支持批量消费” → 不准确。RabbitMQ 通过
prefetchCount预取多条未确认消息可实现准批量,只是无原生批量拉取协议。 - ❌ “Kafka 也支持 Push 消费” → Kafka 原生仅 Pull;所谓推送效果是客户端封装轮询的结果。
- ❌ “增加消费者总能提升消费速度” → 消费者数超过分片数后多出的实例空转,并发上限由分片数决定。
🔀 发散问题
为什么三者都不直接保证全局有序?
全局有序意味着写入与消费只能单点串行,吞吐无法水平扩展,与 MQ 的核心价值冲突;因此都把有序粒度下放到分片,详见本文档『如何保证 MQ 消息的顺序性?』。
拉模式下如何兼顾实时性?
用长轮询:Broker 无消息时挂起拉取请求,有消息立即返回;RocketMQ 默认长轮询 15s,见本文档『MQ 数据传输有哪些模式?』。
【困难】Kafka、RocketMQ、RabbitMQ 副本机制有什么不同之处?⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:15 min | 🏷 标签:MQ / 副本与高可用
💎 关键结论
Kafka 靠 ISR 动态同步副本集合保强一致,RocketMQ 提供 Raft(DLedger)强一致与异步复制两种模式,RabbitMQ 以镜像队列复制、偏最终一致。理由:三者在“一致性强度 vs 自动化程度 vs 使用使利性”上做了不同取舍。
⚡记忆卡片
口诀:Kafka ISR,Rocket 可选 Raft,Rabbit 镜像队
关键词:ISR/Leader 选举/DLedger/主从/镜像队列/仲裁队列
链路:消息写入主副本 → 同步/异步复制到从副本 → 主副本故障时自动选举或提升新主 → 客户端重连或感知新主
📖 核心知识
| 特性维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 核心模型 | 主从复制 | 主从复制 | 镜像队列 (Mirrored Queue) |
| 副本单位 | Partition (分区) | Broker 组 (Broker Group) | Queue (队列) |
| 数据一致性 | ISR (同步副本集合) 机制 Leader 维护一个与其保持同步的 Follower 列表(ISR) 消息只有被 ISR 中的所有副本 应用后才能被消费者读取 提供强一致性保证 | 多数派写入 + 主从同步 同步双写 (DLedger):基于 Raft 协议,消息写入多数节点后才确认,强一致 异步复制:消息在主节点落盘即返回,数据有丢失风险,最终一致 | 镜像队列同步 所有发布到主队列的消息都会同步到所有镜像 支持配置同步镜像(强一致,性能低)和异步镜像(最终一致,性能高) |
| 读写流量 | 主副本读写 所有读写请求只由 Leader 副本处理 Follower 副本只负责从 Leader 异步拉取数据进行同步 | 主副本读写 所有读写请求只由 Master 节点处理 Slave 节点作为热备,只同步数据,不提供服务 | 主队列读写 客户端可与集群中任何节点通信 但针对某个队列的读写请求最终都会被路由到该队列的主节点上 |
| 故障转移 | 自动选举 Controller 从 ISR 集合 中自动选举新的 Leader 服务不可用时间短,自动化程度高 | 自动切换 DLedger 模式:基于 Raft 协议自动选举新 Master 主从模式:依赖 NameServer 和监控脚本进行主从切换,可能需人工介入 | 自动提升 当主队列所在节点宕机时,最老的镜像会被自动提升为新的主队列 客户端需要重连才能感知到新的主节点 |
| 优点 | 高吞吐:主写主读,压力集中 强一致性:ISR 机制保证数据不丢失 自动化:故障转移自动完成 | 灵活可选:提供强一致和最终一致两种模式 金融级可靠:DLedger 模式保证数据强一致和高可用 热备份:Slave 节点数据完整,切换快 | 使用简单:对客户端透明,连接任何节点即可 灵活配置:可针对不同队列设置不同的镜像策略 高可用:队列数据在多个节点冗余 |
| 缺点 | 资源利用率低:Follower 副本只做备份,不服务读写 ISR 抖动:Follower 同步慢会被踢出 ISR,影响可用性 | 主从模式非强一致:异步复制有数据丢失风险 主从切换可能延迟:非 DLedger 模式下,自动化程度不如 Kafka | 性能瓶颈:同步镜像模式下,性能受限于最慢的镜像节点 网络压力大:所有写操作都需要在镜像间同步,网络 IO 高 扩展性限制:队列的镜像数量越多,性能和网络压力越大 |
核心差异总结
- 设计哲学与一致性:
- Kafka:以 ISR 机制为核心,在保证高吞吐的同时,提供了强一致性的折中方案。它不要求所有副本都写入成功,只要求一个动态的、健康的副本集合(ISR)同步即可。
- RocketMQ:提供了 灵活的选择。在金融场景下,可以使用基于 Raft 的 DLedger 模式获得强一致性;在普通场景下,可以使用异步复制获得更高吞吐。
- RabbitMQ:以 镜像队列 为单位进行复制,其一致性级别(强一致或最终一致)取决于镜像配置,更偏向于最终一致性和使用便利性;新版本推荐基于 Raft 的**仲裁队列(Quorum Queue)**替代镜像队列。
- 读写模式:
- Kafka & RocketMQ:都是 “主写主读” 。所有流量都集中在主副本,架构简单,延迟可控。
- RabbitMQ:是 “主写主读,但客户端可连任意节点” 。对客户端更友好,但最终压力还是在主节点上。
- 故障转移:
- Kafka:自动化程度最高,从预设的 ISR 中快速自动选举。
- RocketMQ (DLedger):自动化且强一致,基于 Raft 协议自动选举。
- RocketMQ (主从) & RabbitMQ:自动化但有延迟或依赖。RocketMQ 主从可能需监控脚本;RabbitMQ 客户端需要重连来感知新主节点。
🔬 扩展知识
【L3】ISR 收缩的连锁影响
详情
- Follower 同步落后超过阈值(Kafka 旧版本
replica.lag.time.max.ms)会被踢出 ISR;若acks=all且min.insync.replicas=2,ISR 收缩到只剩 Leader 时写入直接报错,可用性换可靠性。 - 若允许非 ISR 副本选主(
unclean.leader.election.enable=true),故障切换时可能丢失未同步消息,资金链路应禁用。
【L4】仲裁队列对镜像队列的替代
详情
- RabbitMQ 官方已将镜像队列标记为弃用(deprecated),推荐基于 Raft 协议的 Quorum Queue:多数派确认、自动选举,解决了镜像队列主切换时的消息丢失与重复问题。
- 代价是仲裁队列仅支持持久化消息,内存开销高于经典队列。
【L3】Kafka 定制深挖:ISR 同步判定与 LEO/HW
详情


- 同步判定:
replica.lag.time.max.ms(旧版默认 10s,新版 30s)内 Follower 未拉取到新数据即被踢出 ISR;ISR 是动态集合,Leader 天然在集合内。 - LEO/HW:LEO 是每副本下一条待写入的 offset;HW = ISR 中最小 LEO,消费者只能读 HW 之前的已提交消息。HW 更新链:Follower 拉取携 LEO → Leader 更新远程 LEO → Leader 更新 HW → Fetch 响应携带 HW 回传 Follower;因此 HW 更新天然慢一轮,旧版基于 HW 截断可能丢已提交消息。
- Leader Epoch(0.11+):每分区维护
leader-epoch-checkpoint记录<epoch, startOffset>,Follower 重启后用OffsetsForLeaderEpoch询问合法截断点,替代基于 HW 的截断,消除数据丢失与脑裂不一致。 - 首选 Leader:创建分区时指定的首领为首选首领,
kafka-preferred-replica-election可把反复切换后偏移的 Leader 切回,恢复集群负载均匀。 - KIP-392(2.4+):消费者可从就近副本读(
replica.selector),仅限能容忍短暂不一致/跨机房省带宽的场景;默认仍是只读 Leader。
【L3】RocketMQ 定制深挖:CommitLog 复制单位与 brokerRole
详情
- 复制单位是 CommitLog 而非主题/队列:Master 处理读写,Slave 只读并拉取同步;ConsumeQueue 是确定性派生索引,各节点独立重建不复制,省一半复制流量。
brokerRole双模式:SYNC_MASTER等 Slave 写入成功后才返回,可靠性高但延迟增加;ASYNC_MASTER立即返回,Master 宕机丢未同步窗口(通常最后数百毫秒到数秒写入量)。- 同步复制要求 BrokerGroup 内 Slave 在线,Slave 全部不可用时发送失败,Producer 依赖故障规避重试到其他 Broker,消息不丢但会重复。
- 传统主从“复制不切主”:4.5+ DLedger(Raft)补齐自动选主;5.0 Controller 模式引入 ISR 概念,只有同步进度在 ISR 内的 Slave 才有资格选为新 Master,兼顾自动切换与数据不丢。
【L3】RabbitMQ 定制深挖:镜像队列配置与仲裁队列迁移
详情
- 镜像队列用 Policy 启用:
rabbitmqctl set_policy ha-all "^ha\." '{"ha-mode":"all"}';ha-mode可选 all/exactly/nodes,ha-sync-mode: automatic更安全。 - 版本事实:镜像队列自 3.8 起被仲裁队列取代,4.0 已彻底移除;新项目直接声明
x-queue-type=quorum,副本是队列类型内置属性,无需 ha-mode 策略。 - 镜像队列可靠性缺陷(官方弃用的根本原因):自研主从异步复制协议,主宕机丢未同步消息;新主可能选到数据落后节点(
ha-promote-on-failure控制);网络分区时两侧各自主导致消息分叉。 - 故障切换对客户端基本透明但非无感:Channel 未完成操作会失败,未确认消息重投,客户端需重连重试,消息可能重复,消费端必须幂等。
📚 延伸阅读:RabbitMQ 官方文档
🏭 实战场景
详情
(推演案例)某支付平台要求消息零丢失且 RTO 小于 30 秒:若采用 RocketMQ 主从异步复制,Master 宕机可能丢失未同步消息,不满足要求;改用 DLedger 三节点 Raft 组,多数派写入确认,自动选主切换通常在秒级完成,满足要求,代价是写入吞吐相比异步复制有所下降(经验上约下降三成左右,需实测)。若同样要求放在 RabbitMQ,应选仲裁队列而非同步镜像队列,后者在节点故障切换时存在消息丢失的历史缺陷。
⚠️ 常见误区
详情
常见误区:
- ❌ “Kafka 的 Follower 可以分担读流量” → 默认不行,读写集中在 Leader;Follower 只做同步(新版本 KIP-392 提供受限的 Follower 读,非默认主流用法)。
- ❌ “RabbitMQ 镜像队列一定不丢消息” → 主节点宕机而消息尚未同步到镜像时仍会丢;官方已推荐仲裁队列替代。
- ❌ “RocketMQ 主从模式能自动切换” → 传统主从模式切换依赖外部工具/人工,自动选举需 DLedger 模式。
🔀 发散问题
副本数是不是越多越好?
不是。副本数增加写放大与网络开销,常用 3 副本配合 min.insync.replicas=2 在可靠性与性能间取平衡;副本数与确认策略的关系见本文档『如何保证 MQ 消息不丢失?』。
ISR 机制和 Raft 多数派有什么本质区别?
ISR 是“Leader 主导的动态同步集合”,确认要求 ISR 内全部副本同步;Raft 是“多数派写入 + 任期选举”,无固定的全同步集合。前者吞吐更高但依赖 Leader 正确维护 ISR,后者一致性证明更严谨。
【困难】Kafka、RocketMQ、RabbitMQ 分区机制有什么不同之处?⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:15 min | 🏷 标签:MQ / 分区与路由
💎 关键结论
Kafka/RocketMQ 靠分片(Partition/Queue)实现并行与局部有序,RabbitMQ 靠 Exchange 路由实现灵活分发。理由:前者把“顺序与并发”交给分片键,后者把“分发规则”交给路由策略,设计哲学不同。
⚡记忆卡片
口诀:K 哈希轮询,R 选择器,Rabbit 路由规则
关键词:Partition/MessageQueue/Exchange/Routing Key/分片键/独占消费
链路:生产者按分片键路由到分片 → 分片内追加保序 → 组内一个消费者独占一个分片 → 增加分片提升并发
📖 核心知识
核心分区/分片机制对比
| 特性维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 核心概念 | 分区 | 队列 | 队列 |
| 逻辑实体 | Topic | Topic | Exchange + Queue |
| 消息有序性 | 保证 Partition 内有序,Partition 间无序 | 保证 Queue 内有序,Queue 间无序 | 保证 Queue 内有序。如果多个消费者消费同一队列,顺序无法保证 |
| 消息路由方式 | 指定 Key:相同 Key 的消息被哈希到同一 Partition 轮询:保证分区间负载均衡 随机/自定义 | 轮询:默认方式,保证队列间负载均衡 哈希:通过 Sharding Key 选择队列 消息队列选择器:自定义算法 | Exchange 类型 和 Routing Key 决定: Direct:精确匹配 Routing Key Topic:模糊匹配模式 Fanout:广播到所有绑定队列 Headers:匹配消息头 |
| 与消费者的关系 | 一个 Consumer Group 内,一个 Partition 只能被一个 Consumer 独占 | 一个 Consumer Group 内,一个 Queue 只能被一个 Consumer 独占 | 一个 Queue 可以被多个 Consumer 共享竞争。消息以轮询方式分发给这些消费者 |
| 扩展性 | 通过增加 Topic 的 Partition 数量来实现水平扩展 注意:增加 Partition 数量需要重启或使用工具,操作较重 | 通过增加 Topic 的 Queue 数量来实现水平扩展 注意:Queue 数量在创建后通常不建议修改 | 通过增加 Queue 的数量并与 Exchange 绑定来实现扩展,操作相对灵活 |
| 设计哲学 | “日志流” 分片:将巨大的 Topic 日志流切分成多个 Partition,分散存储和压力 | “消息队列” 分片:概念上更接近传统的队列,将一个 Topic 的消息负载均衡到多个队列 | “消息路由”:核心在于 Exchange 如何将消息灵活地路由到不同的队列,队列是消息的终点 |
小结
- Kafka & RocketMQ:
- 通过 分区/队列 进行分片,实现并行。
- 通过 分区/队列内独占消费 来保证消息顺序。
- 提升并发必须增加分片。
- RabbitMQ:
- 通过 Exchange 路由 将消息分发到不同队列。
- 通过 队列内共享消费 来提升并发,但会牺牲消息顺序。
- 路由方式极其灵活。
🔬 扩展知识
【L3】分片数变更的乱序风险
详情
- 扩分片后哈希取模结果变化,同一分片键的新老消息可能路由到不同分片,跨分片无序;因此扩分片需配合双写过渡或业务低峰操作。
- Kafka 支持增加 Partition 但不可减少;RocketMQ Queue 数量创建后通常不建议修改。
【L4】RabbitMQ 一致性哈希交换器
详情
- RabbitMQ 可通过插件
x-consistent-hash交换器把消息按哈希分布到多个队列,近似模拟 Kafka 的分片语义,用于需要按键聚集的场景。 - 但它不具备分区再均衡机制,消费者与队列的绑定需自行管理。
🏭 实战场景
详情
(推演案例)某订单系统要求同一订单的“创建→支付→发货”严格有序,日订单 100 万笔、峰值 5000 TPS:方案是以订单号为分片键哈希到 64 个分片,组内 64 个消费者独占消费,单分片峰值约 80 TPS,容量充足;若后续订单量增长 10 倍,不能直接扩到 640 分片(哈希路由变化导致乱序),应采用新建 Topic 双写迁移或临时转发方案。若同一需求放在 RabbitMQ,需以订单号作 Routing Key 绑定到固定队列,或采用单队列+单消费者,后者吞吐受限。
⚠️ 常见误区
详情
常见误区:
- ❌ “Kafka 的 Topic 整体有序” → 只有单 Partition 内有序,跨 Partition 无序;全局有序需单分区,见本文档『如何保证 MQ 消息的顺序性?』。
- ❌ “RabbitMQ 多消费者消费同一队列也能保序” → 共享竞争消费时不同线程完成时间不定,顺序无法保证。
- ❌ “分片数越多越好” → 分片过多会增加元数据、选举与文件句柄开销(如 Kafka 单机分区过多时负载飙高),应按峰值并发规划。
🔀 发散问题
分片键应该怎么选?
选“需要互相保序的最小业务实体 ID”(如订单号、账户号),同时保证分布均匀避免热点;无业务键的日志类消息用轮询即可。
消费者数超过分片数会怎样?
多出的消费者分不到分片而空转;因此容量规划时先定分片数再定消费者数,经验上分片数按峰值消费并发预留冗余(见本文档『如何处理 MQ 消息积压?』)。
MQ 可靠传输
【困难】如何保证 MQ 消息不丢失?⭐⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:20 min | 🏷 标签:MQ / 可靠性
💎 关键结论
不丢消息 = 生产端确认 + 服务端持久化/多副本 + 消费端处理完再 ACK + 监控对账兜底,三环节缺一不可。理由:丢消息有三条独立路径(网络丢失、存储损坏、过早确认),只加固任何一端都会被其他路径击穿。
⚡记忆卡片
口诀:生产要确认,服务端多副本,消费后 ACK,对账兜底
关键词:acks/min.insync.replicas/同步刷盘/手动提交 Offset/Publisher Confirms/对账
链路:生产端收到确认才视为发送成功 → 服务端多副本+持久化才视为存储成功 → 消费端处理完才提交位点/ACK → 堆积告警+定期对账发现漏网
📖 核心知识
要保证 MQ 中的消息不丢失,需从 生产端、MQ 服务端、消费端 三个环节进行可靠性设计。一言以蔽之,生产端确认+服务端持久化+消费端手动 ACK+监控补偿 是保证消息不丢失的核心逻辑,需根据业务场景权衡性能与可靠性。
一条消息从生产到消费,可以划分三个阶段,每个阶段都可能丢失数据:

- 生产阶段:Producer 创建消息,并通过网络发送给 Broker。
- 存储阶段:Broker 收到消息并存储,如果是集群,还要同步副本给其他 Broker。
- 消费阶段:Consumer 向 Broker 请求消息,Broker 通过网络传输给 Consumer。
(1)生产端防丢失
| Kafka | RocketMQ | RabbitMQ | |
|---|---|---|---|
| 发送方式 | 支持同步、异步、异步回调发送方式 同步性能低;异步无视成功与否;异步回调比较合适 | 支持同步、异步回调发送方式 同步性能低;异步回调比较合适 | 同步、异步发送 通过 Publisher Confirms 确认投递 |
| 生产 ACK | acks 参数确保消息被多少副本写入成功后才返回确认min.insync.replicas 配合 acks 使用,设定最少写入副本数 | (配合刷盘机制)若配置为同步刷盘,消息必须先成功写入磁盘,发送方才能收到写成功的确认 | Publisher Confirms 模式 Broker 收到并持久化后返回 Basic.Ack |
| 失败重试 | retries 参数控制失败重试次数 | retryTimesWhenSendFailed 参数控制失败重试次数 | 监听 Basic.Nack 或超时后业务侧重发 |
| 事务 | 通过事务+生产幂等,实现 exactly once 语义,非真正的分布式事务 | 可实现分布式事务 | Channel 事务(txSelect/txCommit),性能差,不推荐,应使用 Confirms 替代 |
(2)MQ 服务端防丢失
| Kafka | RocketMQ | RabbitMQ | |
|---|---|---|---|
| 持久化 | 消息写入页缓存(PageCache),即可返回生产 ACK,由 OS 负责刷盘(fsync)可设置自动刷盘的间隔时间或消息数阈值 可以通过 AdminClient 手动刷盘 | 默认采用异步刷盘,类似 Kafka,先写入页缓存(PageCache),由 OS 负责刷盘(fsync)也支持同步刷盘,消息必须先成功写入磁盘,发送方才能收到写成功的确认 | 队列 durable=true + 消息 deliveryMode=2 持久化到磁盘支持异步刷盘和同步刷盘 |
| 副本机制 | 通过副本机制保证冗余,避免单点故障 副本的粒度是针对分区 replication.factor 设置分区副本数min.insync.replicas 控制写入到多少个副本才算是“已提交” | 通过副本机制保证冗余,避免单点故障 副本的粒度是针对节点(Broker) | 镜像队列(Mirrored Queue)或仲裁队列(Quorum Queue) 镜像队列为异步复制,仲裁队列基于 Raft 强一致 |
| 故障检测 | 所有 Broker 会向 zk 的/broker路径写临时节点,而 Controller 会监听该目录一旦有 Broker 故障或失联,zk 会话中断,zk 自动删除临时节点,并通知 Controller Controller 将 Broker 视为下线,并将该 Broker 上的分区 Leader 置为无效 | 主节点定期向从节点发送心跳以续活 超时未收到消息,从节点视主节点为下线 | 节点间通过 Erlang 分布式通信维持心跳 GM 协议(镜像队列)或 Raft 协议(仲裁队列)检测节点存活 |
| 故障恢复 | Controller 负责选新的分区 Leader,优先从 ISR 中选 通过 unclean.leader.election.enable 可设置是否允许从非 ISR 中选 Leader,但丢失数据风险增大 | 采用 Raft 来选主 | 镜像队列:最老的镜像自动提升为主 仲裁队列:基于 Raft 选举新主,保证数据完整 |
(3)消费端防丢失
| Kafka | RocketMQ | RabbitMQ | |
|---|---|---|---|
| 消费 ACK | 支持自动、手动提交 Offset(enable.auto.commit=false配置)关闭自动提交,处理完整事务后,再手动提交 Offset | Offset 提交由消费返回状态驱动:消费成功后返回 CONSUME_SUCCESS 才提交位点;失败返回 RECONSUME_LATER 进入重试 | 手动 ACK(autoAck=false)处理成功调用 basicAck;失败调用 basicNack/basicReject |
| 重置 Offset | auto.offset.reset 可以重置消费的 Offset | 支持按时间戳重置消费位点 | 消息未 ACK 会自动重新入队,无需重置 |
| 失败重试 | 不完善:自动提交不支持重试;手动提交需自行处理重试逻辑 | 默认重试 16 次 多次消费失败后会被投递到死信队列,供后续特殊处理 | 手动控制:requeue=true 重新入队,requeue=false 转入死信队列 |
(4)监控与补偿
- 消息堆积告警:监控队列长度,及时发现异常。
- 定期对账:对比生产与消费的记录,修复差异(如定时扫描数据库补发)。
🔬 扩展知识
【L3】可靠性组合与性能代价
详情
| 组合方案 | 可靠性 | 吞吐影响(经验值) | 适用边界 |
|---|---|---|---|
| 异步发送 + 异步刷盘 + 自动提交消费位点 | 低 | 基准 | 日志、埋点等允许丢失的场景 |
| 异步回调发送 + 多副本同步确认 + 手动提交消费位点 | 高 | 吞吐下降约 30%~50% | 订单、库存等核心业务 |
| 同步发送 + 同步刷盘 + 手动 ACK + 对账 | 极高 | 吞吐再降 50% 以上,延迟毫秒级升至数十毫秒级 | 资金类业务,通常仅对部分关键消息启用 |
核心取舍:不丢消息 = 用延迟和吞吐换确认链路。三个环节任一环节的确认缺失都是丢消息的路径,不能只加固一端。
失效场景(保障在什么条件下失效):
- 服务端只配置多副本、但允许非同步副本选主:副本全部异步复制时,Leader 宕机后由落后副本上位,已写入但未同步的消息直接丢失。
- 副本确认要求“多数副本写入”但同步副本集合收缩到只剩 Leader 自己:此时确认退化为单点写入,Leader 宕机即丢。
- 消费端先提交位点/先 ACK 后处理:处理中途宕机,消息永久跳过(丢失路径)。
- 消费端先处理后提交,但处理成功、提交前宕机:不丢消息,转为重复消费问题(需幂等兜底)。
- 消息保留期过短 + 消费者长时间宕机:消息被服务端清理,连重放的机会都没有。
【L4】同步刷盘的残留风险
详情
- 同步刷盘只保证“单机落盘”,若机器整机损毁(磁盘故障、机房断电)仍会丢;残留风险要靠多副本 + 跨机架/跨机房部署消除。
- 同步刷盘使单次写入延迟从亚毫秒级升至数毫秒级,吞吐下降一个数量级,通常只对资金类 Topic 开启。
【L3】Kafka 定制深挖:可靠性参数组合与失效边界
详情
关键配置:
# 生产者
acks=all
enable.idempotence=true
# Broker
replication.factor=3
min.insync.replicas=2
# 消费者
enable.auto.commit=false可靠性三参数的组合权衡
| 组合 | 语义 | 代价 | 结论 |
|---|---|---|---|
acks=1 + min.insync.replicas=1 + unclean.leader.election.enable=true | Leader 写入即确认 | 吞吐最高(基准) | 日志类可接受,业务消息禁用 |
acks=all + min.insync.replicas=1 | 看似严格,实际等同 acks=1 | 无额外保护 | 陷阱配置:ISR 收缩到只剩 Leader 时照样丢 |
acks=all + replication.factor=3 + min.insync.replicas=2 + unclean.leader.election.enable=false | 至少 2 副本写入才确认,且只从 ISR 选主 | 吞吐下降约 30%~50%,延迟上升 | 业务消息的标准答案 |
注意:必须满足 replication.factor > min.insync.replicas,推荐 replication.factor = min.insync.replicas + 1(如 3 副本配 2),两者相等时任一副本挂机整分区即不可写。
acks=all 也会丢消息的三种真实路径
- ISR 收缩到 1:两台 Follower 同时宕机或同步超时(
replica.lag.time.max.ms旧版默认 10s、新版 30s)被踢出 ISR,此时min.insync.replicas=1的默认值使acks=all退化为写单副本,Leader 宕机即丢。 - 页缓存未刷盘 + 整机损毁:Kafka 依赖 OS 页缓存,
acks=all只保证写入多副本的页缓存;同机架整机柜断电仍可能丢失未 fsync 数据,需跨机架部署(broker.rack)兜底。 - HW 截断风险(旧版本):基于 HW 的日志截断在“Leader 宕机 + Follower 重启”组合下可能截掉已提交消息;Kafka 0.11+ 的 Leader Epoch 机制已解决,升级版本是正解。
消费端:先提交后消费 vs 先消费后提交

| 提交策略 | 丢失路径 | 重复路径 |
|---|---|---|
| 先提交 offset 后处理 | 处理中宕机,消息永久跳过(丢消息) | 无 |
先处理后提交(enable.auto.commit=false) | 无 | 处理成功但提交前宕机,重启后重消费(重复,需幂等) |
| 自动提交(默认 5s 间隔) | 提交后处理中宕机,丢失窗口 ≤ 5s | 处理完未提交即宕机,重复窗口 ≤ 5s |
结论:业务消息一律选“先消费后提交”,把丢失风险转化为可用幂等兜底的重复风险;配合 session.timeout.ms=45s、max.poll.interval.ms=300s 避免慢处理误触发 rebalance。
量化代价参考(推演示例):acks=all + 3 副本相比 acks=1,端到端延迟从约 5ms 升至约 15~30ms,吞吐下降约 30%~50%。
【L3】RocketMQ 定制深挖:刷盘×复制四象限与事务消息
详情
刷盘方式 × 复制方式四象限
| 组合 | 可靠性 | TPS 量化代价(SSD 经验值,需实测) | 适用场景 |
|---|---|---|---|
| 同步刷盘 + 同步复制 | 最高,单机与集群故障均不丢 | 相比全异步 TPS 约降 50%,RT 增加数毫秒 | 金融、交易链路 |
| 同步刷盘 + 异步复制 | 单机断电不丢,Master 磁盘损毁丢未同步副本 | TPS 约降 40%~50% | 重要业务 |
| 异步刷盘 + 同步复制 | Master 断电有极小概率丢 PageCache 内存消息 | TPS 约降 10%~20%(复制网络开销) | 一般业务常用 |
| 异步刷盘 + 异步复制 | 性能最高,断电丢失窗口最大 | 基线 | 日志、监控类 |
- 同步刷盘:
flushDiskType=SYNC_FLUSH,刷盘成功才向 Producer 确认;代价是每条消息串行等待 fsync。 - 同步复制:
brokerRole=SYNC_MASTER等待 Slave 写入成功后才返回;异步模式 Master 宕机会丢失未同步窗口(最后几百毫秒到数秒)的消息。 - 磁盘监控:磁盘使用率超阈值(默认 75%)触发强制清理,磁盘满或 IO 挂起时写入被拒绝,需提前告警。
事务消息:核心业务(如下单)通过“半消息 + 本地事务确认 + 回查”机制,避免“业务成功但消息未发”的丢失——这是 RocketMQ 相比另外两者的独有能力。
消费失败兜底:消费失败进入 %RETRY% 重试队列,默认重试 16 次(间隔 10s 递增至 2h),超限进入 %DLQ% 死信队列;死信队列若无监控和补偿流程,消息就“逻辑丢失”。
【L3】RabbitMQ 定制深挖:Confirm + Return 与副本选型
详情
生产端:Confirm + Return 双保险
- Confirm 回答“消息到没到 Broker”:开启
confirmSelect()后异步监听Basic.Ack/Nack,收到 Ack 才认为发送成功;异步确认对吞吐影响很小。 - Return 回答“消息有没有进队列”:
mandatory=true时,无法路由的消息通过basic.return退回生产者;不设置则静默丢弃——这是最隐蔽的丢消息点。 - 失效场景:Confirm 回调与业务写库不是原子操作,回调前进程崩溃会出现“业务成功但消息未标记已发”的不一致,需靠对账兜底。
存储端:持久化三件套 + 副本选型
- 持久化:队列
durable=true+ 消息deliveryMode=2+ 交换机durable=true三者缺一不可;仅消息持久化而队列非持久化,重启后队列不存在,消息仍丢。 - 副本机制选型(关键权衡):
| 维度 | 镜像队列(Mirrored) | 仲裁队列(Quorum,3.8+ 推荐) |
|---|---|---|
| 复制 | 主从异步复制,可配 ha-promote-on-failure | Raft 多数派确认,强一致 |
| 可靠性 | 主宕机可能丢未同步消息 | 确认即安全,绝不丢 |
| 脑裂 | 网络分区时可能消息分叉/丢失(pause_minority 策略需预配) | Raft 天然防脑裂,少数派自动停服 |
| 性能 | 延迟低、吞吐高 | 多数派写盘确认,吞吐约降 30%~50%(经验值) |
| 适用 | 允许微量丢失的低延迟场景(已弃维) | 金融、交易等关键业务(生产首选) |
消费端:业务逻辑成功后才调用 channel.basicAck();失败时 basicNack(requeue=false) 转死信队列,而非无限 requeue(重投风暴会打爆下游);autoAck=true 推送即删除,处理失败直接永久丢失。
本地消息表模式:“业务写库成功”与“消息发送成功”不是原子操作——本地事务内写业务表 + 消息记录表(待发),事务提交后异步发送,Confirm 成功更新记录表状态,失败由定时任务扫表重发;若用 MQ 承载交易链路,也可换 RocketMQ 事务消息把这件事交给中间件。
📚 延伸阅读:Kafka 官方文档
🏭 实战场景
详情
量化数据
- 可靠性配置的经验代价:全链路可靠配置相比纯异步配置,吞吐通常下降 30%~50%,端到端延迟从约 10ms 升至数十毫秒。
- 建议保留期:核心业务消息保留 ≥ 7 天(Kafka 默认 7 天),消费端故障窗口超过保留期即无法恢复。
- 监控告警阈值建议:消息堆积超过 10 万条或堆积增速超过消费速率 2 倍即告警;对账任务以分钟级周期运行。
踩坑案例(推演示例):一次“确认成功但消息蒸发”的资损事故
- 现象:支付系统告警,某小时约 3000 笔支付成功消息未被对账服务消费,用户已扣款但订单状态停留在“待支付”。
- 排查:对账库无这批消息的消费记录,MQ 集群侧按消息 ID 查询无存储痕迹;检查服务端日志发现事故期间两台 Broker 先后宕机,某分区发生了“非同步副本上位”。
- 根因:集群当时为了提升吞吐开启了“允许未同步副本参与选主”,且副本同步确认形同虚设(同步副本集合长期只有 Leader);恰逢 Broker 滚动重启触发分区切换,未同步消息被截断。
- 修复:关闭非同步副本选主;强制要求“消息至少写入 2 个副本才确认”且副本数设为 3;上线“支付流水 vs 消费记录”分钟级对账补偿任务,事故损失当日追回。
场景题
场景:凌晨 3 点,客服集中投诉“用户付款成功但订单一直显示待支付”。监控显示 MQ 消费无堆积、生产端发送成功率 100%。你如何排查与修复?
应急处理:先以支付流水表为准启动人工/脚本补偿,把积压的“已支付未成单”订单批量推进状态,止损优先;同时对受影响 Topic 开启按消息 ID 的追踪。
根因分析:生产端成功率 100% 说明消息已收到确认,消费端无堆积说明不是消费慢;怀疑消息在服务端丢失。沿消息 ID 在集群查询存储痕迹,若查不到则定位事故时间窗内的分区主从切换记录,核对是否存在非同步副本上位、副本收缩至单副本后 Leader 宕机的情况。本案中根因正是副本确认配置与选主策略的组合缺陷。
长期方案:修正集群可靠性配置(禁止非同步选主 + 多副本写入确认 + 副本数 3);上线支付流水与消费记录的分钟级对账任务,差异自动补发;对资金链路单独使用高可靠 Topic 并独立监控。
权衡:可靠性配置让该 Topic 吞吐下降约 40%,但资金链路峰值 TPS 远低于容量上限,用可量化的性能代价换取资损归零,是正确的工程决策。
详情
各 MQ 定制踩坑案例(推演示例)
- Kafka:ISR 收缩引发万笔订单消息丢失——Broker 长时间 Full GC 导致 Follower 被踢出 ISR,
min.insync.replicas=1(默认)使acks=all失效,随后 Leader 宕机且unclean.leader.election.enable=true让落后副本上位截断数据。修复:min.insync.replicas=2+ 禁 unclean 选主 + GC 停顿监控。 - RocketMQ:异步刷盘集群断电丢单——
ASYNC_FLUSH+ASYNC_MASTER集群整机断电,PageCache 中最后约 2 秒的消息全部丢失。修复:交易链路改SYNC_FLUSH+SYNC_MASTER(TPS 约降 50%,水平扩容补回),非核心链路保持异步 + 定时对账兜底。 - RabbitMQ:镜像队列脑裂导致消息分叉丢失——网络分区后同一队列两边各有一个“主”,恢复时少数派数据被直接抛弃,已投递消息重复。修复:短期
cluster_partition_handling=pause_minority,长期核心队列迁移仲裁队列(Raft 天然防脑裂)。
场景题:发送成功但消费端“收不到”(RocketMQ 视角)
Producer 日志显示 SEND_OK,但消费端始终没有收到,Broker 没有重启:
- 应急处理:按 msgId/Key 查询消息并查看消息轨迹,确认卡在“发送-存储-消费”哪一环;同时启动对账任务补偿缺失消息,先恢复业务。
- 根因分析:轨迹显示消息在 Broker 上存在但无消费记录 → 大概率是消费端订阅关系不一致(同一 Group 内不同实例订阅的 Tag 不同);Broker 上查不到 → 检查目标队列所在 Broker 是否主从切换或磁盘写失败;多次失败后进死信队列则是消费失败而非丢失。
- 长期方案:强制开启消息轨迹;死信积压接入告警;发布卡点检查订阅一致性。
⚠️ 常见误区
详情
常见误区:
- ❌ “只要服务端做了持久化就不会丢消息” → 持久化只解决存储环节;生产端未收到确认就放弃、消费端提前 ACK 都是独立的丢失路径。
- ❌ “同步刷盘能彻底杜绝丢消息” → 同步刷盘只保证单机落盘,磁盘损坏/机房断电仍丢,需多副本对冲。
- ❌ “消费端先提交位点再处理效率更高” → 处理中途宕机则消息永久跳过,正确顺序是处理成功后再提交。
🔀 发散问题
- 为什么“生产端确认 + 服务端持久化 + 消费端确认”三者缺一不可?只做好服务端不就够了吗?
丢消息有三条独立路径:生产端没收到确认就放弃(网络丢失)、服务端单点存储损坏(存储丢失)、消费端确认过早(逻辑丢失)。只加固服务端,仍会被生产端静默丢弃和消费端提前 ACK 击穿。可靠性遵循“木桶原理”,最弱一环决定上限。 - 消息保留期与消费位点长期滞后的消费者冲突时,如何做工程决策?
两条路:一是延长保留期(存储成本线性上升,按字节计费可量化);二是为长期离线的下游建立“转发层”,将消息转存到对象存储或数仓后再清理。决策依据是“下游故障窗口 × 消息价值”是否大于存储成本。 - 不丢消息与不重复消费是否冲突?
是同一取舍的两面:先处理后确认会因“处理成功但确认失败”产生重复,见本文档『如何保证 MQ 消息不重复?』,需用幂等兜底。
【困难】如何保证 MQ 消息不重复?⭐⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:20 min | 🏷 标签:MQ / 幂等与去重
💎 关键结论
重复不可能在中间件层彻底根除,端到端的“不重复处理”永远靠消费端幂等:业务幂等为基础,缓存/DB 去重为辅助,监控对账保万一。理由:重复的根源之一是“处理成功但确认失败”后的重投,这个窗口在 Broker 视野之外。
⚡记忆卡片
口诀:幂等打底,去重表辅助,对账兜底
关键词:重试/位点回滚/唯一业务 ID/去重表/SETNX/状态机
链路:生产重试/消费超时重投/副本切换 → 同一消息到达多次 → 消费前用唯一业务 ID 判重 → 幂等执行+对账修复漏网
📖 核心知识
(1)MQ 为什么会出现重复消息?
- 生产端重试(网络抖动时自动重发)
- 消费端超时后 MQ 重新投递(如 RabbitMQ 未及时 ACK)
- 消息队列集群脑裂(如 Kafka 副本切换)

以 Kafka 举例,Kafka 每个 Partition 都是一个有序的、不可变的记录序列,不断追加到结构化的提交日志中。Partition 中为每条记录分配一个连续的 id 号,称为偏移量(Offset),用于唯一标识 Partition 内的记录。
Kafka 的客户端和 Broker 都会保存 Offset。客户端消费消息后,每隔一段时间,就把已消费的 Offset 提交给 Kafka Broker,表示已消费。
在这个过程中,如果客户端应用消费消息后,因为宕机、重启等情况而没有提交已消费的 Offset 。当系统恢复后,会继续消费消息,由于 Offset 未提交,就会出现重复消费的问题。
(2)重复消息通用解决方案
处理重复消息 = “业务幂等为基础,缓存/DB 去重为辅助,监控兜底保万一”。核心思路是 幂等性设计 + 去重机制,确保即使消息被多次消费,业务结果也不会出错。
业务层幂等设计
唯一标识:每条消息携带唯一业务 ID(如订单号、支付流水号),处理前先查库判断是否已执行。
SELECT status FROM orders WHERE order_id = '123'; -- 若已处理则直接跳过状态机控制:业务状态严格流转(如「已支付」订单不允许重复扣款)。
去重表/缓存
数据库去重表:消费前先
INSERT唯一键(消息 ID),利用主键冲突避免重复处理。INSERT INTO message_processed(id) VALUES ('msg_123') ON DUPLICATE KEY IGNORE;Redis 去重:用
SETNX设置消息 ID 过期时间(适合高频场景)。SETNX msg_123 1 EX 3600 # 1 小时内不重复处理
(3)主流 MQ 处理
| 消息队列 | 重复触发场景 | 推荐方案 |
|---|---|---|
| Kafka | 消费者重启导致 offset 回滚 | 业务幂等 + 本地 offset 持久化 |
| RabbitMQ | 未 ACK 导致重新入队 | 手动 ACK + 死信队列监控 |
| RocketMQ | 消息重试机制(16 次后进死信) | 消费日志 + 人工干预 |
极端情况兜底:
- 对账系统:定时扫描业务数据与消息记录,修复不一致(如定时补发短信)。
- 人工告警:监控重复消息频率(如 1 分钟同消息 ID 出现 3 次以上则报警)。
(4)方案选型
- 低频业务:数据库唯一索引(简单可靠)。
- 高频业务:Redis + 过期时间(高性能)。
- 金融级场景:幂等 + 对账 + 人工审核(强一致)。
🔬 扩展知识
【L3】至少一次 + 消费端幂等 vs 中间件级精确一次
详情
| 方案 | 实现复杂度 | 性能代价 | 适用边界 |
|---|---|---|---|
| 至少一次投递 + 消费端幂等 | 低 | 幂等校验增加一次缓存/DB 查询,约 1~5ms | 绝大多数业务(推荐默认) |
| 中间件级精确一次(生产去重 + 事务) | 高 | 吞吐下降约 20%~30%,且只覆盖“生产→存储”段 | 消息源单一、下游可事务化写入 |
关键认知:端到端的“不重复处理”永远靠消费端幂等实现,中间件的去重只能消除部分重复来源,无法消除消费端重试产生的重复。
失效场景(去重保障何时失效):
- Redis 去重的过期窗口陷阱:去重键 TTL 设为 1 小时,但消费端在故障 2 小时后重放,消息再次通过校验 → 重复生效。
- “先查后插”的竞态:两个消费者并发查到“未处理”,同时执行业务,唯一键未覆盖业务动作本身 → 重复扣款。正确做法是把唯一键约束放进业务事务(如插入去重记录与业务写入同一本地事务)。
- 消息 ID 不稳定:生产者重试时重新构造消息体导致 ID 变化,去重表形同虚设;必须用业务唯一键(订单号、支付流水号)而非传输层 ID。
- 脑裂/主从切换:旧主已写入未同步的消息在新主上重新暴露,产生“服务端级”重复,客户端去重是唯一防线。
【L4】Kafka 幂等生产者的去重边界
详情
enable.idempotence=true通过 PID + SequenceNumber 让 Broker 对同一生产者的重试消息去重,但只覆盖“单分区、单会话”内的生产重复。- 跨会话(生产者重启后 PID 变化)、跨分区的重复仍需事务或消费端幂等解决。
【L3】Kafka 定制深挖:三种语义与幂等/事务选型
详情
消息传递的服务质量标准有三种(从低到高):At most once(至多一次,允许丢,适合监控类)、At least once(至少一次,不丢但可能重复,主流 MQ 默认)、Exactly once(恰好一次,最高等级)。绝大部分 MQ(Kafka/RocketMQ/RabbitMQ)默认都是 At least once。
幂等生产者 vs 事务
| 机制 | 去重范围 | 原理 | 代价 | 适用边界 |
|---|---|---|---|---|
幂等生产者(enable.idempotence=true) | 单分区、单会话 | Producer ID(PID)+ Sequence Number,Broker 按 <PID, 分区, SeqNum> 去重 | 吞吐下降约 5%~10%,强制 acks=all | 单分区写入的重试去重(推荐默认开启) |
事务(transactional.id) | 跨分区、跨会话、含 offset | 两阶段提交:initTransactions → beginTransaction → send/sendOffsetsToTransaction → commitTransaction | 吞吐下降约 20%~30% | consume-transform-produce 的端到端 exactly-once |
关键边界:幂等只解决“生产者重试导致的重复”,不解决消费端重复;事务的 exactly-once 也仅限“Kafka→Kafka”链路(如 Flink/Kafka Streams),一旦下游是外部 DB,仍需业务幂等。
失效场景
- PID 换新:Producer 重启后获得新 PID,Broker 无法跨会话去重,滚动发布前后“已写入但回调未收到”的消息被新进程重发仍会重复——必须在消息体携带业务唯一键由下游幂等收口。
transactional.id误配:多实例共用同一transactional.id会互相 fence(ProducerFencedException)并产生悬挂事务;必须实例级唯一(如risk-tx-{instanceId}),消费者统一isolation.level=read_committed。- 事务超时:
transaction.timeout.ms(默认 15 分钟)内未提交的事务被 Broker 主动 abort,若业务实际已成功则产生“消息丢失”假象。
量化代价(推演示例):幂等开销每条消息额外 12 字节头部(PID 8B + SeqNum 4B);事务两阶段提交使端到端延迟增加约 1 个 RTT。
【L3】RocketMQ 定制深挖:四条重复路径与 5s 位点窗口
详情
重复路径
| 重复路径 | 触发条件 |
|---|---|
| Producer 重发 | 发送超时但 Broker 实际已写成功,重试再写一份(内容相同、msgId 不同) |
| 消费再平衡 | 实例上下线触发 Rebalance,offset 定期批量提交(默认 5s),窗口内消息被重复拉取 |
| 消费重试 | 返回 RECONSUME_LATER 后整批重投,业务可能已部分成功 |
| 人工重置位点 | 运维回溯消费,时间范围内消息全量重复 |
关键 offset 语义:集群模式下 ConsumeQueue 的消费位点是定期提交(默认 5 秒)而非逐条提交——即使每条都消费成功,进程崩溃时仍有最多 5 秒的重复窗口。逐条提交会成为吞吐瓶颈,这是“用可预期的重复窗口换吞吐”的设计取舍。
去重键陷阱:用 msgId 去重不可靠——超时重发会产生新 msgId,消息进重试 Topic 后存储位置也会变化;必须使用业务唯一键(如 orderId + 操作类型)。
两层防线:Redis SETNX 前置初筛(十万级 TPS,但主从切换瞬间标记可能丢失)+ DB 唯一键最终兜底(强一致,受 DB 写入能力约束,单机数千 TPS 量级);DB 层必须与业务写入同一事务,否则必然出现“标记成功业务失败(消息被误判已消费而丢失)”或“业务成功标记失败(重复消费)”两种撕裂。
【L3】RabbitMQ 定制深挖:四条重复路径与去重实现
详情
四条重复路径(其中 requeue 与主从切换最容易被忽视):
- 消费者未及时 ACK:处理完但 ACK 丢失(网络/崩溃),Broker 认为未消费重新投递。
- requeue 重复:
basicNack/basicReject带requeue=true时消息重新入队,若失败发生在“业务已部分成功”之后,再投递就是重复执行;且 requeue 后 deliveryTag 变化,无法用 deliveryTag 去重。 - 生产者重试:发送后未收到 Confirm 触发重试,Broker 实际已接收第一份,产生两条内容相同的消息。
- 镜像队列主从切换:主节点宕机时,已投递但未同步的 ACK 状态丢失,从节点提升为主后重新投递已消费消息。
幂等方案补充(更新类业务)
| 方案 | 实现方式 | 适用场景 |
|---|---|---|
| 乐观锁(版本号) | 更新时带 version 条件,重复更新影响行数为 0 | 更新类业务 |
| 状态机控制 | 业务状态严格流转,如“已支付”订单不允许重复扣款 | 有明确状态流转的业务 |
Redis 去重消费(Java)
public void consume(String messageId, String body) {
// 1. Redis 去重检查
Boolean isNew = redis.opsForValue().setIfAbsent(
"msg:processed:" + messageId, "1", Duration.ofHours(1));
if (Boolean.FALSE.equals(isNew)) {
channel.basicAck(deliveryTag, false); // 已处理,直接 ACK
return;
}
// 2. 执行业务逻辑
try {
processBusiness(body);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
redis.delete("msg:processed:" + messageId); // 失败删标记允许重试
channel.basicNack(deliveryTag, false, true);
}
}隐藏坑:
redis.delete与basicNack之间若进程崩溃,标记已删但消息未 requeue;更稳妥的做法是失败时保留标记并转死信队列人工处理,或用 DB 事务内写业务 + 写标记原子化。
幂等标记存储位置决定资金安全:标记若只存 Redis,Redis 与 MQ 同时故障时存在窗口;与业务在同一 DB 事务中提交时,切主后的重复消息会被唯一键挡住。金融级场景幂等必须落在强一致存储(DB)上,Redis 只能作初筛。
🏭 实战场景
详情
量化数据
- 重复率的典型水平:自动提交位点 + 至少一次语义下,每次消费者重启的重复窗口 ≈ 提交间隔(如 5 秒),约等于该窗口内消息量的 1%~3%。
- Redis 去重性能:单实例 SETNX 约 10 万 QPS,单次校验耗时约 1ms;DB 唯一键去重约 1 万 TPS,但强一致。
- 建议:去重键 TTL ≥ 消费端最大故障恢复窗口的 2 倍(如恢复窗口 1 小时则设 2 小时以上)。
踩坑案例(推演示例):重复消费引发的批量重复扣款
- 现象:晚高峰 20 分钟内,800 余笔用户被重复扣款,支付风控告警“同流水多笔成功”。
- 排查:消费日志显示同一批消息 ID 出现两次消费记录,间隔约 5 秒;期间消费者组发生过一次滚动发布。
- 根因:消费者配置为自动定时提交位点(间隔 5 秒),发布重启时“已消费、未提交”的批次被新实例重新投递;而扣款动作只校验了订单状态,未校验支付流水号唯一性。
- 修复:切换为处理成功后手动提交位点;扣款接口增加支付流水号唯一键,重复请求直接返回首次结果;存量重复扣款通过流水比对脚本批量原路退回,2 小时内全部追平。
场景题
场景:库存服务消费“订单取消释放库存”消息时,偶发出现库存被重复加回,导致超卖投诉。你如何排查与修复?
应急处理:立即对受影响 SKU 冻结自动回补,按订单表反查“已取消且已释放库存”的订单,核对释放次数;对确认重复释放的 SKU 扣减多余库存,防止继续超卖。
根因分析:拉取释放库存消息的 ID 与消费日志比对,发现同一消息被消费 2 次。检查消费端发现:释放库存与提交位点不在同一事务,释放成功后提交位点超时,触发重新投递;且释放接口用“库存 += 数量”的增量写法,天然不幂等。
长期方案:释放接口改为基于订单号的幂等操作(释放前校验该订单是否已释放,用唯一键约束防并发);位点提交与业务结果绑定,失败则重试整批;上线“订单取消数 vs 库存释放数”的 T+1 对账。
权衡:幂等校验引入每条消息约 2ms 的额外查询开销,但释放库存属于低频操作,吞吐余量充足;用微小性能代价消除了超卖这一高危资损路径。
⚠️ 常见误区
详情
常见误区:
- ❌ “开启 Kafka 幂等生产者就再也不会重复消费” → 幂等生产者只消除生产端重试重复,消费端重启/位点回滚导致的重复仍在。
- ❌ “用消息 ID 做去重键就够了” → 生产者重试可能重新构造消息导致 ID 变化,应以业务唯一键(订单号/流水号)为准。
- ❌ “先查后插能防并发重复” → 两个消费者可能同时查到未处理,必须用唯一键约束+事务内插入。
🔀 发散问题
- 为什么消息重复不可能在中间件层彻底根除,必须由消费端幂等兜底?
重复的根源之一是消费端“处理成功但确认失败”后的重新投递,这个窗口发生在中间件视野之外,任何 Broker 都无法感知下游业务是否真正执行成功。因此精确一次语义只能收敛重复来源,最终防线永远是消费端的幂等设计。 - 去重表和业务表分库时,如何避免“去重成功但业务失败”?
把去重记录的插入与业务数据写入放在同一个本地事务中(同库同事务),利用唯一键冲突作为并发判重手段;若必须分库,则采用“先落去重记录 → 再执行业务 → 失败则删除/标记去重记录”的补偿流程,并接受短暂的重复窗口。 - Redis 去重和数据库唯一键去重如何选型?
高频低价值消息(日志、埋点)用 Redis,追求 QPS 与低延迟,接受极端情况下的漏判;资金与订单类用 DB 唯一键,强一致、可审计;金融场景常见组合是 Redis 前置快速拦截 + DB 唯一键事务内兜底。
【困难】如何保证 MQ 消息的顺序性?⭐⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:20 min | 🏷 标签:MQ / 顺序性
💎 关键结论
保序的核心是“同一业务 ID 锁定同一分片 + 单线程消费”,从生产、存储、消费三端同时控制;绝大多数业务只需局部有序,全局有序是不必要的性能自残。理由:MQ 只在分片内保序,跨分片乱序是并行扩展的必然代价。
⚡记忆卡片
口诀:同键同分片,单线程消费,失败即暂停
关键词:分片键/哈希路由/max.in.flight/MessageListenerOrderly/pause()
链路:同业务 ID 哈希到同一分片 → 生产端限制未确认请求防重试乱序 → 消费端单线程串行处理 → 失败暂停分片避免跳过
📖 核心知识
要保证 MQ 消息的顺序性,需从 生产、存储、消费 三个环节控制。核心思路是:“同一业务 ID 锁定同一队列 + 单线程消费”,需结合业务需求选择局部顺序或全局顺序方案。

(1)生产端保序
单生产者+单线程发送:同一业务 ID(如订单 ID)的消息由 同一生产者线程 顺序发送,避免多线程并发乱序。
// 示例:相同 orderId 的消息由同一线程发送 mqProducer.send(msg, orderId); // 通过 hash 选择分区禁用异步发送重试:异步发送失败时可能乱序,需同步发送或关闭重试(如 Kafka 配置
max.in.flight.requests.per.connection=1)。
(2)MQ 服务端保序
单分区/队列有序:
Kafka/RocketMQ:同一业务 ID 的消息发送到 同一分区(Partition)。
// 根据 orderId 哈希选择分区 int partition = orderId.hashCode() % partitionNum;RabbitMQ:使用单队列(或一致性哈希交换器绑定唯一队列)。
关闭分区/队列并行:避免服务端多分区/多副本间的顺序混乱(如 Kafka 的
unclean.leader.election.enable=false)。
(3)消费端保序
- 单消费者串行消费:
- 同一队列/分区由 单消费者线程 处理(如 Kafka 单线程消费或
max.poll.records=1)。 - 多消费者时,相同业务 ID 的消息路由到同一消费者(如 RocketMQ 的
MessageQueueSelector)。
- 同一队列/分区由 单消费者线程 处理(如 Kafka 单线程消费或
- 内存队列排序(复杂场景):消费者拉取消息后,按业务 ID 分组存入内存队列,由不同线程分别串行处理。
(4)特殊场景处理
- 全局严格顺序:牺牲性能,全链路单线程(生产→MQ→消费),仅适合低吞吐场景(如 Binlog 同步)。
- 局部顺序:仅保证同一业务 ID 的顺序(如订单的创建→支付→退款),允许不同订单并发。
(5)主流 MQ 处理
| Kafka | RocketMQ | |
|---|---|---|
| 单分区/队列有序 | 单分区追加写入,天然有序 | 单队列有序 |
| 哈希路由 | 采用哈希路由,相同 key 固定发往同一分区 如果不指定 key,则采用轮询方式选择分区 | 使用 MessageQueueSelector 指定队列 |
| 单线程生产 | max.in.flight.requests.per.connection=1 控制 | |
| 单线程消费 | 确保一个分区仅由一个消费者线程处理(Kafka 默认规则) | 使用 MessageListenerOrderly |
| 业务保证 | 避免多线程并发消费同一队列 分布式锁(如 Redis) 关键业务操作加锁,防止并发执行乱序 | 避免多线程并发消费同一队列 分布式锁(如 Redis) 关键业务操作加锁,防止并发执行乱序 |
(6)注意事项
- 性能权衡:顺序性越高,并发性能越低(需根据业务容忍度平衡)。
- 错误处理:消费失败时需暂停当前分区消费(如 Kafka 的
pause()),避免跳过消息导致乱序。 - 监控:定期检查消息积压和顺序偏移(如 Kafka 的
consumer.position())。
🔬 扩展知识
【L3】全局有序 vs 局部有序
详情
| 方案 | 吞吐能力 | 适用边界 |
|---|---|---|
| 全局有序(单分片 + 单消费者) | 极低,通常 < 5000 TPS | Binlog 同步、强序审计日志等低吞吐场景 |
| 局部有序(业务 ID 哈希到分片) | 高,随分片数线性扩展 | 订单、账户状态流转等绝大多数业务(推荐) |
| 无序(轮询/粘性分发) | 最高 | 日志采集、指标上报 |
结论:先问业务到底需要哪一层有序。绝大多数业务只需要同一业务实体内有序,全局有序是不必要的性能自残。
失效场景(有序保障何时被击穿):
- 生产端异步重试乱序:消息 A 发送超时后重试,但消息 B 已成功落队,A 重试成功后反而排在 B 后面。关闭重试会丢消息,限制未确认请求数或启用生产端保序配置才能兼顾。
- 分片数变更:扩分片后哈希取模结果变化,同一业务 ID 的新老消息路由到不同分片,跨分片无序;扩容必须在业务低峰配合双写或停止切换。
- 消费端并行处理:为提升吞吐把拉取到的消息丢进线程池并发处理,顺序立即破坏;若确需并行,必须按业务 ID 二次分发到固定线程。
- 主从切换/非同步副本上位:服务端故障切换可能丢失尾部消息或暴露乱序尾部,见本文档『如何保证 MQ 消息不丢失?』的失效路径。
- 失败消息重试队列旁路:失败消息转入重试/死信后再回注,时间上必然晚于后续消息,破坏原始顺序。
【L4】Kafka 生产端保序配置的演进
详情
- 旧做法:
max.in.flight.requests.per.connection=1严格保序但牺牲约一半发送吞吐。 - 新做法:开启幂等生产者(Kafka 0.11+)后,允许单连接最多 5 个未确认请求仍保证分区内有序(Broker 按序列号排序缓存)。
【L3】Kafka 定制深挖:max.in.flight 与保序的深层关系
详情
max.in.flight.requests.per.connection 的三种配置
| 配置 | 保序能力 | 吞吐 | 适用边界 |
|---|---|---|---|
=1(配合任意 acks) | 严格保序 | 最低,降约 50%(无流水线) | 旧版本客户端 + 严格顺序场景 |
≤5 + enable.idempotence=true(0.11+) | 保序且高吞吐 | 接近默认水平 | 推荐默认 |
>1 且未开幂等 | 不保序 | 最高 | 无序场景 |
原理:该参数控制单连接上未确认请求的并发数。多请求在飞时,若先发批次失败重试、后发批次先落盘,日志顺序即被打乱。幂等生产者通过 <PID, 分区, SeqNum> 让 Broker 拒收乱序批次(OutOfOrderSequenceException 并断连);上限是 5 因为 Broker 的去重/排序窗口容量有限,这是“保序窗口大小”与“流水线深度”的工程折中。
分区内有序也会被击穿的场景
- 未开幂等 + 多在飞请求 + 重试:批次 1 超时重试,批次 2 先写入,重试成功的批次 1 排在批次 2 之后。
- 分区数扩容:哈希取模结果变化,N → 2N 分区时约 50% 的 key 会更换分区;扩分区必须在低峰期配合双写过渡。
- Leader 切换 + unclean 选举:非同步副本上位后尾部消息丢失,衔接处出现“逻辑乱序/缺失”。
- 消费端单分区多线程:Kafka 保证“一个分区只给组内一个消费者”,但消费者内部把 poll 到的批次丢线程池处理,顺序照样破坏。
典型有序场景:金融交易(开户→存款→转账)、日志聚合(启动→运行→异常→终止)、库存管理(入库→出库→盘点)、流媒体(I 帧→P 帧→B 帧)。
【L3】RocketMQ 定制深挖:MessageQueueSelector 与队列锁
详情
生产端:哈希选队列 + 同步发送
producer.send(msg, (mqs, m, arg) -> {
int index = Math.abs(arg.hashCode()) % mqs.size();
return mqs.get(index);
}, orderId); // 同步发送 + 按 orderId 选队列必须用同步发送;某 Broker 不可用时故障规避会把消息发到其他队列导致乱序,顺序消息建议关闭发送重试(retryTimesWhenSendFailed=0),失败交由业务侧重发——重发到别的队列比失败更糟糕。
消费端:MessageListenerOrderly 双重锁串行
- Broker 端队列锁(默认 60s 租约,定期续租)+ 消费端本地锁双重保证同一队列只被一个线程消费;主从切换或重启时锁可能短暂不一致,因此业务幂等/状态机仍是必要双保险。
- 失败返回
SUSPEND_CURRENT_QUEUE_A_MOMENT暂停当前队列(默认 1 秒)后本地原地重试,不进重试 Topic——进重试 Topic 意味着失败消息延迟后从另一个队列回来必然插队破坏顺序;代价是一条毒消息卡死整个队列,需设maxReconsumeTimes超限进死信并告警。
顺序失效四场景:① 队列数变更(哈希取模结果变化,顺序 Topic 队列数禁止在线变更,扩容需新建 Topic 双写迁移);② Broker 故障规避切队列;③ 毒消息反复 SUSPEND 挂起,后续消息延迟数小时;④ 误用 MessageListenerConcurrently 并发消费。
架构降维思路:很多“状态流转”场景其实不需要严格顺序——用状态机 + 幂等容忍乱序(只允许单向状态迁移,重复/滞后的旧状态直接忽略),把顺序消费降级为并发消费,吞吐可提升一个数量级。
【L3】RabbitMQ 定制深挖:保序方案选型与哈希分流
详情
乱序来源:多消费者竞争同一队列(轮询分发、处理速度不同)、requeue=true 后消息被排到队列尾部而非原位置、镜像队列主从切换。
保序方案对比
| 方案 | 原理 | 缺点 | 适用场景 |
|---|---|---|---|
| 单队列单消费者 | 一个队列只配一个消费者,串行处理 | 吞吐低,单点瓶颈 | 消息量极小、绝对严格顺序 |
| 多队列分流(按业务 ID) | 生产者按业务 ID 哈希,每队列单消费者 | 需提前规划队列数 | 推荐。同一业务 ID 有序的场景 |
| 一致性哈希交换机 | rabbitmq_consistent_hash_exchange 插件按 Key 哈希分配队列 | 需安装插件,配置稍复杂 | 大规模有序消费、需动态扩缩队列 |
| 消息组(single active consumer) | 3.8+ 的 x-single-active-consumer + group_id | 仅限 Quorum Queue,配置复杂 | 高可用 + 顺序消费 |
x-single-active-consumer:同一队列多消费者中只有一个活跃,其余热备;活跃消费者宕机后 Broker 自动切换,既保留单消费者串行保序,又解决单点问题,是“高可用 + 保序”的官方方案。
按 orderId 哈希分流(Java)
// 1. 生产者:按 orderId 哈希选择队列
int queueNum = Math.abs(orderId.hashCode()) % QUEUE_COUNT;
channel.basicPublish("", "order_queue_" + queueNum, null, message.getBytes());
// 2. 为每个队列启动一个消费者,串行消费(autoAck=false)
for (int i = 0; i < QUEUE_COUNT; i++) {
channel.basicConsume("order_queue_" + i, false, consumer);
}成本模型:保序的本质代价是把并行度从“消费者数”降为“队列数”:吞吐上限 ≈ 队列数 × 单队列串行吞吐;哈希键选得越细,并行度越高(队列数常见规划 10~100)。
🏭 实战场景
详情
量化数据
- 全局有序的吞吐上限:单分片串行消费通常在 1000~5000 TPS(取决于单条处理耗时),每增加一个分片并行度线性翻倍。
- 消费端二次分发方案:按业务 ID 哈希到 N 个内存队列,N 取 CPU 核数 × 2 较常见,可在保序前提下把消费吞吐提升一个数量级。
- 扩分片导致的乱序窗口:若不做双写过渡,乱序比例 ≈ 扩分片后新路由命中的消息占比(如 8→16 分区约 50% 消息换分片)。
踩坑案例(推演示例):订单状态回退的乱序事故
- 现象:客服反馈订单显示“已发货”后又变回“已支付”,用户投诉物流信息消失。
- 排查:按订单号拉取该订单全部消息,发现“创建→支付→发货”三条消息中,“支付”消息的消费时间反而晚于“发货”;检查发现该 Topic 消费端为提升吞吐引入了线程池并发处理。
- 根因:消费端多线程乱序执行;且下游状态机未做“旧状态不得覆盖新状态”的防御,乱序直接落库。
- 修复:消费端改为按订单号哈希分发到固定内存队列、每队列单线程串行;下游状态机增加状态版本号校验(旧版本直接丢弃);对历史回退订单批量修复。
场景题
场景:支付回调消息乱序导致订单状态回退:部分订单先收到“支付成功”回调后又收到“已受理”回调,订单状态被回退,用户投诉。你如何排查和修复?
应急处理:立即在状态机层增加“禁止向低版本状态流转”的防御开关,先止血;对已回退的订单按支付网关流水批量订正状态。
根因分析:按订单号比对消息到达顺序与消费时间,定位乱序发生在消费端:回调消息被丢进线程池并发处理,不同线程完成时间不确定;且生产端回调重试也放大了乱序。本质是“同订单消息未锁定同一串行通道”。
长期方案:生产端以订单号为分片键保证同单消息同分片;消费端按订单号二次哈希分发到固定线程串行处理;状态机引入状态版本号/时间戳,旧状态覆盖新状态时直接拒绝并告警;对回调消息做幂等(同回调号只生效一次)。
权衡:串行化使单订单链路的处理吞吐受限于单线程,但支付回调 TPS 远低于上限,而状态正确性是不可妥协的红线;用可量化的吞吐余量换取零乱序是合理决策。
⚠️ 常见误区
详情
常见误区:
- ❌ “指定了 key 发消息就一定有序” → 若生产端开启异步重试且不限制未确认请求数,重试成功时可能排在后续消息之后。
- ❌ “消费失败跳过这条继续消费就行” → 跳过会破坏顺序,正确做法是暂停分片(pause)或阻塞重试,配合死信旁路。
- ❌ “扩分片对有序业务无影响” → 扩分片改变哈希路由,新老消息跨分片即乱序,需双写过渡。
🔀 发散问题
- 为什么 MQ 只能保证分片内有序,无法原生保证全局有序?这是设计缺陷吗?
不是缺陷而是取舍:全局有序意味着写入和消费都只能单点串行,吞吐永远无法水平扩展,与 MQ 的核心价值(并行与堆积)直接冲突。因此所有主流 MQ 都把有序粒度下放到分片,把“哪些消息需要互序”的决策权交给业务(通过分片键表达)。 - 生产端“发送失败重试”与“保序”天然冲突,工程上怎么解?
三条路:一是限制单连接未确认请求数为 1,牺牲约一半发送吞吐换严格有序;二是启用生产端幂等保序能力(新版客户端支持多未确认请求下仍保序);三是业务容忍“失败消息进异常队列不重试”,用旁路补偿替代在队重试。 - 消费端既要保序又要高吞吐,常见的架构模式是什么?
“拉取线程 + 二次分发”模式:拉取线程按分片顺序取消息,按业务 ID 哈希到固定数量的内存队列,每个队列由单线程串行消费;既保证同 ID 有序,又获得与队列数成正比的并行度。需注意队列满时的背压与进程崩溃后的重复消费(配合幂等)。
【困难】如何处理 MQ 消息积压?⭐⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:20 min | 🏷 标签:MQ / 消息积压
💎 关键结论
积压治理 = 快速消化存量 + 控制增量:短期扩容+降级恢复服务,长期优化消费+自动化预防;前提是消化速率必须大于生产速率,否则一切扩容无效。理由:积压本质是生产与消费速率失衡,先算清速率再选手段。
⚡记忆卡片
口诀:监控早发现,扩容扛流量,消费改批量,生产限流速
关键词:lag/扩消费者/临时 Topic/批量消费/限流/死信队列
链路:堆积告警发现积压 → 判断瓶颈在并发还是下游 → 扩消费者/临时 Topic 转发消化 → 生产端限流控增量 → 积压清零后回收资源
📖 核心知识
处理 MQ 消息积压的核心思路是 “快速消费存量+优化生产速率”,需结合监控、扩容、降级等手段综合治理:
- 短期:扩容+降级,优先恢复服务。
- 长期:优化消费逻辑+自动化运维,预防再次积压。
(1)快速消费积压消息
- 增加消费者实例:横向扩展消费者服务(如 Kubernetes 动态扩容 Pod),注意分区数限制(Kafka 需提前规划足够分区)。
- 提升消费并行度:
- 调整消费者并发参数(如 Kafka 的
max.poll.records、RabbitMQ 的prefetch_count)。 - 多线程消费(需保证无顺序要求的场景)。
- 调整消费者并发参数(如 Kafka 的
- 临时降级:非核心业务暂停消费(如日志处理),集中资源处理核心业务消息。
(2)优化消费能力
- 批量处理:合并多条消息一次处理(如数据库批量插入)。
- 异步化+削峰:消费者将消息存入内存队列,后台线程异步处理,避免同步阻塞。
- 跳过非关键逻辑:临时关闭日志记录、数据校验等非必要操作。
(3)控制生产端流量
- 限流:生产端启用速率限制(如 Kafka 的
quota、Redis 令牌桶)。 - 削峰填谷:消息先写入缓存层(如 Redis List),再匀速写入 MQ。
- 业务降级:高峰期关闭非核心功能的消息生产(如暂停推荐系统更新)。
(4)监控与告警
- 实时监控指标:队列堆积量(如 Kafka 的
lag)、消费速率(TPS)、消费者状态;设置阈值告警(如积压超过 10W 条触发短信通知)。 - 根因分析工具:日志分析(消费者卡顿、GC 问题)、链路追踪(如 SkyWalking 定位慢消费)。
(5)长期预防措施
- 容量规划:根据业务峰值预先扩容分区/队列(如 Kafka 分区数 = 消费者数 × 1.5)。
- 死信队列+重试机制:处理失败的消息转入死信队列,避免阻塞正常消费。
- 自动化扩缩容:基于积压指标动态调整消费者数量(如 K8s HPA)。
(6)主流 MQ 处理
| 消息队列 | 关键操作 |
|---|---|
| Kafka | 增加分区+消费者,调整 fetch.max.bytes |
| RabbitMQ | 镜像队列扩容,提高 prefetch_count |
| RocketMQ | 消费组扩容,启用定时消息延迟消费 |
🔬 扩展知识
【L3】常规扩容 vs 转发方案
详情
| 方案 | 见效速度 | 代价 | 适用边界 |
|---|---|---|---|
| 扩消费者实例 | 分钟级 | 受分片数上限约束 | 积压 < 分片并行度可消化量(推荐首选) |
| 新建临时 Topic + 分片扩容 N 倍 + 转发程序 | 半小时级 | 需写转发程序,事后要回收 | 分片数不足、存量巨大(千万级以上) |
| 提升单消费者并行度(批处理/多线程) | 立即 | 需改造消费逻辑,有乱序与重复风险 | 无顺序要求且消费逻辑存在批量优化空间 |
| 丢弃/跳过非核心消息 | 立即 | 数据丢失,需业务评估 | 日志、埋点等可弃数据 |
失效场景(扩容手段何时失效):
- 消费者数超过分片数:多出来的消费者空转无消息,扩容无效——这是最常见的“扩了没效果”。
- 瓶颈不在消费并发而在下游:消费逻辑里同步调用慢接口/写库,扩消费者只是把压力转移到下游,甚至把下游打挂。
- 积压期间仍在高速生产:消化速率 < 生产速率时,积压不降反升,必须先限流生产端。
- 转发方案的隐藏前提:转发程序自身吞吐也要够(批量拉取 + 批量写入),否则只是把积压换了个地方。
【L4】积压恢复耗时估算
详情
恢复时长 = 积压量 / (消化速率 - 新增生产速率)- 例:积压 3000 万条,消化 2800 TPS,新增生产 800 TPS → 30000000 / 2000 ≈ 15000s ≈ 4.2 小时。
- 经验值:消费端单实例吞吐通常在 1000~10000 TPS;扩容前先算“消化速率能否大于生产速率”,否则扩容无意义。
【L3】Kafka 定制深挖:分区硬约束与积压陷阱
详情
- 分区数是硬上限且只能增不能减:一个分区同一时刻只能被组内一个消费者线程消费,扩实例超过分区数纯浪费;合并分区需跨日志重排 offset,Kafka 不提供该能力,分区数必须按“未来 1~2 年峰值消费并发 × 1.5”一次性规划到位。
- 积压消费触发 offset 过期:消费者长期追不上时,若
offsets.retention.minutes(默认 7 天)到期且消费组无活跃成员,位点被清理,恢复后从auto.offset.reset位置重新消费,可能漏掉存量。 - 积压期间消息被保留策略清理:
log.retention.ms小于积压消化时长时,存量未消费即被删除——积压治理必须先核对保留期。 max.poll.interval.ms恶性循环:批处理变慢 → 超过阈值(默认 5 分钟)被踢 → rebalance → 重复消费更慢;处置:临时调大到 30 分钟,或减小max.poll.records(默认 500,积压期可调 1000~2000)使单批耗时回到阈值内。- 转发程序不成为瓶颈的要点:只做“拉取→写入”不做业务处理;拉取端调大
fetch.max.bytes/max.poll.records批量拉,写入端回调异步 + 攒批;转发实例数等于原分区数,吞吐可达正常消费的 10 倍以上。
【L3】RocketMQ 定制深挖:读写不对称与消费参数
详情
- 积压不伤写侧,瓶颈在读侧:写侧 CommitLog 顺序追加不受影响;读侧积压消息不在 Page Cache 时,消费按 ConsumeQueue 索引跳转 CommitLog 退化为磁盘随机读,IO util 打满后反过来影响同 Broker 上其他 Topic 的写入。缓解:扩内存扩大 Page Cache、临时 Topic 分散读压力、升级 SSD。
- 扩容上限是队列数:并发消费下实例数可扩到等于队列数(每队列一个),再多空转;
consumeThreadMax(默认 20/64)提升单实例并行度。 - 批量消费参数:
consumeMessageBatchMaxSize+pullBatchSize调大,合并多条一次入库(如 JDBC batch),IO 次数降一个数量级。 - 顺序消息积压不能简单扩容:顺序消费每队列单线程串行,实例超过队列数无效,提升速度只能靠优化单条处理耗时;大面积积压往往说明“用顺序消息解决乱序”的架构决策需要重审。
- 临时 Topic 转发的本质是“用写吞吐换读并行度”:桥接程序只做转发不做业务,把存量搬到更多队列后并行消化;新 Topic 消费仍需幂等,消化完成后要处理双 Topic 切换窗口。
【L3】RabbitMQ 定制深挖:架构短板与内存水位
详情
为什么 RabbitMQ 怕积压(架构根因)
- 内存优先的队列实现:消息默认在内存中排队,积压触发
vm_memory_high_watermark(默认 0.4)后 Broker 阻塞所有生产者连接,积压会反噬全集群。 - 页交换惩罚:持久化消息积压超出内存后换页到磁盘,消费时再换回,退化为随机 IO(积压百万级时消费速率常见下降 5~10 倍,经验值)。
- 与 Kafka 的差距:Kafka 分区独立顺序日志,积压只是顺序读磁盘;RabbitMQ 单队列是单进程串行模型,无法通过加磁盘线性恢复——RabbitMQ 的积压治理重在“预防”,Kafka 重在“快速消化”。
独有手段
- 惰性队列(Lazy Queue):
x-queue-mode=lazy消息直接落盘不占内存,解除内存告警与生产阻塞;代价是消费需逐条读盘,只适合存量慢消化,消化完应转回默认模式;新版本更推荐直接迁移仲裁队列 +delivery-limit。 - 容量与时效约束:
x-max-length限队列上限(超出丢弃最旧,配合 DLX 审计);x-message-ttl使过期消息自动丢弃;disk_free_limit(默认 50 MB)低于阈值全节点阻塞。 - 背压感知:生产者必须监听
connection.blocked/unblocked回调,阻塞期间停发或写本地缓冲/降级,恢复后限速重放。 - connection.blocked 是节点级全局的:一个队列的积压耗尽内存后,同节点所有队列的发布都会被阻塞,拆 vhost 只是逻辑隔离,真正的故障域隔离需节点级拆分;积压告警应设在水位触发前(如 50%)。
🏭 实战场景
详情
告警阈值建议:积压量超过“正常消费速率 × 30 分钟”即触发告警,留出人工介入窗口。
踩坑案例(推演示例):一次千万级积压的夜间事故
- 现象:凌晨风控规则变更引发下游接口大面积超时,消费 TPS 从 5 万跌到 2000,2 小时内积压飙升至 3000 万条,实时风控链路延迟超过 20 分钟。
- 排查:链路追踪定位到消费逻辑中同步调用的风控接口 P99 从 30ms 恶化到 2s;此时消费端已扩到等于分片数上限,继续扩实例无效。
- 根因:单条同步调用是真正瓶颈,而分片数限制了并发上限。
- 修复:立即新建分片数为原来 10 倍的临时 Topic,部署 10 倍并发度的临时消费者组做转发+消化;同时将风控调用改为批量接口(单批 100 条),单条耗时降至约 1ms;积压在 3 小时内消化完毕,事后回收临时 Topic 并回滚规则。
场景题
场景:大促开始后 10 分钟,订单消息积压从 0 飙到 2000 万条,消费 TPS 只有平时的 1/5,下游库存服务 CPU 正常但响应变慢。你如何处置?
应急处理:第一步确认下游库存服务的慢点(大概率是热点行锁或连接池耗尽),对该链路限流保命;第二步把消费者扩到分片数上限;第三步评估是否启动临时 Topic 转发预案。
根因分析:计算消化能力:假设消化速率提升到 1.5 万 TPS、大促生产速率 8000 TPS,净消化 7000 TPS,则 2000 万条需约 48 分钟。若仅靠扩消费者达不到该速率(分片不足),则必须走临时 Topic 方案(分片扩 10 倍,转发 + 消化速率可到 5 万 TPS 以上,约 7~10 分钟追平)。
长期方案:大促前压测定容,按“峰值 TPS × 1.5”预留分片与消费者容量;消费逻辑中下游调用全部批量化并配置熔断降级;建立积压大盘(lag、消化速率、预计清零时间三者联动展示)。
权衡:临时 Topic 方案需要约半小时的工程准备且事后要回收资源,但相比硬扛 48 分钟的业务损失(订单履约延迟),这是明确划算的取舍;预案必须提前演练,事故现场第一次写转发程序是不可接受的。
详情
各 MQ 定制踩坑案例(推演示例)
- RocketMQ:消费者假活导致积压——告警显示积压从 0 涨到 300 万且持续上涨,但实例全部在线无报错;线程 dump 发现消费线程池全部阻塞在一个下游接口的 30s 超时等待上。修复:下游超时降到 2s + 熔断;长期方案:积压告警与“消费 TPS 跌零”告警双监控,才能覆盖假活场景。
- RabbitMQ:内存水位告警引发全集群生产阻塞——下游停消费两小时后积压 800 万条,触发
memory resource limit,Broker 阻塞了全部连接的发布,无关业务连带受损。修复:临时消费者 + 转惰性模式降压;长期方案:关键队列配x-max-length+ DLX,业务队列按重要性节点级隔离,告警设在水位 50% 之前。
⚠️ 常见误区
详情
常见误区:
- ❌ “积压了就多起消费者” → 消费者数超过分片数后多出的实例空转;先确认分片数这个硬上限。
- ❌ “直接给原 Topic 扩分片消化存量” → 存量消息仍在旧分片,且扩分片改变哈希路由破坏同键有序;应新建临时 Topic 转发。
- ❌ “积压消化完就没事了” → 若生产速率仍大于消费速率,积压会再次发生;必须同时限流生产端或永久提升消费能力。
🔀 发散问题
- 为什么“增加消费者”不总是有效?分片数在积压治理中扮演什么角色?
主流 MQ 的分片是并行消费的最小单元,一个分片同一时刻只能被组内一个消费者占用,消费者数超过分片数即空转。因此分片数是消费并发的硬上限,必须在容量规划阶段按峰值预估(常见经验:分片数 ≥ 峰值消费并发 × 1.5)。 - 临时 Topic 转发方案的完整操作链路是什么?为什么它比直接扩分区更安全?
链路:新建 N 倍分片的临时 Topic → 原消费者改为“只拉取、不处理、批量转发”→ 临时消费者组全速消化 → 积压清零后回切并下线临时资源。相比直接给原 Topic 扩分片(扩分片后新消息按新哈希路由,会破坏同键有序,且存量消息仍留在旧分片),转发方案不动原 Topic 元数据,存量与增量都能被高速消化。 - 积压治理中,如何判断瓶颈在消费端还是下游系统?
看消费耗时分布:若单条消息处理耗时陡增(链路追踪可见慢调用),瓶颈在下游,扩消费者只会放大下游压力,正确做法是优化/降级下游(批量化、缓存、熔断非核心依赖);若单条耗时无变化但总吞吐上不去,瓶颈在并发度,扩容/提并行度才有效。
【中等】MQ 如何实现消息重放(重新消费)?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 消息重放
💎 关键结论
消息重放的本质是“回溯消费位点 + 幂等重处理”:重放前确认保留期,重放中幂等+限速,二者缺一不可。理由:重放会让消息被再次消费,不幂等就会产生脏数据;不限速就会冲垮下游。
⚡记忆卡片
口诀:回溯位点,幂等限速,新组隔离
关键词:reset-offsets/seek()/消费位点/保留期/新消费者组
链路:确认消息在保留期内 → 新建独立消费者组 → 按时间戳/offset 重置位点 → 限速重放+监控 lag → 幂等处理重复消息
📖 核心知识
消息重放指将历史消息重新投递给消费者再次消费,常用于:数据修复/补偿、消费逻辑变更后的重算、容灾迁移、数据同步到新的下游系统。
实现方式对比
| MQ | 重放机制 | 特点 |
|---|---|---|
| Kafka | 重置消费者组 offset:kafka-consumer-groups.sh --reset-offsets,或代码调用 seek()/seekToBeginning() | 按 offset/时间戳重置,需先停止消费组 |
| RocketMQ | 控制台或命令重置消费位点(按时间戳);也可新建消费组从头订阅 | 支持按时间重置,不影响其他消费组 |
| RabbitMQ | 无原生 offset 回溯能力,通常靠消息持久化 + 手动重发或归档到外部存储后重新投递 | 重放能力弱 |
关键要点
- 消息必须未过期:重放只能在消息保留期内进行(Kafka
retention.ms、RocketMQfileReservedTime),需提前规划保留时长。 - 幂等是前提:重放会导致消息被再次消费,消费端必须做幂等处理,否则产生重复数据。
- 隔离新消费组:重放建议新建独立消费者组/订阅关系,避免污染线上正常消费进度。
- 限速与监控:重放可能瞬间产生海量请求,需限流防止冲垮下游,并监控重放进度(消费位点/lag)。
- 重放 ≠ 补数:若消息已过保留期,只能依赖数据归档(如落 HDFS/Hive)或数据库反查补数。
一句话总结:消息重放的本质是“回溯消费位点 + 幂等重处理”,重放前保留期、重放中幂等限速,二者缺一不可。
🔬 扩展知识
【L3】重放与存储模型的关联
详情
- Kafka/RocketMQ 按保留策略存储、确认不删,天然支持按位点/时间戳回溯;RabbitMQ 确认即删,重放只能依赖发送端留存或外部归档。
- 因此选型时若业务有“数据回刷/重算”诉求(如实时数仓修数),应优先选择日志型 MQ。
【L4】分层存储与长期重放
详情
- 新一代 MQ(Kafka Tiered Storage、Pulsar 分层存储)把历史数据转存对象存储,保留期从“磁盘容量约束”变为“成本约束”,支持月级甚至更久的重放窗口。
- 大数据体系中更常见的做法是把 MQ 数据落地 HDFS/Hive 后离线重算,MQ 只保留短窗口。
📚 延伸阅读:Kafka 官方文档
🔀 发散问题
重放时为什么建议新建消费者组?
复用线上消费组会把位点拨回,导致正常业务消息被重复消费且进度被污染;新建消费组独立回溯、互不影响,重放完成后可直接下线。
消息已过保留期还能重放吗?
不能从 MQ 重放,只能靠外部手段:生产端留存的数据归档(HDFS/数仓)回刷,或从业务数据库反查重建数据。
【中等】什么是死信队列?如何设计死信队列的处理机制?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 可靠性
💎 关键结论
死信队列(DLQ)是 MQ 的异常处理容器,存放多次重试仍无法消费的消息。它的价值不在于"存"而在于"管":必须配套监控告警、人工分析和补偿重投,否则死信只是把丢消息从明处挪到暗处。
⚡记忆卡片
- 口诀:重试超限进死信,监控告警加补偿
- 关键词:DLQ / %DLQ% / 重试超限 / 监控告警 / 补偿重投
- 链路:消费失败进 %RETRY% 递延重试 → 默认 16 次仍失败 → 转入 %DLQ%消费者组 → 触发告警 → 人工分析后修复重投或归档
📖 核心知识
**死信队列(Dead Letter Queue, DLQ)**是一种用于存储无法被正常消费的消息的队列,本质是 MQ 的异常处理容器。
产生死信的场景
- 消费失败:消费者多次重试后仍失败,消息被投递至 DLQ。
- 消息过期:设置了 TTL 的消息超时后未被消费,转为死信(部分 MQ 支持)。
- 队列超限:队列长度或存储容量达到上限,后续消息被拒绝或丢弃时转入 DLQ。
- 消息被拒绝:消费者明确拒绝(reject)且不重新入队。
死信队列设计要点
- 存储隔离:DLQ 与业务队列分开,避免影响正常消息处理。
- 监控告警:监控 DLQ 堆积量,设置阈值告警,及时发现异常。
- 重新处理机制:
- 手动处理:运维人员通过管理工具查看、导出死信,分析失败原因(如数据格式错误、业务逻辑 bug)。
- 自动重试:定时任务重新投递死信到原队列,或经修复后投递。
- 人工介入:对于需要人工干预的异常(如脏数据),提供处理界面。
- 死信保留策略:设定死信最大存储时间或容量,避免无限堆积。
挑战和应对
- 避免死信堆积:设计合理的重试策略和告警,及时处理。
- 死信风暴:若大量消息同时成为死信,可能压垮 DLQ,需对 DLQ 本身进行限流或容量规划。
- 死信分析:需携带原始消息头、异常堆栈等信息,便于定位问题。
🔬 扩展知识
【L3】在 RocketMQ 中,死信队列是按消费组隔离的:某消费组重试超限的消息进入名为 %DLQ%消费者组名 的 Topic;该 Topic 默认读权限为 2(只写不读),需要人工或工具显式订阅处理,否则消息会静默躺在那里直到过期。
详情
- RocketMQ 并发消费默认重试 16 次(间隔从 10s 递增至 2h),全部失败后自动转入死信;广播消费不支持重试,也不会产生死信。
【L3】在 RabbitMQ 中,死信由死信交换机(DLX)承接:消息因被拒绝(reject/nack 不重入队)、TTL 过期或队列超限成为死信后,按 x-dead-letter-exchange 配置路由到死信队列,见《RabbitMQ面试》『RabbitMQ 中消息什么时候会进入死信交换机?』。
【L4】死信治理的工程化做法:把 DLQ 消费做成独立补偿服务,按异常类型路由(格式错误→人工审核、下游超时→自动重投),并用消息轨迹串联原始消息上下文。
🔀 发散问题
Q:死信队列里的消息会自动重新消费吗?
A:不会。RocketMQ 的死信 Topic 默认不可被正常订阅消费,必须人工介入或通过工具重新投递到原 Topic。
Q:怎么快速发现死信堆积?
A:监控 %DLQ% Topic 的堆积量并接入告警,同时开启消息轨迹,见《RocketMQ面试》『RocketMQ 的消息轨迹如何启用?』。
MQ 事务
【困难】MQ 如何实现分布式事务?⭐⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:20 min | 🏷 标签:MQ / 分布式事务
💎 关键结论
MQ 分布式事务的核心是“本地事务 + 消息可靠投递 + 补偿”,用最终一致性替代强一致性;无论哪种方案,消费端幂等都是方案成立的前提而非可选项。理由:2PC/XA 在跨服务场景因长锁与单点基本不可用,而 MQ 的所有投递路径都天然存在重发。
⚡记忆卡片
口诀:本地事务+消息表/半消息,回查补投递,幂等是前提
关键词:本地消息表/事务消息/半消息/事务回查/最大努力通知/Saga
链路:业务与消息写入同一本地事务 → 扫描/回查驱动投递 MQ → 下游消费并幂等处理 → 对账兜底修复不一致
📖 核心知识
分布式事务是分布式系统中的核心难题。MQ 实现分布式事务的核心思路是 “本地事务 + 消息表 + 补偿机制”,即通过 最终一致性 替代强一致性,避免分布式锁带来的性能损耗。
(1)为什么需要 MQ 事务
在微服务架构下,一个业务操作往往需要跨多个服务完成。例如,电商下单需要同时操作订单服务、库存服务、支付服务。如果使用传统的 2PC/XA 强一致性事务,会带来性能差、可用性低、实现复杂等问题。
MQ 事务的本质是 将大事务拆分为多个本地事务,通过消息驱动各服务,最终达到数据一致性。
(2)MQ 事务的核心方案
| 方案 | 核心思路 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| 本地消息表 | 本地事务 + 消息表(同一事务写入业务数据和消息记录),定时任务扫描消息表投递 MQ | 实现简单,可靠性高 | 定时扫描有延迟,消息表需维护 | 大多数业务场景(推荐) |
| 事务消息(RocketMQ) | 半消息 → 本地事务 → 提交/回滚 → 事务回查 | 无需消息表,原生支持 | 仅 RocketMQ 原生支持,实现较复杂 | RocketMQ 生态,金融级场景 |
| 最大努力通知 | 消息投递后,通过多次通知 + 对账兜底 | 实现简单 | 一致性最弱,可能需要人工介入 | 对一致性要求不高的场景 |
| Saga 模式 | 将长事务拆分为多个本地事务,每步有对应的补偿操作 | 灵活,无全局锁 | 补偿逻辑复杂,可能产生脏读 | 长流程业务 |
方案一:本地消息表(通用方案)
核心流程:
- 业务操作 + 记录消息:在同一个本地事务中,写入业务数据,并插入一条消息记录到
message_table(状态为待发送)。 - 定时任务扫描:后台定时任务扫描
message_table中状态为待发送的记录,投递到 MQ。 - 投递成功更新状态:MQ 投递成功后,将消息状态更新为
已发送。 - 消费者消费:下游服务消费消息,处理业务逻辑。
- 幂等保证:消费者必须实现幂等性,防止消息重复消费。
-- 本地消息表示例
CREATE TABLE message_table (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
business_id VARCHAR(64) NOT NULL COMMENT '业务ID,如订单号',
message_content TEXT NOT NULL COMMENT '消息内容',
status TINYINT DEFAULT 0 COMMENT '0:待发送 1:已发送 2:已消费',
retry_count INT DEFAULT 0 COMMENT '重试次数',
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_business_id (business_id) -- 保证幂等
);方案二:RocketMQ 事务消息
核心流程:
- 发送半消息:生产者先发送一条“半消息”(Half Message)到 Broker,此时消息对消费者不可见。
- 执行本地事务:半消息发送成功后,执行本地业务事务。
- 提交/回滚:
- 本地事务成功 → 发送
Commit,消息变为可消费。 - 本地事务失败 → 发送
Rollback,消息被删除。
- 本地事务成功 → 发送
- 事务回查:若 Broker 长时间未收到
Commit/Rollback,会主动回查生产者的本地事务状态(通过checkLocalTransaction方法)。 - 消费者消费:消息变为可消费后,消费者正常消费。
// RocketMQ 事务消息示例
TransactionMQProducer producer = new TransactionMQProducer("tx_group");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务
try {
// 业务操作...
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 事务回查:检查本地事务是否成功
// 根据业务ID查询数据库,返回 COMMIT 或 ROLLBACK
return LocalTransactionState.COMMIT_MESSAGE;
}
});
producer.sendMessageInTransaction(msg, null);(3)各 MQ 对事务的支持
| 特性 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 事务机制 | 生产者事务(幂等 + 事务,保证精确一次语义) | 事务消息(半消息 + 本地事务 + 事务回查) | Channel 事务(txSelect/Commit/Rollback) |
| 事务范围 | 仅生产者到 Broker | 生产者本地事务 + Broker | 仅生产者到 Broker(同一 Channel) |
| 性能影响 | 中等 | 中等 | 极大(吞吐量可下降百倍量级) |
| 是否推荐 | 配合幂等使用 | 推荐,原生支持,金融级场景首选 | 不推荐,应使用 Publisher Confirms 替代 |
| 分布式事务 | 不支持 | 支持(最终一致性) | 不支持 |
(4)注意事项
- 幂等性是前提:无论哪种方案,消费者都必须实现幂等性,防止消息重复消费。
- 事务回查配置:RocketMQ Broker 默认每 60s(
transactionCheckInterval)扫描一次超时半消息;首次回查默认在事务超时(transactionTimeout,默认 6s)后进行,回查次数上限默认 15 次(transactionCheckMax,RocketMQ 4.x),超过后消息被丢弃,需监控悬挂消息并合理设置。 - 消息表性能:本地消息表方案中,消息表可能成为性能瓶颈,需考虑分表或使用 Redis 缓存。
- 监控与对账:无论哪种方案,都应有对账机制兜底,发现不一致时及时修复。
🔬 扩展知识
【L3】本地消息表 vs 事务消息 vs 最大努力通知
详情
| 维度 | 本地消息表 | 事务消息(半消息+回查) | 最大努力通知 |
|---|---|---|---|
| 一致性强度 | 最终一致,可靠 | 最终一致,可靠 | 最弱 |
| 消息投递延迟 | 扫描周期决定(秒级~分钟级) | 实时 | 实时 |
| 对业务的侵入 | 高(要建表、要扫描任务) | 中(需实现回查接口) | 低 |
| 额外存储成本 | 消息表(可能成为瓶颈) | 无 | 无 |
| 适用边界 | 任意 MQ,存量改造首选 | 需中间件原生支持 | 跨平台通知类弱一致场景 |
补充:强一致的 2PC/XA 方案在跨服务场景基本不可用(锁资源时间长、协调者单点、吞吐下降一个量级),MQ 方案的本质就是用最终一致性换可用性。
失效场景(事务保障何时失效):
- 本地消息表:业务写成功但消息表插入失败——若两者不在同一本地事务内,就会出现“业务完成但永远没人发消息”;必须同库同事务。
- 事务消息:回查接口写错——回查永远返回“提交”或永远返回“未知”,导致未成功的业务被下游消费,或消息永远悬挂;回查必须真正查库判断业务结果。
- 消费端不幂等——投递方重试 + 回查重发天然产生重复消息,下游不幂等则一致性被击穿。
- 消息表无限膨胀——“已发送”记录不清理,扫描任务越来越慢,新消息投递延迟从秒级恶化到分钟级。
- 最大努力通知的对账缺失——通知失败后无对账兜底,不一致永远不被发现。
【L4】MQ 事务与 TCC/Saga 的边界
详情
- MQ 事务只解决“本地业务成功后一定能通知到下游”,解决不了下游执行失败的撤销问题。
- 需要“预留资源 + 可撤销”语义(冻结库存/额度)时用 TCC;流程长、参与方多、每步可补偿时用 Saga。
【L3】Kafka 定制深挖:事务三件套
详情
- Kafka 0.11+ 的事务 = 事务协调器 + 幂等生产者 + read_committed 消费:协调器管理事务生命周期并把状态持久化到内部主题
__transaction_state;幂等生产者用 PID + 序列号防重试重复;isolation.level=read_committed的消费者只见已提交消息。 - 工作流程:
initTransactions获取事务 ID → 事务内写入带 PID+序列号的消息 → commit/abort 写协调器 → read_committed 消费者只见已提交。 - 它不是数据库式通用事务:只原子化“写 Kafka 分区 + 提交位点”这类内部动作,隔离级别仅 read_committed,下游外部系统仍需自己幂等。
- 两个坑:
transactional.id必须实例级唯一且稳定,多实例共用会互相 fence 抛ProducerFencedException;事务超过transaction.timeout.ms(默认 15 分钟)未提交会被 Broker 主动 abort。
【L3】RabbitMQ 定制深挖:Channel 事务与 Confirm 替代
详情
- 事务 API:
channel.txSelect()开启 → 发布消息暂存信道缓冲区 →txCommit()同步等 Broker 确认后才真正入队,失败则txRollback()丢弃。 - 性能代价极大:txCommit 是同步 RPC,管道完全串行化,吞吐相比异步 Confirm 下降数十到上百倍;且只覆盖生产者到 Broker 一段,不涉及消费端与下游。
- 原子性边界:只保证“本 Channel 内本次事务的多条发布一起入队或不入队”,不覆盖业务写库、消费端处理、跨 Channel/跨 Broker 投递——把它理解成“批量发布的原子提交”而非分布式事务。
- 事务 vs Confirm:事务强一致但同步阻塞,适合低频关键消息;Confirm 异步回调吞吐高且提供等效“到达可知”语义,配合未确认集维护 + 失败重发,是生产默认选择(需 v3.3+)。“写库 + 发消息”原子性需求用本地消息表或 RocketMQ 事务消息,RabbitMQ 事务解决不了。
📚 延伸阅读:RocketMQ 官方文档
🏭 实战场景
详情
量化数据
- 本地消息表:扫描周期常见 1~10s,消息投递延迟 ≈ 扫描周期 + 发送耗时;消息表建议按天分表,单表控制在千万级以内。
- 事务消息:首次回查默认在事务超时(约 6s)后触发,Broker 默认每 60s 扫描一轮,回查次数上限默认 15 次,超过后消息被丢弃或入异常队列,必须监控悬挂消息数。
- 性能代价:事务消息相比普通消息吞吐下降约 20%~40%;RabbitMQ 的 Channel 事务更极端,吞吐可下降百倍量级,生产中不应使用。
踩坑案例(推演示例):回查接口“恒返回提交”引发的资金不平
- 现象:对账发现一批“账户已扣款但权益未发放”的差异单,金额累计数十万元。
- 排查:这批差异均走事务消息链路;检查生产者回查逻辑,发现回查接口未接入业务查询,默认返回了“提交”;再查本地事务日志,发现这批订单本地事务实际因唯一键冲突回滚。
- 根因:本地事务回滚后应返回“回滚”,但回查实现是空壳;Broker 超时未收到确认时回查得到“提交”,把未成功业务的消息投递给了下游。
- 修复:回查接口改为真正查询业务表状态返回结果;对存量差异单批量冲正补发;增加“事务消息提交数 vs 业务成功数”的小时级对账监控。
场景题
场景:营销系统“下单送积分”链路:订单创建成功后发积分消息,最近发现约 0.1% 的订单没发积分,用户投诉。排查发现是生产端“先写库后发消息”,发消息时服务被发布重启。你如何设计修复方案?
应急处理:先跑一次性补偿脚本:按订单表反查积分表,对缺失积分的订单补发;同时暂停该服务的发布窗口,防止继续产生新差异。
根因分析:“本地事务成功 + 发消息是另一步”的方案存在天然窗口:两步之间宕机则消息永久丢失,且无任何重试机制。这正是本地消息表/事务消息要解决的“本地事务与消息发送的原子性”问题。
长期方案:改造为本地消息表:下单事务内同时写入订单与“待发送”积分消息记录,后台扫描任务秒级投递;或若中间件支持事务消息,改为半消息 + 本地事务 + 提交/回查,回查接口必须真实查库。消费端积分接口按订单号幂等,防重复加分。
权衡:本地消息表引入扫描延迟(秒级)与一张新表的维护成本,但积分发放本身非毫秒级敏感业务;若团队已使用支持事务消息的中间件,可省去消息表,但要为回查接口补充单测与监控,避免重蹈“空壳回查”覆辙。
⚠️ 常见误区
详情
常见误区:
- ❌ “用了事务消息就自动幂等了” → 事务消息只保证消息可靠投递,重试/回查仍会产生重复,消费端必须自行幂等。
- ❌ “回查接口返回 COMMIT 总归保险” → 空壳回查会把未成功业务的消息投递给下游,直接造成资损;回查必须真实查库。
- ❌ “本地消息表和事务消息能做到强一致” → 二者都是最终一致性,差别在投递延迟与侵入性,不提供跨服务强一致。
🔀 发散问题
- 本地消息表和事务消息都号称“最终一致”,两者在“失败窗口”上的本质区别是什么?
本地消息表把“消息投递”从业务事务中剥离交给后台扫描,失败窗口被“定时重试”兼容,代价是投递延迟;事务消息把投递时机交给 Broker 回查驱动,投递实时,但一致性上限依赖回查接口的正确性。前者更土但更稳,后者更优雅但多一个“回查写错”的故障面。 - 为什么消费端幂等是 MQ 事务方案的硬前提,而不是可选项?
所有方案都存在重发路径:扫描任务重复投递、回查后重发、消费确认失败重投。若下游不幂等,重发一次就多一条脏数据,“最终一致”直接变成“最终不一致”。因此幂等不是加固项,而是方案成立的前提条件。 - 什么场景下应该放弃 MQ 事务,改用 TCC 或 Saga?
当下游需要“预留资源 + 可撤销”的语义(如冻结库存、冻结额度)时用 TCC;当流程长、参与方多、每步可补偿(如跨多服务的履约链路)时用 Saga。MQ 事务只解决“本地业务成功后一定能通知到下游”的问题,解决不了下游执行失败的撤销问题。
【中等】什么是 Exactly-Once 语义?MQ 如何实现?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:MQ / 投递语义
💎 关键结论
三种投递语义中,绝大多数业务用 At Least Once + 业务幂等即可;真正的端到端 Exactly-Once 只能靠消费端幂等兑现,中间件能力只能收敛重复来源。理由:“处理成功但确认失败”的重投窗口在 Broker 视野之外。
⚡记忆卡片
口诀:至多一次会丢,至少一次会重,精确一次靠幂等
关键词:At Most Once/At Least Once/Exactly-Once/幂等生产者/事务消息
链路:先确认再处理→至多一次;先处理再确认→至少一次;幂等生产者+事务收敛重复→Kafka 精确一次;端到端仍需消费幂等
📖 核心知识
消息传递语义是 MQ 可靠性的核心概念,分为三种:
| 语义 | 含义 | 实现难度 | 典型场景 |
|---|---|---|---|
| At Most Once | 最多一次,消息可能丢失,但不会重复 | 最简单 | 日志收集、监控数据 |
| At Least Once | 至少一次,消息不会丢失,但可能重复 | 中等 | 大多数业务场景(配合幂等) |
| Exactly Once | 精确一次,消息不丢失也不重复 | 最难 | 金融交易、计费 |
各 MQ 的 Exactly-Once 实现
- Kafka:通过 幂等生产者(
enable.idempotence=true)+ 事务(transactional.id)实现 端到端 Exactly-Once。- 幂等生产者:为每条消息分配
PID和SequenceNumber,Broker 去重。 - 事务:保证“消费-处理-生产”链路的原子性。
- 幂等生产者:为每条消息分配
- RocketMQ:通过 事务消息 实现业务级的 Exactly-Once(依赖事务回查)。
- RabbitMQ:不原生支持 Exactly-Once,需业务侧通过幂等性设计保证。
实践建议
- 大多数场景:使用 At Least Once + 业务幂等 即可满足需求,实现简单且可靠。
- 严格 Exactly-Once:仅在金融等关键场景使用,需付出较高的性能和复杂度代价。
- 幂等性是关键:无论哪种语义,业务侧幂等设计都是兜底手段。
🔬 扩展知识
【L3】Kafka 精确一次的实现层次
详情
- 幂等生产者解决“单分区、单会话”的生产重复;事务(
transactional.id)把多分区写入与位点提交打包为原子操作,支撑 consume-transform-produce 场景。 - 严格地说,Kafka 的端到端 EOS 需配合支持事务的下游(如 Kafka Streams/Flink);写入普通数据库仍需业务幂等。
【L4】流计算框架与 EOS
详情
- Flink 通过 checkpoint + Kafka 事务实现两阶段提交,达成流处理链路端到端精确一次。
- 代价是 checkpoint 周期带来的延迟与吞吐下降,金融场景需实测评估。
【L3】Kafka 定制深挖:EOS 的构成与失效场景
详情
- Kafka 的 Exactly Once = 幂等生产 + 事务 + 精准 offset 控制:幂等生产者(PID + 序列号)去重单分区重试;事务(
commitTransaction/abortTransaction)跨分区原子写入;sendOffsetsToTransaction把消费位点提交纳入同一事务。 - consume-transform-produce 链路:读入 → 转换 → 写结果 Topic,用事务把三个动作原子化,配合
isolation.level=read_committed,任一环节失败整体 abort,重试时不重不漏。 - 失效场景:链路边界仅限 Kafka→Kafka,下游是 MySQL 等外部系统时需 Flink 两阶段提交把外部写入纳入 checkpoint,或直接外部幂等;
transaction.timeout.ms(默认 15 分钟)内未提交被 Broker 主动 abort,业务实际已成功时产生“消息丢失”假象;多实例共用同一transactional.id互相 fence 产生悬挂事务。 - CAP 权衡:Kafka 优先保 AP,通过事务补充一致性;Kafka Streams 用状态存储 + 检查点实现流处理 EOS。
📚 延伸阅读:Kafka 官方文档
🔀 发散问题
为什么 RabbitMQ 不支持 Exactly-Once?
其确认模型是“投递-ACK/重入队”二态状态机,无生产端去重与事务化位点机制,重复只能由业务幂等消化;这是“灵活路由优先”设计取舍的结果。
Exactly-Once 和幂等有什么关系?
Exactly-Once 是结果语义(业务效果只发生一次),幂等是实现手段;中间件层面的精确一次只能消除传输重复,端到端结果正确永远落在消费端幂等上,见本文档『如何保证 MQ 消息不重复?』。
MQ 架构
【困难】Kafka、RocketMQ、RabbitMQ 有什么区别?如何选型?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:12 min | 🏷 标签:MQ / 技术选型
💎 关键结论
三款 MQ 没有绝对优劣,只有定位差异:Kafka 拼吞吐(日志/大数据管道),RocketMQ 拼业务特性(事务/延迟/顺序消息),RabbitMQ 拼灵活路由与低延迟。选型先看场景再看指标,而不是背对比表。
⚡记忆卡片
- 口诀:Kafka 吞吐王,RocketMQ 业务全,RabbitMQ 路由灵
- 关键词:高吞吐 / 事务消息 / 延迟消息 / 灵活路由 / 消息堆积
- 链路:确认场景(日志/交易/业务路由)→ 对齐关键能力(堆积/事务/延迟)→ 结合团队技术栈与生态 → 定选型
📖 核心知识
这三款主流消息队列各有定位和适用场景,没有绝对优劣,只有是否适合业务场景。
核心对比表
| 对比维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 开发语言 | Scala/Java | Java | Erlang |
| 出身 | LinkedIn(Apache 顶级项目) | 阿里巴巴(Apache 顶级项目) | RabbitMQ Technologies(VMware) |
| 吞吐量 | 极高(百万级 TPS) | 高(十万级 TPS) | 中等(万级 TPS) |
| 延迟 | 毫秒级 | 毫秒级 | 微秒级(最低) |
| 消息可靠性 | 高(ISR + 副本) | 高(同步刷盘 + 主从复制) | 极高(多种确认机制 + 持久化) |
| 消息有序 | 分区内有序 | 队列内有序(支持严格顺序消息) | 队列内有序 |
| 事务消息 | 支持(0.11+,跨分区原子性) | 原生支持(半消息 + 事务回查) | 不支持(需 AMQP 事务,性能差) |
| 定时/延迟消息 | 不支持(需业务实现) | 原生支持(18 个延迟级别,5.0+ 支持任意时间) | 插件支持(rabbitmq_delayed_message_exchange) |
| 消息回溯 | 支持(按 offset) | 支持(按时间) | 不支持 |
| 消息堆积 | 极强(磁盘顺序写,亿级) | 极强(设计支持亿级堆积) | 弱(堆积后性能急剧下降) |
| 协议 | 自研协议 | 自研协议 | AMQP/STOMP/MQTT(标准协议) |
| 路由 | Topic/Partition | Topic/Tag/MessageQueue | 灵活(Exchange + Routing Key) |
| 管理界面 | 第三方(Kafka Manager 等) | 自带 Dashboard | 自带 Web 管理界面 |
| 元数据管理 | ZooKeeper/KRaft | NameServer | Mnesia 数据库 |
| 生态 | 大数据生态(Kafka Connect/Streams/ksqlDB) | 阿里云生态 | Spring AMQP 生态 |
选型建议
- 选 Kafka:
- 大数据场景(日志收集、用户行为追踪、流计算)。
- 超高吞吐量需求(日志、监控、埋点)。
- 与 Hadoop/Spark/Flink 等大数据组件集成。
- 消息回溯需求。
- 选 RocketMQ:
- 金融级业务(订单、交易、支付)。
- 需要事务消息、定时消息、顺序消息等高级特性。
- 国内技术栈(社区活跃、中文文档完善)。
- 消息堆积能力强但需要业务级可靠性保证。
- 选 RabbitMQ:
- 复杂路由需求(路由键、主题匹配)。
- 低延迟、高可靠场景(<1 万 TPS 但要求极高可靠性)。
- 需要标准协议(AMQP)。
- 小规模系统、微服务异步通信。
不适用场景速查
- Kafka 不适用:需要复杂路由、低延迟(毫秒级)响应的业务。
- RocketMQ 不适用:需要与 Kafka 生态紧密集成的项目。
- RabbitMQ 不适用:需要海量消息堆积和高吞吐量的场景。
架构差异本质
- Kafka:以吞吐量为核心,存储采用 Partition 独立日志,适合流式数据。
- RocketMQ:以业务可靠性为核心,存储采用 CommitLog 统一日志 + ConsumeQueue 索引,融合了 Kafka 的高吞吐和 RabbitMQ 的业务特性。
- RabbitMQ:以灵活路由和协议兼容为核心,基于 Erlang 的 Actor 模型,适合业务消息。
🔬 扩展知识
【L3】存储模型差异如何决定行为差异
详情
- Kafka:每分区一个独立日志,追加写 + 页缓存 + 零拷贝,堆积能力极强,但分区元数据随规模膨胀。
- RocketMQ:所有 Topic 写同一个 CommitLog,顺序写一次落盘,ConsumeQueue 作为消费索引,读写分离,业务消息友好。
- RabbitMQ:消息存内存 + 磁盘队列,确认后删除,无堆积设计,堆积时内存压力导致性能急剧下降。
【L4】选型中的隐性成本
详情
- 运维复杂度:Kafka 需关注分区/副本/再均衡,RocketMQ 关注 CommitLog 磁盘容量,RabbitMQ 关注内存水位与镜像队列。
- 团队技能:Erlang 栈排障门槛高;Java 团队对 RocketMQ/Kafka 源码级排查更友好。
- 迁移成本:协议不兼容(Kafka 自研协议 vs AMQP),迁移基本等于重写接入层。
📚 延伸阅读:Kafka 官方文档
🏭 实战场景
详情
场景:日志平台 + 订单交易共存,一套 MQ 还是两套?(场景为推演示例)
- 背景:日志埋点峰值约 50 万条/秒(单条 200B),订单交易峰值约 2000 TPS 但要求事务与延迟消息。
- 决策:日志链路用 Kafka(吞吐与堆积是硬需求,允许秒级延迟);交易链路用 RocketMQ(原生事务/延迟/顺序消息,吞吐在容量内)。两套 MQ 各取所长,避免单套迁就导致两头受限。
- 代价与兜底:多一套中间件的运维成本;通过统一监控告警与双链路演练对冲,日志与交易互不干扰(日志洪峰不影响交易 SLA)。
⚠️ 常见误区
详情
常见误区:
- ❌ “Kafka 吞吐高,所有场景都用 Kafka 就行” → Kafka 不支持原生延迟消息,事务链路也只限 Kafka→Kafka,业务消息场景 RocketMQ 的特性更省心。
- ❌ “RabbitMQ 吞吐低是因为 Erlang 不行” → 主要是“确认后删除 + 内存优先”的队列模型不为堆积设计;低延迟和灵活路由正是它的主场。
- ❌ “事务消息选 Kafka 还是 RocketMQ 都一样” → RocketMQ 的半消息 + 事务回查面向“本地事务 + 发消息”的业务一致性;Kafka 事务面向 Kafka→Kafka 链路的原子性,解决的不是同一个问题。
🔀 发散问题
Q1:如果团队只会 Java,但需要 AMQP 协议怎么办?
A:协议兼容和团队技能冲突时,优先看是否真的需要 AMQP(通常只有跨语言/跨厂商互通才需要);若只是想要路由能力,RocketMQ 的 Tag/SQL 过滤或 Kafka 多 Topic 设计也能覆盖。
Q2:亿级堆积场景为什么 RabbitMQ 撑不住而 Kafka 可以?
A:Kafka 是磁盘顺序日志,堆积只是多占磁盘,消费速率不受存量影响;RabbitMQ 消息在队列中等待确认,堆积时内存水位告警会阻塞发布,磁盘队列性能也急剧下降。
Q3:Kafka 的消息回溯能力为什么对选型影响很大?
A:可按 offset/时间重读历史数据,意味着下游故障修复、数据重算、新消费组接入都不需要生产端重发;这是日志与数据管道场景选 Kafka 的重要理由。
【困难】Kafka、ActiveMQ、RabbitMQ、RocketMQ 有什么优缺点?⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:15 min | 🏷 标签:MQ / 产品对比
💎 关键结论
四款 MQ 按吞吐分两档:Kafka/RocketMQ 十万级,ActiveMQ/RabbitMQ 万级;按定位分三类:Kafka 日志流、RocketMQ 业务交易、RabbitMQ 灵活路由,ActiveMQ 已逐渐退出主流。理由:单机吞吐、topic 数量敏感度、延迟与生态活跃度是拉开差距的关键维度。
⚡记忆卡片
口诀:K 拼吞吐,R 拼交易,Rabbit 拼路由,Active 渐退场
关键词:单机吞吐量/时效性/可用性/消息可靠性/生态活跃度
链路:明确吞吐量级与功能诉求 → 对比四款产品的吞吐/延迟/可靠性/生态 → 按公司规模与场景定候选 → 小步压测验证
📖 核心知识
ActiveMQ
ActiveMQ 是由 Apache 出品,ActiveMQ 是一个完全支持JMS1.1 和 J2EE 1.4 规范的 JMS Provider 实现。它非常快速,支持 多种语言的客户端 和 协议,而且可以非常容易的嵌入到企业的应用环境中,并有许多高级功能。
(a) 主要特性
- 服从 JMS 规范:
JMS规范提供了良好的标准和保证,包括:同步 或 异步 的消息分发,一次和仅一次的消息分发,消息接收 和 订阅 等等。遵从JMS规范的好处在于,不论使用什么JMS实现提供者,这些基础特性都是可用的; - 连接灵活性:
ActiveMQ提供了广泛的 连接协议,支持的协议有:HTTP/S,IP多播,SSL,TCP,UDP等等。对众多协议的支持让ActiveMQ拥有了很好的灵活性; - 支持的协议种类多:
OpenWire、STOMP、REST、XMPP、AMQP; - 持久化插件和安全插件:
ActiveMQ提供了 多种持久化 选择。而且,ActiveMQ的安全性也可以完全依据用户需求进行 自定义鉴权 和 授权; - 支持的客户端语言种类多:除了
Java之外,还有:C/C++,.NET,Perl,PHP,Python,Ruby; - 代理集群:多个
ActiveMQ代理 可以组成一个 集群 来提供服务; - 异常简单的管理:
ActiveMQ是以开发者思维被设计的。所以,它并不需要专门的管理员,因为它提供了简单又使用的管理特性。有很多中方法可以 监控ActiveMQ不同层面的数据,包括使用在JConsole或者在ActiveMQ的Web Console中使用JMX。通过处理JMX的告警消息,通过使用 命令行脚本,甚至可以通过监控各种类型的 日志。
(b) 部署环境
ActiveMQ 可以运行在 Java 语言所支持的平台之上。使用 ActiveMQ 需要:
Java JDKActiveMQ安装包
(c) 优点
- 跨平台 (
JAVA编写与平台无关,ActiveMQ几乎可以运行在任何的JVM上); - 可以用
JDBC:可以将 数据持久化 到数据库。虽然使用JDBC会降低ActiveMQ的性能,但是数据库一直都是开发人员最熟悉的存储介质; - 支持
JMS规范:支持JMS规范提供的 统一接口; - 支持 自动重连 和 错误重试机制;
- 有安全机制:支持基于
shiro,jaas等多种 安全配置机制,可以对Queue/Topic进行 认证和授权; - 监控完善:拥有完善的 监控,包括
Web Console,JMX,Shell命令行,Jolokia的RESTful API; - 界面友善:提供的
Web Console可以满足大部分情况,还有很多 第三方的组件 可以使用,比如hawtio;
(d) 缺点
- 社区活跃度不及
RabbitMQ高; - 根据其他用户反馈,会出莫名其妙的问题,会 丢失消息;
- 目前重心放到
activemq 6.0产品Apollo,对5.x的维护较少; - 不适合用于 上千个队列 的应用场景;
RabbitMQ
RabbitMQ 于 2007 年发布,是一个在 AMQP (高级消息队列协议) 基础上完成的,可复用的企业消息系统,是当前最主流的消息中间件之一。
(a) 主要特性
- 可靠性:提供了多种技术可以让你在 性能 和 可靠性 之间进行 权衡。这些技术包括 持久性机制、投递确认、发布者证实 和 高可用性机制;
- 灵活的路由:消息在到达队列前是通过 交换机 进行 路由 的。
RabbitMQ为典型的路由逻辑提供了 多种内置交换机 类型。如果你有更复杂的路由需求,可以将这些交换机组合起来使用,你甚至可以实现自己的交换机类型,并且当做RabbitMQ的 插件 来使用; - 消息集群:在相同局域网中的多个
RabbitMQ服务器可以 聚合 在一起,作为一个独立的逻辑代理来使用; - 队列高可用:队列可以在集群中的机器上 进行镜像,以确保在硬件问题下还保证 消息安全;
- 支持多种协议:支持 多种消息队列协议;
- 支持多种语言:用
Erlang语言编写,支持只要是你能想到的 所有编程语言; - 管理界面:
RabbitMQ有一个易用的 用户界面,使得用户可以 监控 和 管理 消息Broker的许多方面; - 跟踪机制:如果 消息异常,
RabbitMQ提供消息跟踪机制,使用者可以找出发生了什么; - 插件机制:提供了许多 插件,来从多方面进行扩展,也可以编写自己的插件。
(b) 部署环境
RabbitMQ 可以运行在 Erlang 语言所支持的平台之上,包括 Solaris,BSD,Linux,MacOSX,TRU64,Windows 等。使用 RabbitMQ 需要:
ErLang语言包RabbitMQ安装包
(c) 优点
- 由于
Erlang语言的特性,消息队列性能较好,支持 高并发; - 健壮、稳定、易用、跨平台、支持 多种语言、文档齐全;
- 有消息 确认机制 和 持久化机制,可靠性高;
- 高度可定制的 路由;
- 管理界面 较丰富,在互联网公司也有较大规模的应用,社区活跃度高。
(d) 缺点
- 尽管结合
Erlang语言本身的并发优势,性能较好,但是不利于做 二次开发和维护; - 实现了 代理架构,意味着消息在发送到客户端之前可以在 中央节点 上排队。此特性使得
RabbitMQ易于使用和部署,但是使得其 运行速度较慢,因为中央节点 增加了延迟,消息封装后 也比较大; - 需要学习 比较复杂 的 接口和协议,学习和维护成本较高。
RocketMQ
RocketMQ 出自 阿里 的开源产品,用 Java 语言实现,在设计时参考了 Kafka,并做出了自己的一些改进,消息可靠性上 比 Kafka 更好。RocketMQ 在阿里内部 被广泛应用在 订单,交易,充值,流计算,消息推送,日志流式处理,binglog 分发 等场景。
(a) 主要特性
- 基于 队列模型:具有 高性能、高可靠、高实时、分布式 等特点;
Producer、Consumer、队列 都支持 分布式;Producer向一些队列轮流发送消息,队列集合 称为Topic。Consumer如果做 广播消费,则一个Consumer实例消费这个Topic对应的 所有队列;如果做 集群消费,则 多个Consumer实例 平均消费 这个Topic对应的队列集合;- 能够保证 严格的消息顺序;
- 提供丰富的 消息拉取模式;
- 高效的订阅者 水平扩展能力;
- 实时 的 消息订阅机制;
- 亿级 消息堆积 能力;
- 较少的外部依赖。
(b) 部署环境
RocketMQ 可以运行在 Java 语言所支持的平台之上。使用 RocketMQ 需要:
Java JDK- 安装
git、Maven RocketMQ安装包
(c) 优点
- 单机 支持
1万以上 持久化队列; RocketMQ的所有消息都是 持久化的,先写入系统PAGECACHE,然后 刷盘,可以保证 内存 与 磁盘 都有一份数据,而 访问 时,直接 从内存读取。- 模型简单,接口易用(
JMS的接口很多场合并不太实用); - 性能非常好,可以允许 大量堆积消息 在
Broker中; - 支持 多种消费模式,包括 集群消费、广播消费等;
- 各个环节 分布式扩展设计,支持 主从 和 高可用;
- 开发度较活跃,版本更新很快。
(d) 缺点
- 支持的 客户端语言 不多,目前是
Java及C++,其中C++还不成熟; RocketMQ社区关注度及成熟度也不及前两者;- 没有
Web管理界面,提供了一个CLI(命令行界面) 管理工具带来 查询、管理 和 诊断各种问题; - 没有在
MQ核心里实现JMS等接口;
Kafka
Apache Kafka 是一个 分布式消息发布订阅 系统。它最初由 LinkedIn 公司基于独特的设计实现为一个 分布式的日志提交系统 (a distributed commit log),之后成为 Apache 项目的一部分。Kafka 性能高效、可扩展良好 并且 可持久化。它的 分区特性,可复制 和 可容错 都是其不错的特性。
(a) 主要特性
- 快速持久化:可以在
O(1)的系统开销下进行 消息持久化; - 高吞吐:在一台普通的服务器上既可以达到
10W/s的 吞吐速率; - 完全的分布式系统:
Broker、Producer和Consumer都原生自动支持 分布式,自动实现 负载均衡; - 支持 同步 和 异步 复制两种 高可用机制;
- 支持 数据批量发送 和 拉取;
- 零拷贝技术 (zero-copy):减少
IO操作步骤,提高 系统吞吐量; - 数据迁移、扩容 对用户透明;
- 无需停机 即可扩展机器;
- 其他特性:丰富的 消息拉取模型、高效 订阅者水平扩展、实时的 消息订阅、亿级的 消息堆积能力、定期删除机制;
(b) 部署环境
使用 Kafka 需要:
Java JDKKafka安装包
(c) 优点
- 客户端语言丰富:支持
Java、.Net、PHP、Ruby、Python、Go等多种语言; - 高性能:单机写入
TPS约在100万条/秒,消息大小10个字节; - 提供 完全分布式架构,并有
replica机制,拥有较高的 可用性 和 可靠性,理论上支持 消息无限堆积; - 支持批量操作;
- 消费者 采用
Pull方式获取消息。消息有序,通过控制 能够保证所有消息被消费且仅被消费 一次; - 有优秀的第三方
Kafka Web管理界面Kafka-Manager; - 在 日志领域 比较成熟,被多家公司和多个开源项目使用。
(d) 缺点
Kafka单机超过64个 队列/分区 时,Load时会发生明显的飙高现象。队列 越多,负载 越高,发送消息 响应时间变长;- 使用 短轮询方式,实时性 取决于 轮询间隔时间;
- 消费失败 不支持重试;
- 支持 消息顺序,但是 一台代理宕机 后,就会产生 消息乱序;
- 社区更新较慢。
技术选型对比
| 特性 | ActiveMQ | RabbitMQ | RocketMQ | Kafka |
|---|---|---|---|---|
| 单机吞吐量 | 万级,比 RocketMQ、Kafka 低一个数量级 | 同 ActiveMQ | 10 万级,支撑高吞吐 | 10 万级,高吞吐,一般配合大数据类的系统来进行实时数据计算、日志采集等场景 |
| topic 数量对吞吐量的影响 | topic 可以达到几百/几千的级别,吞吐量会有较小幅度的下降,这是 RocketMQ 的一大优势,在同等机器下,可以支撑大量的 topic | topic 从几十到几百个时候,吞吐量会大幅度下降,在同等机器下,Kafka 尽量保证 topic 数量不要过多,如果要支撑大规模的 topic,需要增加更多的机器资源 | ||
| 时效性 | ms 级 | 微秒级,这是 RabbitMQ 的一大特点,延迟最低 | ms 级 | 延迟在 ms 级以内 |
| 可用性 | 高,基于主从架构实现高可用 | 同 ActiveMQ | 非常高,分布式架构 | 非常高,分布式,一个数据多个副本,少数机器宕机,不会丢失数据,不会导致不可用 |
| 消息可靠性 | 有较低的概率丢失数据 | 基本不丢 | 经过参数优化配置,可以做到 0 丢失 | 同 RocketMQ |
| 功能支持 | MQ 领域的功能极其完备 | 基于 erlang 开发,并发能力很强,性能极好,延时很低 | MQ 功能较为完善,还是分布式的,扩展性好 | 功能较为简单,主要支持简单的 MQ 功能,在大数据领域的实时计算以及日志采集被大规模使用 |
综上,各种对比之后,有如下建议:
- 一般的业务系统要引入 MQ,最早大家都用 ActiveMQ,但现在确实大家用的不多了,没经过大规模吞吐量场景的验证,社区也不是很活跃,所以大家还是算了吧,我个人不推荐用这个了;
- 后来大家开始用 RabbitMQ,但是确实 erlang 语言阻止了大量的 Java 工程师去深入研究和掌控它,对公司而言,几乎处于不可控的状态,但是确实人家是开源的,比较稳定的支持,活跃度也高;
- 不过现在确实越来越多的公司会去用 RocketMQ,确实很不错,毕竟是阿里出品,但社区可能有突然黄掉的风险(目前 RocketMQ 已捐给 Apache,但 GitHub 上的活跃度其实不算高)对自己公司技术实力有绝对自信的,推荐用 RocketMQ,否则回去老老实实用 RabbitMQ 吧,人家有活跃的开源社区,绝对不会黄。
- 所以中小型公司,技术实力较为一般,技术挑战不是特别高,用 RabbitMQ 是不错的选择;大型公司,基础架构研发实力较强,用 RocketMQ 是很好的选择。
- 如果是大数据领域的实时计算、日志采集等场景,用 Kafka 是业内标准的,绝对没问题,社区活跃度很高,绝对不会黄,何况几乎是全世界这个领域的事实性规范。
🔬 扩展知识
【L3】吞吐差异的根因
详情
- Kafka/RocketMQ 靠顺序追加写 + 页缓存 + 批量 + 零拷贝/mmap,把磁盘当内存用,吞吐达十万级;ActiveMQ/RabbitMQ 的内存队列模型与确认即删设计,堆积与吞吐都受内存约束,停留在万级。
- topic/队列数量敏感度:Kafka 按分区建文件,分区过多产生随机写;RocketMQ 混写 CommitLog,对队列数量不敏感。
【L4】生态活跃度的选型权重
详情
- 中间件的长期风险不在功能而在社区:文档质量、issue 响应、版本节奏决定了遇到问题能否自救。
- 对 Erlang 技术栈缺失的团队,RabbitMQ 的二次开发能力几乎为零,选型时应把“可控性”计入成本。
🏭 实战场景
详情
(推演案例)某中型电商选型评估:日常订单峰值 3000 TPS,需要订单事件被库存、营销、积分三个系统同时消费,且要求事务性通知。候选对比:ActiveMQ 万级吞吐够用但社区活跃度低且存在丢消息口碑风险,排除;RabbitMQ 延迟最低但堆积能力弱、无原生事务消息;RocketMQ 十万级吞吐 + 事务消息 + 顺序消息直接命中诉求,最终选 RocketMQ;同期日志链路(日增约 50GB、峰值 5 万条/秒)单独用 Kafka 承接,与业务链路隔离。关键结论:不要用一个 MQ 硬吃所有场景,按链路拆分选型更稳。
⚠️ 常见误区
详情
常见误区:
- ❌ “吞吐量看官方宣传值就行” → 官方值多为特定消息大小与配置下的极限,选型必须用自己的真实消息体与可靠配置压测。
- ❌ “RabbitMQ 延迟低所以适合一切业务” → 它的堆积能力弱,海量积压场景会迅速恶化。
- ❌ “ActiveMQ 成熟稳定可以随便用” → 其社区活跃度与大规模验证不足,新项目不建议首选。
🔀 发散问题
为什么不推荐用 RabbitMQ 做海量日志采集?
日志场景吞吐高、可容忍少量丢失、需要长保留与重放,正是 Kafka 的优势区间;RabbitMQ 万级吞吐与确认即删模型会导致堆积即性能崩塌。
同一公司能否同时用两套 MQ?
可以且常见:交易链路用 RocketMQ(事务/顺序),数据链路用 Kafka(吞吐/生态);代价是运维两套集群,需评估团队规模。
【简单】什么是 JMS?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:MQ / 协议规范
💎 关键结论
JMS 是 Java 平台的消息服务标准/规范(API),不是具体产品;它定义了 P2P 与 Pub/Sub 两种模型和同步/异步两种消费方式,ActiveMQ 是典型实现。理由:标准化的 API 让应用与具体消息实现解耦,切换 Provider 不改业务代码。
⚡记忆卡片
口诀:JMS 是规范,P2P 独占,Pub/Sub 广播
关键词:JMS Provider/Queue/Topic/Sender/Receiver/持久订阅
链路:应用通过 JMS API 发消息 → Provider(如 ActiveMQ)投递 → P2P 由唯一接收者取走 / Pub/Sub 广播给所有订阅者
📖 核心知识
提到 MQ,就顺便提一下 JMS 。
JMS(JAVA Message Service,java 消息服务)API 是一个消息服务的标准/规范,允许应用程序组件基于 JavaEE 平台创建、发送、接收和读取消息。它使分布式通信耦合度更低,消息服务更加可靠以及异步性。
在 EJB 架构中,有消息 bean 可以无缝的与 JMS 消息服务集成。在 J2EE 架构模式中,有消息服务者模式,用于实现消息与应用直接的解耦。
JMS 消息模型:在 JMS 标准中,有两种消息模型:P2P(Point to Point)、Pub/Sub(Publish/Subscribe)。
P2P 模式

P2P 模式包含三个角色:MQ(Queue),发送者 (Sender),接收者 (Receiver)。每个消息都被发送到一个特定的队列,接收者从队列中获取消息。队列保留着消息,直到他们被消费或超时。
P2P 的特点:
- 每个消息只有一个消费者(Consumer)(即一旦被消费,消息就不再在 MQ 中)
- 发送者和接收者之间在时间上没有依赖性,也就是说当发送者发送了消息之后,不管接收者有没有正在运行,它不会影响到消息被发送到队列
- 接收者在成功接收消息之后需向队列应答成功
如果希望发送的每个消息都会被成功处理的话,那么需要 P2P 模式。
Pub/Sub 模式

包含三个角色:主题(Topic),发布者(Publisher),订阅者(Subscriber)。多个发布者将消息发送到 Topic,系统将这些消息传递给多个订阅者。
Pub/Sub 的特点:
- 每个消息可以有多个消费者
- 发布者和订阅者之间有时间上的依赖性。针对某个主题(Topic)的订阅者,它必须创建一个订阅者之后,才能消费发布者的消息。
- 为了消费消息,订阅者必须保持运行的状态。
为了缓和这样严格的时间相关性,JMS 允许订阅者创建一个可持久化的订阅。这样,即使订阅者没有被激活(运行),它也能接收到发布者的消息。
如果希望发送的消息可以不被做任何处理、或者只被一个消息者处理、或者可以被多个消费者处理的话,那么可以采用 Pub/Sub 模型。
JMS 消息消费:在 JMS 中,消息的产生和消费都是异步的。对于消费来说,JMS 的消息者可以通过两种方式来消费消息:
- 同步 - 订阅者或接收者通过
receive方法来接收消息,receive方法在接收到消息之前(或超时之前)将一直阻塞; - 异步 - 订阅者或接收者可以注册为一个消息监听器。当消息到达之后,系统自动调用监听器的
onMessage方法。
JNDI - Java 命名和目录接口,是一种标准的 Java 命名系统接口。可以在网络上查找和访问服务。通过指定一个资源名称,该名称对应于数据库或命名服务中的一个记录,同时返回资源连接建立所必须的信息。JNDI 在 JMS 中起到查找和访问发送目标或消息来源的作用。
【简单】什么是 AMQP?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:MQ / 协议规范
💎 关键结论
AMQP 是应用层的开放消息协议,定义了消息格式、传输与处理机制,RabbitMQ 是其最知名实现。理由:它把 Exchange/Binding 路由、确认、持久化等语义写进协议本身,天然具备高可靠与跨平台能力。
⚡记忆卡片
口诀:AMQP 定协议,Exchange 路由,确认+持久化保可靠
关键词:Connection/Channel/Exchange/Queue/Binding/Routing Key
链路:客户端建 Connection 开 Channel → 消息发到 Exchange → 按 Binding 规则路由到 Queue → 消费者取走并确认
📖 核心知识
AMQP(Advanced Message Queuing Protocol)是一种应用层协议,用于在消息队列系统中定义消息的格式、传输方式和处理机制。
AMQP 是一个面向消息的、异步传输的协议,具有高可靠性、可拓展性、跨平台的特性,适合在分布式系统中传输重要数据。它是 RabbitMQ、ActiveMQ 等消息中间件的底层协议。
核心组件
| 组件 | 作用 |
|---|---|
| Connection | 客户端与消息代理(如 RabbitMQ)的物理连接 |
| Channel | 逻辑通信链路(多路复用连接,减少开销) |
| Exchange | 消息路由枢纽(支持多种路由策略) |
| Queue | 存储消息的容器,消费者从中拉取数据 |
| Binding | 定义 Exchange 与 Queue 的映射关系(含路由规则) |
路由模型
- Direct:精确匹配路由键(如
order.create) - Fanout:广播到所有绑定队列(无视路由键)
- Topic:通配符匹配(如
order.*) - Headers:基于消息头键值匹配(非路由键)
可靠性保障
- 消息确认:消费者手动确认,失败则重入队列
- 持久化:队列/消息持久化防丢失
- 事务:支持原子性提交/回滚(批量操作)
协议对比
| 协议 | 优势场景 | 局限性 |
|---|---|---|
| AMQP | 企业级应用(强事务、高可靠) | 略重(不适合 IoT 轻量场景) |
| MQTT | 物联网(低功耗、低带宽) | 功能简单(无复杂路由) |
| JMS | Java 生态集成 | 仅限 Java,跨平台性弱 |
【困难】RocketMQ 与 Kafka 有什么区别?各自的优势是什么?⭐⭐⭐
🎯 目标等级:L3 | ⏱ 建议用时:15 min | 🏷 标签:MQ / 架构对比
💎 关键结论
RocketMQ 最初借鉴 Kafka 设计,但走向不同:Kafka 用 Partition 独立日志 + 批处理追求极致吞吐,适合大数据管道;RocketMQ 用全局 CommitLog + 二级索引追求金融级业务能力(事务/延迟/顺序/轨迹内置),适合业务消息。选型看场景:日志埋点选 Kafka,订单交易选 RocketMQ。
⚡记忆卡片
- 口诀:Kafka 吞吐强适合管道,RocketMQ 功能全适合交易
- 关键词:CommitLog vs Partition / NameServer vs ZooKeeper / 事务消息 / 延迟消息 / 技术选型
- 链路:场景决定选型 → 大数据/日志 → Kafka(Partition 独立日志、批量压缩)→ 金融/业务消息 → RocketMQ(全局顺序写、内置事务延迟轨迹)
📖 核心知识
核心架构对比
| 对比维度 | Kafka | RocketMQ |
|---|---|---|
| 元数据管理 | ZooKeeper / KRaft | NameServer(无状态,AP) |
| 存储架构 | 每个 Partition 独立 Log 文件 | 所有 Topic 共用 CommitLog + ConsumeQueue 索引 |
| 存储方式 | Partition 级别顺序写 | 全局顺序写 CommitLog,异步构建索引 |
| 副本同步 | ISR 机制 + Leader/Follower | 主从复制 / DLedger / Controller |
| 分区/队列 | Partition | MessageQueue |
| 消费模式 | 拉模式(Pull) | 推模式(Push,封装的 Pull)+ 拉模式 |
| 消息索引 | 偏移量索引 + 时间戳索引 | ConsumeQueue + IndexFile(支持按 Key 查询) |
功能特性对比
| 特性 | Kafka | RocketMQ |
|---|---|---|
| 事务消息 | 支持(0.11+,跨分区原子写入) | 原生支持(半消息+回查,业务侵入小) |
| 延迟消息 | 不支持 | 原生支持(4.x 18 级,5.0+ 支持任意时间) |
| 顺序消息 | 分区内有序 | 队列内严格有序(专门优化) |
| 消息回溯 | 按 offset | 按时间 |
| 消息过滤 | 客户端过滤 | Broker 端过滤(Tag/SQL92) |
| 死信队列 | 无内置 | 内置(自动转入 %DLQ% Topic) |
| 消息轨迹 | 无内置 | 内置(RMQ_SYS_TRACE_TOPIC) |
| 消息堆积 | 强(磁盘顺序写) | 更强(设计支持亿级堆积) |
存储设计的本质差异
- Kafka 存储(Partition 独立日志):每个 Partition 一个独立目录,消息只属于该 Partition。优点:删除/清理时按 Partition 独立操作,简单高效;缺点:Topic/Partition 数量多时,磁盘随机写严重(多 Partition 并发写不同文件)。
- RocketMQ 存储(全局 CommitLog):所有 Topic 的消息混合写入同一个 CommitLog,ConsumeQueue 作为索引。优点:完全顺序写,即使 Topic 数量多也能保持高吞吐,堆积能力强;缺点:消费时需根据 ConsumeQueue 索引跳转读取,读取有额外开销;删除旧消息按文件粒度较粗。
性能对比
| 指标 | Kafka | RocketMQ |
|---|---|---|
| 单机吞吐 | 更高(百万 TPS 量级) | 高(十万 TPS 量级) |
| Topic 数量影响 | 大(Partition 多则随机写严重) | 小(CommitLog 统一写) |
| 延迟 | 毫秒级 | 毫秒级 |
| 消息大小 | 适合大批量 | 适合中小消息 |
选型建议
- 选 Kafka:大数据场景(日志、监控、埋点)、超高吞吐、与 Hadoop/Spark/Flink 集成。
- 选 RocketMQ:金融业务(订单、交易)、需要事务/延迟/顺序消息、Topic 数量多、国内技术栈。
- 协议兼容:RocketMQ 兼容 JMS、MQTT 等协议;Kafka 仅原生协议,多协议互通需求是选 RocketMQ 的加分项。
🔬 扩展知识
【L3】副本机制差异决定切主行为
详情
- Kafka 用 ISR 集合 + Controller 自动切主(KRaft 后更轻量);RocketMQ 传统主从需人工切换,4.5+ DLedger、5.0+ Controller 才补齐自动切主能力,选型时要把高可用方案一并评估。
- 混合架构在大型互联网公司很常见:数据管道用 Kafka,业务交易链路用 RocketMQ,两者并存而非二选一。
【L4】性能对比的正确姿势
详情
- 单机吞吐数字受硬件、批量大小、ack 策略影响巨大,官方基准与实际业务差异可达数倍;选型前应在同配置环境用真实消息模型分别压测,而不是直接引用宣传数字。
【L4】两者边界的融合趋势
详情
- Kafka 通过 KRaft 去 ZooKeeper 化、引入更灵活的元数据管理;RocketMQ 5.0 引入存算分离、gRPC 协议与任意时长定时消息,两者在架构上互相借鉴。
- 选型时建议以当前稳定版本的功能矩阵为准,避免拿未落地特性做决策。
📚 延伸阅读:RocketMQ 官方文档、Kafka 官方文档
🏭 实战场景
详情
选型推演:电商订单链路 vs 埋点链路(以下为架构推演,非生产实测数据)
- 假设订单链路峰值 5 万 TPS、需要事务消息保下单与消息原子性、需要按订单号查轨迹:RocketMQ 内置事务消息 + 消息轨迹 + 按 Key 查询可直接满足;Kafka 需自研本地消息表 + 外部链路系统补齐。
- 假设埋点链路峰值 50 万 TPS、下游是 Flink 实时计算、容忍少量丢失:Kafka 的批量压缩与 Flink 生态集成优势明显;RocketMQ 吞吐余量紧张且与大数据生态对接成本高。
- 结论:同一系统内两条链路分别选型、并存使用,是大型架构的常见形态。
⚠️ 常见误区
详情
常见误区:
- ❌ "RocketMQ 就是 Kafka 的仿制品" → RocketMQ 早期借鉴了日志存储思想,但为金融场景自研了二级索引、事务消息、回查机制等大量能力,架构已显著分化。
- ❌ "吞吐高的就一定更好" → 吞吐只是维度之一;业务消息场景更看重事务、延迟、轨迹、运维能力,盲目选 Kafka 会背大量自研成本。
- ❌ "Kafka 完全不支持事务" → Kafka 0.11+ 支持事务(跨分区原子写入 + 精确一次语义),但面向流处理管道,与 RocketMQ 面向业务一致性的半消息回查定位不同。
🔀 发散问题
Q1:为什么 RocketMQ 在海量 Topic 下写入不劣化而 Kafka 会?
A:RocketMQ 所有 Topic 混写单个 CommitLog,写入始终是单一文件的顺序追加;Kafka 每 Partition 独立日志,分区多时写路径变多文件并发写。
Q2:两者的消费模型有什么差异?
A:Kafka 纯拉模式,按 Partition 分配;RocketMQ 提供 Push(长轮询封装)与 Pull,5.0 还新增 Pop 无状态消费,见《RocketMQ面试》『RocketMQ 5.0 有哪些新特性?』。
Q3:如果团队已有 Kafka,还要引 RocketMQ 吗?
A:看业务需求:若需要事务/延迟/顺序消息且自研成本不可接受,引入 RocketMQ 专门承载业务链路是合理选择;否则优先用本地消息表等通用方案避免多套 MQ 的运维成本。
Q4:为什么 RocketMQ 用 NameServer 而不用 ZooKeeper?
A:NameServer 无状态、各自独立、AP 模型,部署运维极简,满足路由发现这类“最终一致即可”的需求;ZooKeeper 的强一致选举对消息路由场景是过重依赖,详见《RocketMQ面试》『为什么 RocketMQ 不用 ZooKeeper,而是自己开发 NameServer?』。
Q5:两者都支持事务,能互相替代吗?
A:不能:Kafka 事务面向“内部消息原子写入+精确一次”,RocketMQ 事务消息面向“本地业务事务与消息投递的原子性(半消息+回查)”,解决的是不同层面的问题,见本文档『MQ 如何实现分布式事务?』。