ROCKETMQ 5 / 从一条消息开始

RocketMQ 架构、路由与消息生命周期,一张张图看懂。

这不是脱离代码的概念图。页面严格对应当前 compose.yaml、producer.py 和 consumer.py:先看本机实验拓扑,再跟踪一条消息如何找到 Broker、落盘、被消费并 ACK。

约 25 分钟5 张架构 / 流程图RocketMQ 5.3.2Python gRPC SDK可离线打开
01 / LOCAL LAB TOPOLOGY

先定位:每个进程到底跑在哪里?

Python 程序运行在宿主机;RocketMQ 服务运行在 Docker 网络。Producer 与 Consumer 不直接使用 NameServer 或 Broker 的 Remoting 端口,而是访问 Proxy 的 gRPC 端口 8081。Dashboard 通过 NameServer 发现集群,并经本地桥接访问 Broker。

RocketMQ 本地实验环境全景宿主机上的 Python Producer 和 Consumer 访问 Docker 中的 Proxy;Proxy 与 Broker 同进程,Broker 注册到 NameServer 并将数据写入 Docker Volume;Dashboard 查询 NameServer 和 Broker。 宿主机 · macOS / Linux Docker Compose 网络 · rocketmq-lab_rocketmq Producerproducer.py · SEND Consumerconsumer.py · ACK RocketMQ ProxygRPC :8081路由查询 / 发送 / 消费 / ACK与 Broker 同一个容器进程 Broker · broker-aRemoting :10911队列、消息存储、消费进度DefaultCluster / ASYNC_MASTER broker-storeDocker named volumeCommitLog / ConsumeQueue / Index NameServer路由注册 / 查询 · :9876不保存业务消息 Dashboard Bridgesocat · 127.0.0.1:10911 → broker Dashboard宿主机 :18080集群 / Topic / 消费组 / 消息 gRPC 请求 :8081消息 + receipt handle内部调用顺序写入注册 Broker查询运行状态

客户端只连 Proxy

Python 5.x SDK 使用 gRPC,入口是 127.0.0.1:8081。不要把 NameServer 的 9876 当作 SDK endpoint。

NameServer 是“地图”

Broker 定期注册路由;客户端查询 Topic 在哪个 Broker。它不是消息存储,也不承载消息正文。

消息最终进入 Broker

Broker 维护 Topic 队列、顺序写 CommitLog,并派生 ConsumeQueue 供消费定位。

02 / CORE CONCEPTS

五个核心对象:别把“主题、队列、组”混在一起。

RocketMQ 的逻辑模型可以从“大类、分片、协作、实例、消息”五层理解。这里的 Queue 是 Topic 的分区,不是另一个独立服务。

RocketMQ 核心对象关系Topic 包含多个 MessageQueue,Producer Group 下的生产者发送消息,Consumer Group 下的消费者共同分配队列。Producer Group · 逻辑分组Producer A发送实例Producer B发送实例Topic · HelloWorldTopicMessageQueue 0MessageQueue 1… Queue 2 ~ 7同一 Queue 内可保持顺序Consumer Group · hello-world-groupConsumer A分配部分 QueueConsumer B分担消费Message负载分配Tag 是 Topic 内的二级分类;Key 是业务检索标识;Message ID 是系统标识。

Topic

消息的一级分类。Producer 向 Topic 发送,Consumer 订阅 Topic。本实验使用 HelloWorldTopic。

MessageQueue

Topic 在 Broker 内的逻辑分片。提高并行度;有序消息需要把同一业务键路由到同一 Queue。

Consumer Group

同组消费者协作分摊消息;不同组可以各自消费同一条消息,进度彼此独立。

Tag / Key

示例用 Tag hello 做服务端过滤,用 Key 标记业务消息;两者都不是 Topic。

Producer Group

生产者的逻辑身份,常用于事务等场景。普通消息的关键目标仍是 Topic。

消费不是“取走即删除”

Broker 按保留策略保存消息。ACK 更新消费状态,不会马上从 CommitLog 删除正文。

03 / PRODUCE PATH

发送一条消息,要先问路,再把正文交给 Broker。

producer.startup() 建立客户端;Producer 查询路由、选择 Queue,然后发送。成功回执中的 Message ID 表示服务端已接受消息;当前 Broker 使用异步刷盘,回执不等于磁盘介质已经同步完成。

RocketMQ Producer 发送链路Producer 经 Proxy 查询 NameServer 路由,选择消息队列,Proxy 将消息发送给 Broker,Broker 追加 CommitLog 并返回 Message ID。① 构造 MessageTopic / Tag / Key / Body② Proxy 查询路由Topic 在哪里?③ NameServer返回 broker-a + Queue④ Broker 接收校验并分配物理偏移⑤ CommitLog顺序追加消息正文⑥ SendReceiptmessage_idQueryRouteSendMessageASYNC_FLUSH成功回执producer.py 循环发送时,每一轮都会得到独立 Message ID;Queue 选择决定并行度与局部顺序。
为什么不直接连 NameServer:9876 发送?

