RocketMQ的运行架构和消息模型可以总结为"一个中心化的存储集群 + 一套标准化的发布订阅模型"。它通过NameServer做服务发现,Broker做消息存储,再配合生产者和消费者这两大客户端,构成了一个高吞吐、低延迟的分布式消息系统。

一、核心架构组件

  • NameServer轻量级路由中心。它负责管理Broker集群的元数据(如Broker地址、Topic路由信息),为生产者和消费者提供服务发现功能。。

  • Broker消息存储核心。这是RocketMQ最关键的组件,负责消息的接收、存储、查询和消费投递。Broker通常采用主从(Master-Slave)架构部署,Master负责处理写入请求,Slave负责高可用和部分读请求。单个Broker节点可以承载成千上万个Topic,这得益于其优秀的存储设计。

  • Producer(生产者)消息发布者。它是业务系统中负责创建并发送消息的客户端。生产者会定期从NameServer拉取Topic的路由信息,然后与目标Broker建立长连接,并将消息发送到该Topic下的某个队列中。

  • Consumer(消费者)消息订阅者。它是业务系统中负责处理消息的客户端。消费者同样从NameServer获取路由,然后向Broker发起拉取请求,获取消息进行业务处理。同一个消费组(ConsumerGroup)内的多个消费者会共同分担消息,实现水平扩展。

二、核心消息模型

  • 主题(Topic)逻辑上的消息分类。它是消息的第一级容器,用于区分不同的业务,例如订单Topic、商品Topic。一个Topic下会包含多个消息队列。

  • 消息队列(MessageQueue)物理存储分片。Topic由多个Queue组成,每个Queue内部的消息是严格有序的。通过增加Queue的数量,可以水平提升Topic的吞吐能力。

  • 消息(Message)数据传输的最小单元。每条消息都归属于一个Topic和一个Queue,拥有唯一的ID(Message ID)和可选的Key(业务键,用于快速查找)与Tag(标签,用于消息过滤)。

  • 消费者分组(ConsumerGroup)消费身份逻辑分组。这是发布订阅模型的核心。同一个Group内的多个消费者实例共同消费Topic中的消息,彼此是竞争关系(一条消息只会被Group内的一个消费者处理);而不同的Group订阅同一个Topic,彼此独立,都能收到全量消息,从而实现一对多的广播效果

三、消息发送方式

发送方式 特点 可靠性 适用场景
同步发送 (Sync) 发送后阻塞,等待Broker返回结果。 最高。可通过SendResult判断消息是否成功。 重要的业务通知、订单状态更新等对可靠性要求极高的场景。
异步发送 (Async) 发送后不阻塞,通过回调接口处理结果。 。失败时可重试,且不阻塞主线程,吞吐量高。 高并发、对响应时间敏感的业务,如电商下单、秒杀系统。
单向发送 (Oneway) 只管发送,不等待任何响应。 。消息可能丢失,无任何重试机制。 日志收集、监控数据上报等允许少量丢失、追求极致性能的场景。

四、可靠性机制:确认与重试

  • 生产者确认:同步发送返回的 SendResult 包含 SendStatus,如 SEND_OK、FLUSH_DISK_TIMEOUT 等,告知发送结果。异步发送则通过 SendCallback 的 onSuccess 和 onException 方法来感知结果。

  • 消费者确认:消息处理回调必须返回 ConsumeConcurrentlyStatus。
    CONSUME_SUCCESS:消费成功,Broker会更新消费位点。
    RECONSUME_LATER:消费失败,Broker会稍后重试推送该消息。

  • 重试与死信
    RocketMQ会自动为每个消费者组创建一个重试队列。消费失败的消息会被先移入重试队列,避免阻塞正常消息。
    每条消息默认最多重试 16次。超过后,消息会被移入死信队列(DLQ),需要人工介入处理。

五、高级特性

消息消费模式

  • 集群模式(Clustering):默认模式。同一个消费者组内的多个消费者共同消费队列中的消息,每条消息只会被组内一个消费者实例消费。这是负载均衡的模式。

  • 广播模式(Broadcasting)每条消息都会被同一个消费者组内的所有消费者实例消费一次。即消息会“广播”给组内的每一个消费者。

顺序消息

要保证消息的局部顺序,需在生产者和消费者两端协作完成。

  • 生产者:使用 MessageQueueSelector 将同一业务键(如订单ID)的消息都发送到同一个队列中。

  • 消费者:使用 MessageListenerOrderly 进行消费,它会确保一个队列同时只有一个线程在处理。

