《消息队列高手课》笔记
《消息队列高手课》笔记
为什么需要消息队列?
核心应用:异步处理、系统解耦、流量削峰、系统间通信、数据缓冲、数据一致性。
该如何选择消息队列?
选型维度:是否开源(能否商用)> 社区活跃度 > 技术生态适配性 > 高可用 > 高性能 > 可靠传输。
主流 MQ 对比
| 特性 | ActiveMQ | RabbitMQ | RocketMQ | Kafka |
|---|---|---|---|---|
| 单机吞吐量 | 万级 | 万级 | 十万级 | 十万级+ |
| 时效性 | 毫秒级 | 微秒级 | 毫秒级 | 毫秒级以内 |
| 可用性 | 高(主从) | 高(主从) | 非常高(分布式) | 非常高(多副本) |
| 消息可靠性 | 有较低概率丢失 | 基本不丢 | 参数优化后可做到不丢 | 同 RocketMQ |
| 应用场景 | MQ 功能完备 | 并发强、延时低 | 在线业务 | 大数据/实时计算/日志 |
| 语言支持 | - | 非常多 | Java | Scala/Java |
- RabbitMQ:优点——语言支持多、路由配置灵活;缺点——堆积能力差、性能一般
- RocketMQ:优点——性能稳定、支持事务;缺点——国外认同度弱
- Kafka:优点——可靠稳定、大数据生态最健全;缺点——同步延迟高,不太适合在线业务
消息模型:主题和队列有什么区别?
队列模型
早期消息队列按"队列"数据结构设计:生产者入队,消费者出队。多消费者之间竞争消息,每条消息只能被一个消费者收到。
发布/订阅模型
消息发送给主题,订阅者均可收到全量消息。与队列模型的核心区别:一份消息数据能否被消费多次。
RabbitMQ 的消息模型
Exchange 位于生产者和队列之间,根据配置策略将消息投递到不同队列。
RocketMQ 的消息模型
标准发布/订阅模型。每个主题包含多个队列,通过队列实现多实例并行生产和消费。有序性只在队列层面保证。
消费位置(Consumer Offset):每个消费组在每个队列上维护一个消费位置,消费完的消息不立即删除。
Kafka 消息模型与 RocketMQ 相同,只是将 Queue 称为分区(Partition)。
如何利用事务消息实现分布式事务?
Kafka:直接抛出异常,用户自行处理(重试提交或事务补偿)。

RocketMQ:通过事务反查机制解决。Broker 定期反查 Producer 本地事务状态,根据结果决定提交/回滚。

MQ 事务方案总结
优点:消息独立存储降低耦合、吞吐量优于本地消息表方案。
缺点:一次发送需两次网络请求(half 消息 + commit/rollback)、需实现消息状态回查接口。
如何确保消息不会丢失?
检测方法:利用消息队列有序性,给消息附加递增序号,Consumer 端检查序号连续性。
确保不丢失:
- 生产阶段:捕获发送错误,针对性容错
- 存储阶段:数据设置副本,所有副本写入成功才视为提交
- 消费阶段:数据处理完毕后手动提交消费偏移量
如何处理消费过程中的重复消息?
三种服务质量:
- At most once:至多一次,允许丢失
- At least once:至少一次,不允许丢失,允许重复
- Exactly once:恰好一次,不丢失不重复
主流 MQ(RocketMQ、RabbitMQ、Kafka)均为 At least once。
解决方案:消费端幂等设计。At least once + 幂等消费 = Exactly once。
常用幂等方法:
- 数据库唯一约束:INSERT IF NOT EXIST
- 前置条件检查:乐观锁版本号机制
- 记录并检查操作:全局唯一 ID + 消费状态记录(GUID、雪花算法)
消息积压了该如何处理?
发送端性能优化
优先检查发消息前的业务逻辑耗时,设置合适的并发和批量大小。
消费端性能优化
消费性能必须高于生产性能,才能持续健康运行。可通过水平扩容增加并发数,但扩容 Consumer 实例必须同步扩容分区(队列)数量,确保两者相等。
消息积压处理
- 分析原因:发送变快 vs 消费变慢(通过监控数据判断)
- 发送陡增:扩容消费实例
- 资源不足:降级关闭非核心业务
- 消费变慢:检查日志错误、打印堆栈定位阻塞原因
学习开源代码该如何入手?
- 先看官方文档:了解项目定位、用途、使用方式、适用场景、优缺点
- 由点及面阅读源码:带着问题和目的性阅读,避免泛泛而读
如何使用异步设计提升系统性能?
异步编程减少线程等待,但增加程序复杂度。常用框架:CompletableFuture(Java8 内置)、RxJava。
如何实现高性能的异步网络传输?
- IO 密集型应用关键:减少 IO 等待时间
- 理想模型:并发 IO + 线程复用(NIO 多路复用)
- 代表框架:Netty
传输协议:应用程序之间对话的语言
传输协议定义信息规则,使收发双方能互相交流。复杂协议可定义为 TLV 结构。
Kafka 如何实现高性能 IO?
- 批量消息:提升服务端处理能力
- 顺序读写:省去大部分磁盘寻址时间
- PageCache:利用操作系统文件缓存,大部分消费读命中 PageCache,减少磁盘 IO
- 零拷贝:调用
sendfileAPI,减少一次数据复制,DMA 直接完成,无需 CPU 参与
数据压缩:时间换空间的游戏
压缩不仅节省存储空间,还可提升网络传输性能。本质是 CPU 时间换磁盘/网络资源。
常用算法:ZIP、GZIP、SNAPPY、LZ4。选择时需综合考虑压缩速度和压缩率。
Kafka 压缩策略:Producer 端压缩 → Broker 端保持 → Consumer 端解压。
RocketMQ Producer 源码分析
核心类:
DefaultMQProducerImpl:封装大部分 Producer 业务逻辑MQClientInstance:封装客户端通用业务逻辑MQClientAPIImpl:封装客户端与服务端 RPCNettyRemotingClient:实现底层网络通信
发送方式:单向、同步、异步。
Kafka Consumer 源码分析
消费模型要点:
- 每个 Consumer 属于一个 ConsumerGroup
- ConsumerGroup 中每个 Consumer 独占一个或多个 Partition
- 每个 Partition 至多有 1 个 Consumer 在消费
- Coordinator 负责分配 Consumer 和 Partition 对应关系
- Consumer 维护心跳,故障时触发 rebalance
请求构建:暂存发送队列,批量发送(poll() 方法中统一处理)。
Kafka 和 RocketMQ 的消息复制差异
RocketMQ 复制
- 复制单位:Broker(主从模式,一主一从或一主多从)
- 传统主从模式:性能好,灵活性和可用性差
- Dledger 模式:半数以上节点复制成功才返回,支持自动选举,可用性更好
Kafka 复制
- 复制单位:分区(每个分区的副本构成小复制集群)
- Broker 不分主从,只是分区副本的容器
- ISR(In Sync Replicas):保持数据同步的副本集合(含主节点)
- 使用 ZooKeeper 监控和选举
RocketMQ NameServer

