客户端只连 Proxy
Python 5.x SDK 使用 gRPC,入口是 127.0.0.1:8081。不要把 NameServer 的 9876 当作 SDK endpoint。
这不是脱离代码的概念图。页面严格对应当前 compose.yaml、producer.py 和 consumer.py:先看本机实验拓扑,再跟踪一条消息如何找到 Broker、落盘、被消费并 ACK。
Python 程序运行在宿主机;RocketMQ 服务运行在 Docker 网络。Producer 与 Consumer 不直接使用 NameServer 或 Broker 的 Remoting 端口,而是访问 Proxy 的 gRPC 端口 8081。Dashboard 通过 NameServer 发现集群,并经本地桥接访问 Broker。
Python 5.x SDK 使用 gRPC,入口是 127.0.0.1:8081。不要把 NameServer 的 9876 当作 SDK endpoint。
Broker 定期注册路由;客户端查询 Topic 在哪个 Broker。它不是消息存储,也不承载消息正文。
Broker 维护 Topic 队列、顺序写 CommitLog,并派生 ConsumeQueue 供消费定位。
RocketMQ 的逻辑模型可以从“大类、分片、协作、实例、消息”五层理解。这里的 Queue 是 Topic 的分区,不是另一个独立服务。
消息的一级分类。Producer 向 Topic 发送,Consumer 订阅 Topic。本实验使用 HelloWorldTopic。
Topic 在 Broker 内的逻辑分片。提高并行度;有序消息需要把同一业务键路由到同一 Queue。
同组消费者协作分摊消息;不同组可以各自消费同一条消息,进度彼此独立。
示例用 Tag hello 做服务端过滤,用 Key 标记业务消息;两者都不是 Topic。
生产者的逻辑身份,常用于事务等场景。普通消息的关键目标仍是 Topic。
Broker 按保留策略保存消息。ACK 更新消费状态,不会马上从 CommitLog 删除正文。
producer.startup() 建立客户端;Producer 查询路由、选择 Queue,然后发送。成功回执中的 Message ID 表示服务端已接受消息;当前 Broker 使用异步刷盘,回执不等于磁盘介质已经同步完成。
NameServer 只提供路由。5.x Python SDK 使用 gRPC,因此实际入口是 Proxy :8081;Proxy 再协调 Broker。
可靠性取决于刷盘、复制、重试与业务幂等。本实验是单 Broker + ASYNC_FLUSH,适合学习,不是生产高可用配置。
SimpleConsumer 发起长轮询。Broker 找到消息后返回正文和 receipt handle,并让消息在 15 秒内对同组其他消费者不可见。业务成功后 ACK;若进程崩溃或未 ACK,超时后消息可再次投递,所以消费逻辑要幂等。
先 ACK 再处理会在业务失败时丢失重试机会。示例先打印正文,再调用 consumer.ack(message)。
网络超时可能让 ACK 已成功但客户端不知道。用业务 Key、唯一约束或幂等表防止重复副作用。
换一个 Consumer Group,会得到独立消费视角;同一个 Group 的多个实例则分摊队列。
Broker 采用顺序写 CommitLog 保存完整消息,再构建面向 Topic/Queue 的 ConsumeQueue。IndexFile 用于按 Key 等条件查询。Docker named volume 让容器重建后数据仍保留;make clean 才会连 volume 一起删除。
观察 NameServer 先启动、Broker 注册、Topic 和消费组初始化。
对照终端中的 Message ID,看 Producer 发送、Consumer 接收并 ACK。
体会长轮询:Consumer 可以先等待,消息到达后立即返回。
查看 DefaultCluster、HelloWorldTopic、hello-world-group 和消息查询。
临时把 consumer.ack(message) 注释掉,并缩短实验等待。观察不可见窗口过后是否再次收到。
区分“容器被删”和“持久数据被删”。
Python 使用 RocketMQ 5.x gRPC SDK,直接连接 Proxy。9876 是 NameServer 的路由服务端口,不接收 SDK 的消息发送和消费请求。
不会。消息正文位于 Broker 存储。Broker 会重新向 NameServer 注册路由;NameServer 不保存业务消息。
通常不会,它们共同分摊队列与消息;换成两个不同 Consumer Group,才会获得两份独立消费视角。
ACK 更新的是消费状态,不是立即删除 CommitLog 中的正文。消息由 Broker 的保留和清理策略管理。
至少一次投递允许重试。处理成功但 ACK 响应丢失、进程崩溃或超时,都可能产生重复投递。