Skip to content

RocketMQ

This content is not available in your language yet.

RocketMQ 是以日志存储为核心的消息系统:消息先进入物理 CommitLog,再异步构建消费队列和索引,消费端通过逻辑队列定位物理消息。

把 RocketMQ 当成一条“消息从写入到被确认”的流水线来读:

  1. 为什么先写 CommitLog,再生成 ConsumeQueue? 因为顺序追加适合吞吐,而按 Topic/Queue 查找需要轻量索引;两者分开,写入热路径不用为每个消费队列做随机写。
  2. 消息什么时候算“写成功”,什么时候算“消费者可见”? 写入物理日志、刷盘、索引重放、拉取和消费位点是不同边界,任何一步滞后都会造成“已写入但暂时拉不到”。
  3. 长轮询到底解决了什么? 它不是让服务端主动推送,而是让 Broker 暂存没有匹配消息的请求,直到消息到达、超时或连接异常再返回。
  4. 重平衡如何改变消费权? 消费者组重新分配队列后,旧消费者必须停止对应 ProcessQueue;否则同一队列可能出现重复消费或位点覆盖。
  5. 事务消息为什么要先写半消息? 先把业务消息隐藏起来,再依据本地事务结果提交、回滚或接受回查,避免消费者看到尚未确定的状态。
  6. DLedger 与传统主从的差别是什么? 传统复制更像主节点向从节点同步,DLedger 把日志复制、选主和提交位置纳入更严格的多数派协议,故障切换边界不同。
项 值
仓库 apache/rocketmq
本地路径 E:\source\java\mq\rocketmq
分支 develop
Commit 293f5885719fc4aa3619446a1900f58ccfcfdd29(2026-08-14)
最近 tag rocketmq-all-5.5.0

本册所有源码坐标均基于上述 commit。

  • 把高吞吐追加写与按 Topic/Queue 消费的随机访问拆开。
  • 在消费者暂时没有消息时,用长轮询减少空拉取。
  • 用半消息、回查和副本协议处理事务可见性与节点故障。
模块 职责 关键抽象
store 物理日志、消费索引、刷盘与复制 CommitLog、ConsumeQueue、MappedFileQueue
broker 拉取、长轮询、事务和请求协议 PullMessageProcessor、PullRequestHoldService
client 生产、消费、重平衡 RebalanceImpl、ProcessQueue
dledger 复制与选主 DLedgerCommitLog
  1. DefaultMessageStore#asyncPutMessage → CommitLog#asyncPutMessage。
  2. ReputMessageService → ConsumeQueue / IndexFile。
  3. DefaultMessageStore#getMessage → PullMessageProcessor#processRequest。
  4. PullRequestHoldService → NotifyMessageArrivingListener。
  5. TransactionalMessageServiceImpl → EndTransactionProcessor。
  6. DLedgerCommitLog → DLedgerRoleChangeHandler。
对象 对照点
Kafka 都以追加日志为核心;RocketMQ 把消费队列作为独立索引
Redis Streams 都保存消费位点;RocketMQ 的物理日志和逻辑队列分层更明显
Spring Cloud Alibaba RocketMQ 上层 Binder 负责 Spring 适配,本册关注 Broker、Store 和 Client 内核