RocketMQ
RocketMQ 是以日志存储为核心的消息系统:消息先进入物理 CommitLog,再异步构建消费队列和索引,消费端通过逻辑队列定位物理消息。
本册问题地图
Section titled “本册问题地图”把 RocketMQ 当成一条“消息从写入到被确认”的流水线来读:
- 为什么先写 CommitLog,再生成 ConsumeQueue? 因为顺序追加适合吞吐,而按 Topic/Queue 查找需要轻量索引;两者分开,写入热路径不用为每个消费队列做随机写。
- 消息什么时候算“写成功”,什么时候算“消费者可见”? 写入物理日志、刷盘、索引重放、拉取和消费位点是不同边界,任何一步滞后都会造成“已写入但暂时拉不到”。
- 长轮询到底解决了什么? 它不是让服务端主动推送,而是让 Broker 暂存没有匹配消息的请求,直到消息到达、超时或连接异常再返回。
- 重平衡如何改变消费权? 消费者组重新分配队列后,旧消费者必须停止对应 ProcessQueue;否则同一队列可能出现重复消费或位点覆盖。
- 事务消息为什么要先写半消息? 先把业务消息隐藏起来,再依据本地事务结果提交、回滚或接受回查,避免消费者看到尚未确定的状态。
- DLedger 与传统主从的差别是什么? 传统复制更像主节点向从节点同步,DLedger 把日志复制、选主和提交位置纳入更严格的多数派协议,故障切换边界不同。
| 项 | 值 |
|---|---|
| 仓库 | apache/rocketmq |
| 本地路径 | E:\source\java\mq\rocketmq |
| 分支 | develop |
| Commit | 293f5885719fc4aa3619446a1900f58ccfcfdd29(2026-08-14) |
| 最近 tag | rocketmq-all-5.5.0 |
本册所有源码坐标均基于上述 commit。
它解决什么问题
Section titled “它解决什么问题”- 把高吞吐追加写与按 Topic/Queue 消费的随机访问拆开。
- 在消费者暂时没有消息时,用长轮询减少空拉取。
- 用半消息、回查和副本协议处理事务可见性与节点故障。
| 模块 | 职责 | 关键抽象 |
|---|---|---|
store |
物理日志、消费索引、刷盘与复制 | CommitLog、ConsumeQueue、MappedFileQueue |
broker |
拉取、长轮询、事务和请求协议 | PullMessageProcessor、PullRequestHoldService |
client |
生产、消费、重平衡 | RebalanceImpl、ProcessQueue |
dledger |
复制与选主 | DLedgerCommitLog |
读码入口顺序
Section titled “读码入口顺序”DefaultMessageStore#asyncPutMessage→CommitLog#asyncPutMessage。ReputMessageService→ConsumeQueue/IndexFile。DefaultMessageStore#getMessage→PullMessageProcessor#processRequest。PullRequestHoldService→NotifyMessageArrivingListener。TransactionalMessageServiceImpl→EndTransactionProcessor。DLedgerCommitLog→DLedgerRoleChangeHandler。
| 对象 | 对照点 |
|---|---|
| Kafka | 都以追加日志为核心;RocketMQ 把消费队列作为独立索引 |
| Redis Streams | 都保存消费位点;RocketMQ 的物理日志和逻辑队列分层更明显 |
| Spring Cloud Alibaba RocketMQ | 上层 Binder 负责 Spring 适配,本册关注 Broker、Store 和 Client 内核 |