Skip to content

领域模型

领域模型 如图rocketMQ中的消息的生命周期主要分为消息生产、消息存储、消息消费三部分。

生产者生产消息并发送至服务端,消息被存储在服务端的主题(Topic)中,消费者通过订阅主题来消费消息。

零基础可以先这样理解:

RocketMQ 的核心链路就是:生产者把消息发到 Topic,Broker 按队列存储,消费者组按订阅关系消费。

为什么要先理解领域模型

RocketMQ 的很多问题不是 API 写错,而是模型理解错。例如:以为增加消费者实例就能让每个实例都收到全量消息;以为 Topic 就等于队列;以为同一个消费组里可以随意配置不同过滤条件。这些都会导致消息少消费、重复消费、订阅关系异常或排查方向错误。

领域模型要解决的是三个问题:

  1. 消息发到哪里:由 Topic、Tag、Key 描述。
  2. 消息存在哪里:由 Broker 和 MessageQueue 承担。
  3. 谁能消费到:由 ConsumerGroup 和订阅关系决定。

工作原理图

mermaid
flowchart TD
    A["Producer<br/>订单服务"] --> B["Topic<br/>ORDER_EVENT"]
    B --> C["MessageQueue 0"]
    B --> D["MessageQueue 1"]
    B --> E["MessageQueue 2"]
    C --> F["ConsumerGroup: SMS_CG"]
    D --> F
    E --> F
    C --> G["ConsumerGroup: POINTS_CG"]
    D --> G
    E --> G
    F --> H["短信消费者实例<br/>组内分摊消费"]
    G --> I["积分消费者实例<br/>独立消费进度"]

这张图里最重要的边界是:Topic 负责消息分类,MessageQueue 负责存储和并发,ConsumerGroup 负责消费身份和进度。

生产者

用于产生消息的运行实体,一般集成于业务调用链的上游,轻量级、匿名、无身份。

消息

主题(Topic)

消息传输和存储的分组容器,内部由多个队列组成,消息的存储和水平拓展实际是通过主题内队列实现的。

队列(MessageQueue)

消息传输和存储的实际单元容器,类比于其他消息队列中的分区,rocketMQ通过流式特性的无限队列结构来存储消息,消息在队列内具备顺序性存储特征

消息(Message)

rocketMQ的最小传输单元,消息具备不可变性,在初始化发送和完成存储后即不可变。

消费消息

消费者分组(ConsumerGroup)

rocketMQ发布订阅模型中定义的独立的消费身份分组,用于统一管理底层运行的多个消费者(Consumer)。同一个消费组的多个消费者必须保持消费逻辑和配置一致,共同分担消费组订阅的消息。实现消费水平能力的水平拓展。

消费者(Consumer)

消息消费的运行实体,一般集成于业务调用链的下游,消费者必须被指定到某一个消费组中。

订阅关系(Subscription)

发布订阅模型中消息过滤、重试、消费进度的配置规则。订阅关系以消费组粒度进行管理,消费组通过定义订阅关系控制指定消费组下的消费者如何实现消息过滤、消费重试及消费进度恢复等。rocketMQ的订阅关系除过滤表达式之外都是持久化的,即服务端重启或请求断开,订阅关系依然保留。

通信方式

分布式架构思想下,经常将复杂系统拆分为多个独立的子模块,此时就需要考虑子模块之间的远程通信。典型的的通信方式分为以下两种,一种是同步的RPC远程调用;一种是基于中间件的异步通信方式。

同步RPC调用模型

通信方式 同步RPC调用模型下,不同系统之间直接进行调用通信,每个请求直接从调用方发送到被调用方,然后要求被调用方立即返回响应结果给调用方,以确定本次调用结果是否成功。

此处的同步并不代表RPC的编程接口方式,RPC也可以有异步非阻塞调用的编程方式,但本质上仍然是需要在指定时间内得到目标端的直接响应。

异步通信模型

通信方式

