Rocketmq 本质上是一个分布式消息队列,用来让不同系统之间通过“消息”异步协作。
一、为什么需要消息队列
假设“患者建档”完成后,需要:
如果建档接口直接依次调用这些系统,任何一个系统变慢或故障,都会影响主流程。
使用 RocketMQ 后:
|
1 2 3 4 5 |
建档服务 -> 发送“患者已建档”消息 -> RocketMQ | +-------------+-------------+ | | | 日志服务 短信服务 统计服务 |
它主要解决:
二、核心组件
|
1 2 |
Producer -> NameServer -> Broker -> Consumer 生产者 路由中心 消息服务器 消费者 |
生产消息的一方,例如:
|
1 2 3 4 |
rocketMQTemplate.convertAndSend( "chronic-log-topic", chronicLog ); |
Producer 将消息发送给 Broker。
消费消息的一方,例如你的:
|
1 2 3 4 5 6 7 8 9 10 11 |
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("chronic-log-consumer-group"); consumer.subscribe("chronic-log-topic", "*"); consumer.registerMessageListener( (MessageListenerConcurrently) (messages, context) -> { // 处理消息 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } ); |
负责:
生产环境通常部署多个 Broker,并可以配置主从复制。
NameServer 保存 Topic 和 Broker 的路由关系。
Producer 和 Consumer 先查询 NameServer,知道消息应该发到或从哪个 Broker 获取。NameServer 自身不存储业务消息。
三、Topic、Tag 和 Key
Topic 是消息的业务分类。
|
1 2 3 |
chronic-log-topic patient-event-topic payment-topic |
生产者和消费者通过 Topic 联系起来。
Tag 是 Topic 下的二级分类。
例如患者事件统一放在一个 Topic:
|
1 2 3 4 |
Topic: patient-event-topic Tag: CREATED Tag: UPDATED Tag: DELETED |
消费者可以只订阅其中一类:
|
1 |
consumer.subscribe("patient-event-topic", "CREATED"); |
订阅多个 Tag:
|
1 2 3 4 |
consumer.subscribe( "patient-event-topic", "CREATED || UPDATED" ); |
订阅全部:
|
1 |
consumer.subscribe("patient-event-topic", "*"); |
Key 用于标识和查询消息,常放业务唯一编号:
|
1 |
message.setKeys(patientId); |
Topic 负责分类,Tag 负责过滤,Key 主要用于定位和排查消息。
四、Consumer Group
Consumer Group 表示一组承担相同消费职责的消费者。
|
1 2 3 4 5 6 |
Topic: chronic-log-topic Group: chronic-log-group 实例 A ─┐ 实例 B ─┼─ 共同分摊消息 实例 C ─┘ |
在集群模式下,同一条消息通常只会交给其中一个实例处理。
这适合服务部署多个副本:
|
1 2 3 |
chronic-service-1 chronic-service-2 chronic-service-3 |
三个实例使用同一个 Group,可以共同消费、提高吞吐量。
|
1 2 3 4 5 6 7 |
chronic-log-topic | +-> Group: log-storage-group | +-> Group: audit-group | +-> Group: statistics-group |
每个 Group 都会收到一份消息。
因此,两个业务都需要处理同一条消息时,必须使用不同 Group。
同一个 Group 的消费者应保持以下配置一致:
不要为了节省 Group 名称,让多个不相关业务共用一个 Group。
五、MessageQueue
一个 Topic 内部通常包含多个 MessageQueue:
|
1 2 3 4 5 6 |
Topic: chronic-log-topic Queue 0 Queue 1 Queue 2 Queue 3 |
MessageQueue 是 RocketMQ 分配消息和并行消费的基本单位。
假设一个 Topic 有 4 个队列:
| 消费实例数 | 效果 |
|---|---|
| 1 | 一个实例负责4个队列 |
| 2 | 每个实例大约负责2个队列 |
| 4 | 每个实例大约负责1个队列 |
| 6 | 可能有2个实例分不到队列 |
所以增加 Consumer 实例不一定能无限提高消费能力,还需要足够的队列数量。
六、消费模式
CLUSTERING|
1 |
consumer.setMessageModel(MessageModel.CLUSTERING); |
同一个 Group 中,一条消息只由一个实例处理。这是最常用的模式。
BROADCASTING|
1 |
consumer.setMessageModel(MessageModel.BROADCASTING); |
同一个 Group 中,每个实例都会处理消息。
广播模式通常不适合需要可靠重试、统一消费进度的关键业务。业务系统需要多方处理时,一般优先使用不同 Group,而不是广播模式。
七、消费成功、重试和死信
并发消费监听器会返回状态:
|
1 |
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; |
表示处理成功,消费进度可以向前推进。
处理失败时返回:
|
1 |
return ConsumeConcurrentlyStatus.RECONSUME_LATER; |
RocketMQ 会在稍后重新投递。
超过最大重试次数后,消息通常会进入死信队列:
|
1 |
%DLQ%consumer-group-name |
需要监控并人工或自动补偿死信消息。
八、RocketMQ 的消息可靠性
RocketMQ 通常提供“至少一次”投递语义:
|
1 |
消息可能重复,但尽量不丢 |
例如:
因此消费者必须考虑幂等性。
常见方法:
|
1 2 3 4 |
消息唯一 ID + 数据库唯一索引 业务单号 + 唯一约束 消费记录表 Redis SETNX |
不要假设一条消息绝对只会被消费一次。
九、常见消息类型
最常用,发送后尽快消费。
消息延迟一段时间再投递,例如:
不同 RocketMQ 版本对延迟级别和任意延迟时间的支持有所区别。
要求消息按顺序处理,例如:
|
1 |
患者建档 -> 修改档案 -> 注销档案 |
通常需要把相同业务 ID 的消息发送到同一个 MessageQueue,并使用顺序消费监听器。
RocketMQ 主要保证队列内顺序,不保证整个 Topic 的全局顺序。
用于协调本地数据库事务和消息发送,例如:
|
1 |
创建订单成功,并且必须可靠地发出“订单已创建”消息 |
它通过半消息、本地事务执行和事务状态回查降低数据库与消息状态不一致的风险,但不能直接替代完整的业务幂等和补偿机制。
十、Push Consumer 并不是真正由 Broker 主动推送
DefaultMQPushConsumer 名称中虽然有 Push,但客户端内部本质上仍然是主动向 Broker 拉取消息。
RocketMQ 客户端把拉取、线程池、队列分配和消费进度管理封装起来,对业务代码表现得像消息被推送到监听器。
十一、结合你的 ChronicLogConsumer
当前代码的运行过程是:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 |
Spring 创建 ChronicLogConsumer | @PostConstruct 创建 DefaultMQPushConsumer | 设置 NameServer、Namespace、消费组 | 订阅 log.rocketmq.topic | 注册一个 MessageListenerConcurrently | consumer.start() | RocketMQ 分配 MessageQueue | 收到消息后保存 MySQL 或 MongoDB |
这里需要重点检查:
rocketmq.topic.group 是否被其他业务监听器共用shutdown()@RefreshScope 刷新配置后是否可能重新创建 Consumer
最核心的记忆方式是:
|
1 2 3 4 5 6 |
Topic 决定消息属于哪类业务 Tag 决定消费者筛选哪些消息 Group 决定哪些消费者共同承担一份消费职责 MessageQueue 决定并行度和分配单位 返回状态决定成功还是重试 幂等性负责解决重复消费 |