- 独立进程,为 Broker/Producer/Consumer 提供寻址服务
- 多节点部署,节点间互不通信,每个节点独立提供全部服务
- Broker 主动通知所有 NameServer 更新路由信息,同时充当心跳
核心类 RouteInfoManager 中保存所有路由信息(内存存储,无持久化):
public class BrokerData implements Comparable<BrokerData> {
private final HashMap<String/* topic */, List<QueueData>> topicQueueTable;
private final HashMap<String/* brokerName */, BrokerData> brokerAddrTable;
private final HashMap<String/* clusterName */, Set<String/* brokerName */>> clusterAddrTable;
private final HashMap<String/* brokerAddr */, BrokerLiveInfo> brokerLiveTable;
private final HashMap<String/* brokerAddr */, List<String>/* Filter Server */> filterServerTable;
}Kafka 的协调服务 ZooKeeper
ZooKeeper 核心服务:高可用、高可靠的一致性存储,提供元数据读写、节点监控、选举、分布式锁等功能。
使用注意:
- 不要写入大量数据(超过几百 MB 性能和稳定性下降)
- 不要让业务可用性依赖 ZK 可用性(ZK 选举慢、对网络抖动敏感)
Kafka 使用 ZK:保存元数据(Broker 列表 + 主题分区信息)、监控 Broker 存活、选举。
RocketMQ 与 Kafka 事务对比
- 均基于两阶段提交实现,利用特殊主题队列/分区记录事务日志
- RocketMQ:消息暂存特殊队列,提交后移到业务队列;适用于本地事务 + 消息一致性
- Kafka:消息直接放业务分区,配合客户端过滤;适用于Exactly Once 机制和实时计算
MQTT 协议:支持海量 IoT 设备
专为物联网设计,IoT 设备性能差、网络不稳定。服务端挑战:支撑海量客户端和主题。
自建集群关键:前置 Proxy 集群解决海量连接、会话管理、海量主题三个问题。


Pulsar 的存储计算分离设计

核心区别:存储职责从 Broker 分离到 BookKeeper 存储集群,Broker 变为无状态节点。
优点:调度灵活、故障转移简单快速。缺点:系统复杂度更高、性能略差。
流计算与消息
Flink 流计算原理
JobGraph(有向无环图)→ Task → SubTask(线程级并行),由 JobManager 分配给 TaskManager 执行。
端到端 Exactly Once
- Flink:CheckPoint + Barrier 机制恢复状态
- Kafka:事务 + 生产幂等
- 配合方式:每个 CheckPoint 对应一个 Kafka 事务,通过两阶段提交保证状态一致性
主流消息队列消息存储对比
Kafka 存储
- 以 Partition 为单位,每个 Partition 包含消息文件(Segment file)+ 索引文件(Index)
- 稀疏索引:每隔几条消息创建一条索引,节省存储空间
- 写入:尾部连续追加;查找:二分法遍历索引 + 顺序遍历消息文件
RocketMQ 存储
- 以 Broker 为单位,所有主题消息写入同一组消息文件
- 定长稠密索引(ConsumerQueue):每条消息都有索引,固定 20 字节
- 查找:直接计算索引全局位置(索引序号 × 20),两次绝对位置寻址
存储结构对比

- Kafka:分区粒度细,灵活,易数据迁移和扩容;小消息场景稀疏索引节省空间
- RocketMQ:Broker 粒度粗,写入批量更大,上千活动主题场景写入性能更优