Skip to content

订阅关系

什么是订阅关系?

RocketMQ系统中订阅关系是消费者获取消息、处理消息的规则和配置。 订阅关系由消费者分组动态注册到服务端系统,并在后续的消息传输中按照订阅关系定义的过滤规则进行消息匹配和消费进度维护。 通过配置订阅关系可控制如下传输行为:

  • 消息过滤规则:用于控制消费者在消费消息时,选择主题(Topic)内的哪些消息进行消费,设置消费过滤规则可以高效的过滤消费者需要的消息集合,灵活根据不同的业务场景设置不同的消息接收范围。
  • 消费状态:RocketMQ服务端默认提供订阅关系持久化能力,即消费者分组(ConsumerGroup)在服务端注册订阅关系后,当消费者离线再次上线后,可以获取离线前的消费进度并继续消费。

零基础可以先这样理解:

订阅关系 = ConsumerGroup 对某个 Topic 的消费规则。

它回答三个问题:消费哪个 Topic、过滤哪些 Tag、从哪里继续消费。

工作原理图

mermaid
flowchart TD
    A["Producer 发送消息"] --> B["Topic: ORDER_EVENT"]
    B --> C{"Broker 根据订阅关系匹配"}
    C --> D["ConsumerGroup: SMS_CG\n订阅 ORDER_PAID"]
    C --> E["ConsumerGroup: POINTS_CG\n订阅 ORDER_PAID"]
    C --> F["ConsumerGroup: STOCK_CG\n订阅 ORDER_CANCELED"]
    D --> G["短信系统实例 1/2 负载均衡消费"]
    E --> H["积分系统实例 1/2 负载均衡消费"]
    F --> I["库存系统实例消费取消事件"]

同一个 Topic 的同一条消息,可以被多个 ConsumerGroup 各消费一份;但同一个 ConsumerGroup 内有多个 Consumer 实例时,它们是分摊消费,不是每个实例都收到一份。

订阅关系判断原则

RocketMQ的订阅关系按照消费者分组(ConsumerGroup)和主题(Topic)颗粒度设计,因此,一个订阅关系指的是某个消费者分组对于某个主题的订阅(重在ConsumerGroup和Topic之间),规则如下:

  • 不同消费者分组对于同一个主题的订阅相互独立如下图所示,消费者分组GroupA和GroupB分别以不同的订阅关系订阅了同一个主题Topic A,这两个订阅关系相互独立,可以各自定义,不受影响。 订阅相互独立
  • 同一个消费者分组对于不同主题的订阅也相互独立如下图所示,消费者分组Group A订阅了主题Topic A和主题Topic B,对于Group A来说订阅的Topic A为一个订阅关系,订阅的Topic B为另一个订阅关系,且这两个订阅关系相互独立,可以各自定义,不受影响。 不同主题订阅相互独立

订阅关系一致

同一消费者组(ConsumeGroup)内的消费者(Consumer)在消费逻辑上必须保持一致即订阅关系一致,参考订阅关系判断原则并结合图文来思考。

订阅关系示例

订阅关系一致

  1. 同一消费者组(ConsumeGroup)内的所有消费者订阅主题(Topic)过滤表达式一致(如Tag都相同)。如图: 三个相同

同一ConsumerGroup下的三个Consumer实例C1、C2和C3分别都订阅了TopicA,且订阅TopicA的Tag也都是Tag1,符合订阅关系一致原则。

订阅关系不一致

  1. 同一消费者组(ConsumeGroup)下的消费者(Consumer)实例订阅的主题(Topic)不同。

官网说以上问题适用于3.x、4.x的SDK,5.x版本SDK已经支持同一个ConsumerGroup下的Consumer实例订阅不同的Topic。

后半句同一个ConsumerGroup下的Consumer实例订阅不同的Topic,有些许疑惑。

疑惑已解决: 由于网上的文章大多数打着版本5的使用,但实际使用的仍然是旧的api(着实很坑),因此新版的订阅关系一致性无法实现,其实官网也给了版本5的简单示例SDK测试消息收发

快速连接版本5.xAPI使用

代码 Demo:同组订阅保持一致

同一个 ConsumerGroup 内的多个实例应该订阅相同 Topic 和 Tag。

java
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("SMS_ORDER_EVENT_CG");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("ORDER_EVENT", "ORDER_PAID || ORDER_CANCELED");
consumer.registerMessageListener((MessageListenerConcurrently) (messages, context) -> {
    for (MessageExt message : messages) {
        System.out.println("短信系统处理订单事件:" + message.getTags());
    }
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();

不要让同一个 SMS_ORDER_EVENT_CG 的 A 实例订阅 ORDER_PAID,B 实例订阅 ORDER_CANCELED。如果消费逻辑不同,应拆成不同 ConsumerGroup。

订阅关系设计建议

需求推荐做法原因
短信和积分都要处理支付事件建两个 ConsumerGroup两个系统各自保留消费进度
同一个短信服务部署 3 个实例3 个实例使用同一个 ConsumerGroup实例之间负载均衡,提高吞吐
支付和取消逻辑完全不同可以拆不同 ConsumerGroup 或不同消费服务避免同组订阅表达式不一致
临时排查消息用单独 ConsumerGroup不影响线上消费者进度