RocketMQ 订阅不一致:一次消费者组设计问题排查记录

RocketMQ 订阅不一致:别把注解监听器当成一个 Consumer

简单记录一次 RocketMQ 订阅关系治理。真正的坑不在 Topic 数量,而在于误把 Spring 注解监听器理解成了 Java Client 的“一个 Consumer 多次订阅”。

问题现象

同一个应用里有多个监听器,分别消费不同 Topic,却复用了服务级 Consumer Group:

@RocketMQMessageListener(
        topic = "topic_a",
        consumerGroup = "consumer_group_service"
)
class TopicAListener {}

@RocketMQMessageListener(
        topic = "topic_b",
        consumerGroup = "consumer_group_service"
)
class TopicBListener {}

单实例、低流量时不一定马上暴露。扩容或滚动发布后,RocketMQ 可能提示订阅关系不一致,表现为消费者反复注册、部分 Topic 消费异常或消费状态不稳定。

最容易忽略的实现差异

Java Client 的常见写法是:手动创建一个 DefaultMQPushConsumer,然后在同一个 Consumer 上调用多次 subscribe。因此,一个进程里通常只有一个 Consumer 实例,对外上报的是一份完整订阅关系:

DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer_group_service");
consumer.subscribe("topic_a", "*");
consumer.subscribe("topic_b", "*");
consumer.start();

@RocketMQMessageListener 不是“给一个 Consumer 增加一个订阅”。Spring RocketMQ Starter 会扫描每个注解,为每个监听器创建独立的监听容器,并由容器创建对应的 RocketMQ Consumer。也就是说:

2 个 @RocketMQMessageListener
        ↓
2 个监听容器
        ↓
2 个 Consumer 实例

如果这两个注解的 consumerGroup 相同,那么同一个应用进程内已经存在两个使用同组名称、却上报不同 Topic 的 Consumer。多实例部署后,这种差异会被放大,最终触发订阅不一致。

根因分析

Consumer Group 代表一组稳定的消费语义,不是简单的服务名称。同一组消费者需要保持一致的 Topic、Tag 或订阅表达式、消费模式和业务语义。

原有设计把“同一个服务”误当成了“同一个消费者组”,同时忽略了注解模式下“一注解一容器一 Consumer”的生命周期模型。因此,问题本质是消费者组边界设计错误,而不是 RocketMQ 集群本身故障。

修复方案

将消费者组从“一个服务一个组”调整为“一个服务 + 一个 Topic 一个组”:

@RocketMQMessageListener(
        topic = "topic_a",
        consumerGroup = "consumer_group_service_topic_a"
)
class TopicAListener {}

@RocketMQMessageListener(
        topic = "topic_b",
        consumerGroup = "consumer_group_service_topic_b"
)
class TopicBListener {}

统一命名格式:

consumer_group_{service}_{topic}

本次先处理导出类、分配类、活动订单类、批量任务类等监听器,随后将同样规则推广到其他 Topic。消费者组常量集中管理,避免 Listener 中出现散落且容易重复的字符串。

发布与验证

改造按以下顺序执行:

  1. 扫描所有 @RocketMQMessageListener,按注解统计 Consumer Group、Topic、Tag 和消费模式,检查同组是否存在多份订阅。
  2. 为每个 Topic 配置独立消费者组,确认同一组的所有实例订阅参数完全一致。
  3. 滚动发布,观察监听容器启动、消费者注册、消息堆积和消费速率。
  4. 在 RocketMQ 控制台核对每个组对应的 Topic、Tag 和实例列表。
  5. 对新组的消费起点、历史消息处理和重复消费风险,按实际 RocketMQ 配置与业务用例验证。

验证重点不是“服务是否启动”,而是“每个注解创建的 Consumer 是否只承担清晰、稳定的一种消费语义”。

总结

这次问题可以归纳为:

Java Client:一个 Consumer,多次 subscribe
注解监听器:一个注解,一个容器,一个 Consumer

因此,在 Spring 注解模式下,不同 Topic 默认拆分消费者组;即使 Topic 相同,Tag、消费模式或业务语义不一致时,也应重新评估是否需要拆组。Consumer Group 的命名应围绕“实际 Consumer 的订阅边界”设计,而不是简单复用服务名称。