ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

RocketMQ架构精讲:NameServer、Broker、Producer、Consumer核心流程

RocketMQ架构精讲:NameServer、Broker、Producer、Consumer核心流程 RocketMQ架构精讲NameServer、Broker、Producer、Consumer核心流程作者黒漂技术佬系列专栏RocketMQ核心原理与无人售货柜项目实战一、先看全景RocketMQ整体架构把RocketMQ的架构想象成一个快递物流网络有四个核心角色┌─────────────────┐ │ NameServer │ ← 注册中心物流信息中心 │ (路由信息存储) │ └────────┬────────┘ │ 心跳注册/路由查询 ┌──────────────┼──────────────┐ │ │ │ ┌─────▼─────┐ ┌─────▼─────┐ ┌────▼──────┐ │ Broker-A │ │ Broker-B │ │ Broker-C │ ← 消息存储服务器快递分拣中心 │ (Master) │ │ (Master) │ │ (Master) │ │ (Slave) │ │ (Slave) │ │ (Slave) │ └─────┬─────┘ └─────┬─────┘ └────┬──────┘ │ │ │ │ 发送消息 │ │ ┌─────┴──────────────┴──────────────┴─────┐ │ Producer │ ← 消息生产者寄件人 │ (订单服务/支付服务/...) │ └──────────────────────────────────────────┘ ┌──────────────────────────────────────────┐ │ Consumer │ ← 消息消费者收件人 │ (库存服务/推送服务/告警服务) │ └──────────────────────────────────────────┘ │ │ │ └──────────────┼──────────────┘ │ 拉取消息 ┌────────▼────────┐ │ Broker │ └─────────────────┘四个角色的分工一句话概括NameServer路由信息中心告诉Producer消息该往哪发告诉Consumer消息该从哪取Broker消息存储服务器真正存消息和转发消息的地方Producer消息生产者发消息的Consumer消息消费者收消息的下面逐个拆解。二、NameServer轻量级注册中心2.1 它是干什么的NameServer是RocketMQ的通讯录。Broker启动时把自己的地址、Topic信息、队列信息注册到NameServerProducer和Consumer启动时从NameServer拉取这份路由表才知道消息往哪发、从哪取。2.2 为什么不用Zookeeper很多中间件Kafka、Dubbo用Zookeeper做注册中心RocketMQ早期版本也用过后来自己搞了个NameServer原因是一个经典的分布式理论取舍AP vs CP。先解释下CAP定理的三个字母CConsistency一致性所有节点看到的数据是一致的AAvailability可用性每个请求都能收到响应不保证是最新数据PPartition tolerance分区容错网络分区时系统仍能运行分布式系统三者只能选其二由于网络分区P不可避免实际选择是CP还是AP。对比项Zookeeper (CP)NameServer (AP)一致性强一致最终一致可用性写入需多数派投票leader选举期间不可用各节点独立互不通信复杂度高ZAB协议、leader选举低就是个内存HashMap部署至少3节点集群可单节点也可多节点互不关联RocketMQ的选择逻辑在消息中间件场景里可用性比强一致更重要。Broker挂了NameServer感知到就行不追求所有NameServer瞬间数据一致。每个NameServer节点独立维护路由表Producer/Consumer连任意一个都能拿到数据偶尔短暂不一致在业务上可接受Broker心跳30秒上报一次最多30秒感知到变化。2.3 心跳机制Broker每隔30秒向所有NameServer发心跳NameServer收到后更新Broker的存活时间戳。NameServer每隔10秒扫描一次路由表如果发现某个Broker超过120秒没心跳就判定它下线从路由表里移除。这套机制简单粗暴但有效不搞复杂的一致性协议用心跳超时来维护状态。三、Broker消息存储和转发服务器3.1 Master/Slave架构Broker分Master和Slave两种角色Master Broker负责读写Producer发消息到MasterConsumer也从Master拉消息默认Slave Broker只负责备份从Master同步数据Master挂了可接管读请求同步方式有两种异步复制ASYNC_MASTERMaster收到消息后立即返回成功然后异步复制给Slave。优点是延迟低缺点是Master挂了可能丢少量未复制的数据。同步双写SYNC_MASTERMaster收到消息后等Slave也写入成功才返回。优点是数据不丢缺点是延迟略高。售货柜场景怎么选支付相关消息用同步双写保数据安全日志类消息用异步复制换性能。3.2 Broker的存储结构Broker存消息不是随便往磁盘一堆而是有精心的设计CommitLog统一存储文件 ├── 所有Topic的消息混在一起按写入顺序追加 ├── 每个文件固定1GB写满新建 └── 顺序写磁盘性能接近内存写 ConsumeQueue消费队列/逻辑索引 ├── 每个Topic的每个Queue对应一个ConsumeQueue文件 ├── 存储的是消息在CommitLog中的偏移量(offset)、大小、Tag的hashcode └── Consumer通过ConsumeQueue快速定位消息 IndexFile索引文件 └── 支持按MessageId和Key查询消息这个设计的精妙之处在于写入时所有消息追加到同一个CommitLog顺序写性能极高读取时通过ConsumeQueue索引快速定位读取效率也不低。写入和读取两不误。3.3 Broker注册到NameServerBroker启动时把自己的信息打包发给所有NameServer注册信息包括Broker名称brokerNameBroker地址IP:Port集群名称clusterNameMaster/Slave角色标识该Broker上所有的Topic和Queue信息注册后每30秒心跳续约心跳时也会带上最新的Topic路由信息。四、Producer消息生产者4.1 启动流程Producer启动时做三件事从NameServer拉取路由信息哪些Broker有哪些Topic的哪些Queue建立到目标Broker的网络连接Netty长连接启动定时任务每30秒从NameServer更新路由信息路由信息缓存在本地发消息时不用每次都问NameServer30秒更新一次够用了。4.2 消息发送流程Producer发送一条消息的完整流程 1. Producer校验消息Topic、Body非空等 2. 查找Topic的路由信息本地缓存 3. 选择目标Queue负载均衡策略 ├── 轮询默认Round Robin ├── 故障隔离上次发送失败的Queue暂时排除 └── 指定QueueMessageQueueSelector自定义 4. 构建请求通过网络发送给目标Broker 5. 等待Broker返回确认根据发送方式决定是否阻塞 ├── 同步发送阻塞等待ACK ├── 异步发送不阻塞回调通知结果 └── 单向发送直接返回不等 6. 根据结果做后续处理重试/记录日志等售货柜场景举例关门后发送订单创建消息用同步发送确保消息不丢发送用户行为日志用单向发送丢了也无所谓。4.3 发送失败重试Producer发送失败时会自动重试默认重试2次共3次尝试。重试时会避开上次失败的Broker选择其他Broker的Queue发送。这个机制配合Broker的Master/Slave架构能在单节点故障时自动切换对用户透明。五、Consumer消息消费者5.1 两种消费模式集群消费Clustering同一个ConsumerGroup下的多个Consumer实例分摊消费所有消息。比如100条消息2个Consumer各消费50条。这是最常用的模式适合横向扩展消费能力。ConsumerGroup-A (集群模式) ├── Consumer-1 ← 消费 Queue-0, Queue-1 ├── Consumer-2 ← 消费 Queue-2, Queue-3 └── Consumer-3 ← (如果只有4个Queue这个会空闲)广播消费Broadcasting同一个ConsumerGroup下的每个Consumer实例都消费全量消息。比如100条消息2个Consumer各消费100条。适合所有节点都要感知同一份数据的场景比如本地缓存刷新。售货柜场景库存扣减用集群消费一台柜子的消息只需一个消费者处理设备配置全量下发用广播消费所有柜子都要收到配置更新。5.2 ConsumerGroup的概念ConsumerGroup是逻辑上的一组消费者实例同一个Group必须消费同一个Topic且消费逻辑必须一致。你可以把它理解为一个消费团队团队里的人分工干活集群模式或者各自干全部活广播模式。一个典型的微服务部署库存服务部署3个实例都配置同一个ConsumerGroupinventory_consumer_group3个实例自动分摊消费库存Topic的消息。某个实例挂了Rebalance机制会自动把它负责的Queue分配给其他实例消费不中断。5.3 消费者拉取机制RocketMQ的Consumer本质上是Pull模式——消费者主动从Broker拉取消息不是Broker推过来的。但用起来感觉像PushDefaultPushConsumer因为它内部用长轮询做了封装Consumer向Broker发拉取请求 ├── Broker有消息 → 立即返回 └── Broker没消息 → Hold住请求挂起5秒默认 ├── 5秒内有新消息 → 立即返回 └── 5秒后还没消息 → 返回空Consumer再次发起拉取长轮询的好处既避免了Push模式下Consumer处理不过来被压垮又避免了Pull模式下频繁空轮询浪费资源。六、消息流转全链路把上面四个角色串起来一条消息从产生到消费的完整链路1. Producer启动 → 从NameServer获取Topic路由信息 2. Producer发送消息 → 根据路由选择Queue → 发送到对应Broker 3. Broker收到消息 → 写入CommitLog → 构建ConsumeQueue索引 → 返回ACK 4. Consumer启动 → 从NameServer获取Topic路由信息 5. Consumer向Broker发拉取请求 → 通过ConsumeQueue定位消息 → 从CommitLog读取 6. Consumer处理消息 → 返回消费状态(CONSUME_SUCCESS/RECONSUME_LATER) 7. 消费失败 → Broker按延迟等级重新投递 → 超过重试上限进入死信队列七、结合售货柜场景的角色映射把架构映射到我们的无人售货柜项目RocketMQ角色售货柜项目对应具体说明NameServer部署在总部机房所有服务都连总部NameServer获取路由Broker-A总部主Broker存储订单、支付、库存等核心消息Broker-B区域分中心Broker按区域分Topic减少跨地域网络延迟Producer订单服务、支付服务、柜子网关各业务服务发消息到BrokerConsumer库存服务、推送服务、告警服务、ERP同步服务各服务消费各自Topic比如一个完整的支付流程用户关门 →订单服务(Producer)发送订单创建消息到Broker库存服务(Consumer)消费消息扣减本地库存用户支付完成 →支付服务(Producer)发送支付成功消息柜子网关(Consumer)消费消息通知对应柜子出货ERP同步服务(Consumer)消费同一消息更新总部库存推送服务(Consumer)消费同一消息给用户发扣款通知一条支付消息被4个消费者服务各自消费不同ConsumerGroup各司其职互不干扰。这就是解耦的威力——支付服务只管发消息谁来消费、怎么消费它一概不关心。八、小结这一篇我们拆解了RocketMQ的四大核心组件NameServer作为轻量级AP注册中心Broker负责消息存储和转发Master/Slave架构CommitLog/ConsumeQueue存储设计Producer和Consumer通过长轮询完成消息的发送和拉取。最后把这些角色映射到无人售货柜项目的实际部署中。下一篇我们动手实操搭建RocketMQ开发环境。
返回列表