异步消息通信模式下,各子系统之间无需强耦合直接连接,调用方只需要将请求转化成异步事件(消息)发送给中间代理,发送成功即可认为该异步链路调用完成,剩下的工作中间代理会负责将事件可靠通知到下游的调用系统,确保任务执行完成。该中间代理一般就是消息中间件。

异步通信的优势

  • 系统拓扑简单。由于调用方和被调用方统一和中间代理通信,系统是星型结构,易于维护和管理。
  • 上下游耦合性弱。上下游系统之间弱耦合,结构更灵活,由中间代理负责缓冲和异步恢复。 上下游系统间可以独立升级和变更,不会互相影响。
  • 容量削峰填谷。基于消息的中间代理往往具备很强的流量缓冲和整形能力,业务流量高峰到来时不会击垮下游。

消息传输模型

主流的消息中间件的传输模型主要为点对点模型和发布订阅模型。

点对点模型

点对点模型] 点对点模型也叫队列模型,具有如下特点:

  • 消费匿名:消息上下游沟通的唯一的身份就是队列,下游消费者从队列获取消息无法申明独立身份。
  • 一对一通信:基于消费匿名特点,下游消费者即使有多个,但都没有自己独立的身份,因此共享队列中的消息,每一条消息都只会被唯一一个消费者处理。因此点对点模型只能实现一对一通信。

发布订阅模型

发布订阅模型 发布订阅模型具有如下特点:

  • 消费独立:相比队列模型的匿名消费方式,发布订阅模型中消费方都会具备的身份,一般叫做订阅组(订阅关系),不同订阅组之间相互独立不会相互影响。
  • 一对多通信:基于独立身份的设计,同一个主题内的消息可以被多个订阅组处理,每个订阅组都可以拿到全量消息。因此发布订阅模型可以实现一对多通信。

传输模型对比

点对点模型和发布订阅模型各有优势,点对点模型更为简单,而发布订阅模型的扩展性更高。Apache RocketMQ使用的传输模型为发布订阅模型,因此也具有发布订阅模型的特点。

mermaid
flowchart TD
    A["一条订单支付消息"] --> B{"消费模型"}
    B -- "同一个 ConsumerGroup" --> C["多个实例负载均衡\n只处理一次"]
    B -- "不同 ConsumerGroup" --> D["每个组各处理一次\n短信、积分、物流都能收到"]

所以,不要用“有几个消费者实例”判断消息会被处理几次,要看它们是不是属于同一个 ConsumerGroup。

如果模型用错会怎样

错误理解直接后果正确理解
一个 Topic 只能给一个系统消费业务被迫复制多份消息不同 ConsumerGroup 可以独立消费同一 Topic
同一个消费组的实例都会收到全量消息实际只有一个实例处理某条消息,导致以为“丢消息”组内负载均衡,组间广播语义
Tag 可以随意改订阅关系变化后部分消息收不到Tag 是过滤语义,要稳定设计
Key 可有可无排查消息时无法按业务主键检索Key 应该放订单号、支付单号等业务唯一标识
Topic 混放所有业务权限、监控、重试和排查都混乱按业务域和消息语义拆 Topic

代码 Demo:Topic、Tag、Key 的使用

下面示例展示 RocketMQ 最常见的消息字段。初学时先记住:Topic 做大类,Tag 做小类,Key 放业务唯一标识。

java
Message message = new Message(
    "ORDER_EVENT",
    "ORDER_PAID",
    "pay-order-1001",
    "{\"orderId\":1001,\"paid\":true}".getBytes(StandardCharsets.UTF_8)
);

消费者按 Topic 和 Tag 订阅:

java
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("POINTS_ORDER_PAID_CG");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("ORDER_EVENT", "ORDER_PAID");
consumer.registerMessageListener((MessageListenerConcurrently) (messages, context) -> {
    for (MessageExt msg : messages) {
        String body = new String(msg.getBody(), StandardCharsets.UTF_8);
        System.out.println("收到消息:" + msg.getKeys() + ",内容:" + body);
    }
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();

不同消费组可以订阅同一个 Topic,例如积分系统和短信系统都消费 ORDER_PAID,但它们的消费进度互不影响。