事务消息

RocketMQ解决分布式事务问题的核心特性,它通过两阶段提交事务回查机制,保证了本地事务消息发送的原子性——要么两者都成功,要么都不成功。在保证消息可靠投递的同时,实现了最终一致性

典型场景如“下单支付扣库存”:

  1. 生产者发送半消息到Broker,消费者不可见

  2. 生产者执行本地事务,executeLocalTransaction() 插入订单(待支付)
    成功→UNKNOW,失败→ROLLBACK

  3. Broker回查本地事务,checkLocalTransaction() 查询订单支付状态
    已支付→COMMIT,待支付→UNKNOW,已取消→ROLLBACK

  4. 最终本地事务状态
    COMMIT → 消息对消费者可见
    多次UNKNOW后超时 → ROLLBACK 消息删除

通过两阶段提交+回查机制,将“支付结果”作为事务最终状态的判断依据。

  • 第一阶段:生产者向Broker发送一条半消息(Half Message)
    消息成功写入Broker并持久化
    消费者不可见(无法消费)
    消息状态为UN_KNOW

  • 第二阶段:生产者执行本地事务后,根据结果向Broker发送提交(Commit) 或回滚(Rollback) 指令
    提交 → 半消息转为普通消息,消费者可见
    回滚 → 半消息被删除,消费者永远看不到
     

消息过滤

RocketMQ在Broker端提供两种过滤方式,可有效减少无效网络传输:

  • Tag过滤:简单、高效,适合标签明确的场景。例如:consumer.subscribe("TOPIC", "TAG_A || TAG_B");

  • SQL92过滤:功能强大,支持自定义属性过滤。需要Broker开启enablePropertyFilter配置。例如:consumer.subscribe("TOPIC", MessageSelector.bySql("a > 5 and b = 'abc'"));

延迟消息

通过设定延迟级别,让消息在一段时间后投递,适用于订单超时未支付等场景。RocketMQ支持18个延迟级别,如 1s, 5s, 10s, ... 2h。设置方法:

message.setDelayTimeLevel(3); // 3级对应10秒延迟

RocketMQ5.0版本实现延迟消息的核心机制:时间轮 + TimerLog

  • TimerWheel:本质是一个环形数组,数组的每个槽位(slot)代表一个时间刻度(例如1秒)。指针按固定频率(如50ms/次)移动,指向当前时间对应的槽位。

  • TimerLog:一个顺序追加的日志文件,用于存储消息的索引(如消息在CommitLog中的物理偏移量)。

完整流程

  1. 发送与持久化
    生产者发送带投递时间的延迟消息到Broker。Broker收到后,正常持久化到CommitLog,Topic不变。

  2. 索引构建
    持久化成功后,Broker不修改原消息,而是提取元数据(物理偏移量、投递时间、原始Topic等),形成一条索引写入 TimerLog,并用时间轮(TimerWheel)管理该索引。
    此时消费者拉取Topic 时:Broker会检查TimerLog中是否存在该消息的未到期索引。由于索引存在且未到期,Broker主动过滤,不返回这条消息。因此消息虽然躺在CommitLog中,但消费者看不到。

  3. 精准投递
    时间轮指针按固定频率(如50ms)转动。当指针转到投递时间对应的刻度时,触发该索引。Broker根据索引中的物理偏移量,从CommitLog中找到原消息,清除索引(或标记为已到期)。此后消费者再次拉取Topic,Broker检查发现无未到期索引,正常返回消息,完成消费。

消费幂等控制

在分布式消息系统中,消费幂等是一个核心设计问题。由于网络抖动、消费者重启、Broker重试等原因,消息可能会被重复投递,因此消费者必须实现幂等处理,确保重复消息不会导致业务数据异常。消息消费幂等的核心:根据唯一标识判断消息是否已被处理过,避免重复执行业务逻辑。

  • Message ID:RocketMQ为每条消息生成的全局唯一标识符,主要用于消息的唯一性标识链路追踪

  • Key业务层面的自定义标识符,由生产者在发送消息时设置,用于业务关联和消息过滤。

不建议仅依赖Message ID做幂等,因为极端情况下可能重复(Broker重启、重试等)。优先使用业务Key作为幂等ID,Message ID作为兜底。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