NameServer 只提供路由。5.x Python SDK 使用 gRPC,因此实际入口是 Proxy :8081;Proxy 再协调 Broker。

发送成功就是绝对不丢吗?

可靠性取决于刷盘、复制、重试与业务幂等。本实验是单 Broker + ASYNC_FLUSH,适合学习,不是生产高可用配置。

04 / CONSUME & ACK

消费是“领取一段可见性租约”,ACK 才完成本次处理。

SimpleConsumer 发起长轮询。Broker 找到消息后返回正文和 receipt handle,并让消息在 15 秒内对同组其他消费者不可见。业务成功后 ACK;若进程崩溃或未 ACK,超时后消息可再次投递,所以消费逻辑要幂等。

SimpleConsumer 消费和 ACK 生命周期Consumer 长轮询 Broker,Broker 返回消息并设置不可见时间,业务处理成功后 ACK;失败或超时则消息重新可见并可能重复投递。① Receive最多 16 条 · 长轮询② Broker POP按 Group + Queue 查找返回 BodyMessage IDReceipt Handle③ 业务处理打印 / 写库 / 调接口④ ACK确认本次消费完成⑤ 更新消费状态同组不再收到该投递重新可见可能重复投递长轮询15 秒 invisible处理成功确认 receipt handle未 ACK / 超时 / 崩溃至少一次语义意味着:宁可重投,也不能假设每条业务只执行一次。

ACK 在业务成功之后

先 ACK 再处理会在业务失败时丢失重试机会。示例先打印正文,再调用 consumer.ack(message)。

必须考虑重复

网络超时可能让 ACK 已成功但客户端不知道。用业务 Key、唯一约束或幂等表防止重复副作用。

不同 Group 各消费一次

换一个 Consumer Group,会得到独立消费视角;同一个 Group 的多个实例则分摊队列。

05 / STORAGE & FAILURE BOUNDARY

消息正文、消费索引和进度,不是同一份东西。

Broker 采用顺序写 CommitLog 保存完整消息,再构建面向 Topic/Queue 的 ConsumeQueue。IndexFile 用于按 Key 等条件查询。Docker named volume 让容器重建后数据仍保留;make clean 才会连 volume 一起删除。

RocketMQ Broker 存储结构消息先顺序写入 CommitLog,再异步分发为 ConsumeQueue 与 IndexFile;消费进度和消息正文分离。Message完整正文 + 属性CommitLog所有 Topic 顺序追加物理偏移 / 完整消息ConsumeQueueTopic + Queue 逻辑索引指向 CommitLog 偏移IndexFileKey / 时间等查询索引不是消费的主路径Consumer StateGroup 的消费进度 / POP 状态ACK 改状态,不删正文appenddispatch构建索引定位消息broker-store volume → /home/rocketmq/store → 容器重建后仍保留
06 / HANDS-ON PATH

按这个顺序动手,边看图边验证。

1. 启动与观察

make up make status make logs

观察 NameServer 先启动、Broker 注册、Topic 和消费组初始化。

2. 跑通消息闭环

make demo COUNT=3

对照终端中的 Message ID,看 Producer 发送、Consumer 接收并 ACK。

3. 分开两个终端

# 终端 A make consumer # 终端 B make producer COUNT=5

体会长轮询:Consumer 可以先等待,消息到达后立即返回。

4. 打开 Dashboard

open http://localhost:18080

查看 DefaultCluster、HelloWorldTopic、hello-world-group 和消息查询。

5. 模拟未 ACK

临时把 consumer.ack(message) 注释掉,并缩短实验等待。观察不可见窗口过后是否再次收到。

6. 理解数据生命周期

make down # 保留 volume make up # 消息数据仍在 make clean # 删除 volume

区分“容器被删”和“持久数据被删”。

RocketMQ 实验端口地图9876 是 NameServer,8081 是 Python gRPC Proxy,10911 是 Broker Remoting,18080 是 Dashboard。9876NameServer 路由mqadmin / Dashboard8081Proxy gRPCPython SDK endpoint10911Broker RemotingDashboard / mqadmin18080Dashboard HTTP浏览器访问
07 / CHECK YOUR UNDERSTANDING

能回答这五题,就完成了第一轮。

Python endpoint 为什么是 8081,而不是 9876?

Python 使用 RocketMQ 5.x gRPC SDK,直接连接 Proxy。9876 是 NameServer 的路由服务端口,不接收 SDK 的消息发送和消费请求。

NameServer 重启后,已经落盘的消息会消失吗?

不会。消息正文位于 Broker 存储。Broker 会重新向 NameServer 注册路由;NameServer 不保存业务消息。

两个同组 Consumer 会各收到一份消息吗?

通常不会,它们共同分摊队列与消息;换成两个不同 Consumer Group,才会获得两份独立消费视角。

ACK 后消息为什么还能在 Dashboard 查询?

ACK 更新的是消费状态,不是立即删除 CommitLog 中的正文。消息由 Broker 的保留和清理策略管理。

为什么业务代码仍要做幂等?

至少一次投递允许重试。处理成功但 ACK 响应丢失、进程崩溃或超时,都可能产生重复投递。