RocketMQ 全面使用文档
1. RocketMQ 概述
1.1 什么是 RocketMQ
RocketMQ 是阿里巴巴开源的一款分布式消息中间件,基于 Java 语言开发,具有高并发、高可靠、高可用、低延迟等特性。它最初诞生于阿里巴巴内部的 MetaQ,2012 年正式开源,2017 年成为 Apache 顶级项目。
RocketMQ 主要用于解决分布式系统中的异步解耦、流量削峰、日志收集、数据同步等问题,广泛应用于电商、金融、物流等对消息可靠性要求极高的场景。
1.2 为什么选择 RocketMQ
与其他消息中间件(如 Kafka、RabbitMQ)相比,RocketMQ 具有以下优势:
- 高可靠性:支持同步刷盘、同步复制,确保消息不丢失;
- 高可用性:NameServer 无状态集群,Broker 主从架构,故障自动转移;
- 高性能:基于零拷贝、顺序写盘等技术,单机支持百万级 TPS;
- 丰富的消息类型:支持普通消息、顺序消息、事务消息、延迟消息等;
- 灵活的消息过滤:支持 Tag 过滤、SQL92 过滤,满足复杂业务需求;
- 完善的运维工具:提供 RocketMQ Console 可视化控制台,方便监控和管理。
1.3 应用场景
- 异步解耦:将核心业务与非核心业务分离,提升系统响应速度(如订单创建后异步发送短信、邮件);
- 流量削峰:在高并发场景下缓冲流量,保护下游系统(如秒杀活动);
- 日志收集:收集系统日志,统一存储和分析(如 ELK 架构中的日志传输);
- 数据同步:实现不同系统或数据库之间的数据同步;
- 事件驱动:构建事件驱动架构,实现业务流程的解耦和编排。
2. 核心概念详解
2.1 NameServer
NameServer 是 RocketMQ 的服务发现组件,类似于 Kafka 的 Zookeeper,但更轻量级。它的主要作用是:
- 管理 Broker 的路由信息:Broker 启动时会向所有 NameServer 注册自己的信息(包括 Broker 地址、Topic 配置等);
- 为 Producer 和 Consumer 提供路由查询:Producer 和 Consumer 通过 NameServer 获取 Broker 的地址,从而与 Broker 通信。
NameServer 是无状态的,多个 NameServer 之间互不通信,每个 NameServer 都保存完整的路由信息。Producer 和 Consumer 可以配置多个 NameServer 地址,只要有一个 NameServer 可用,就能正常工作。
2.2 Broker
Broker 是 RocketMQ 的核心存储和转发组件,负责接收、存储和投递消息。它的主要作用是:
- 接收 Producer 发送的消息,并持久化到磁盘;
- 维护 ConsumeQueue 和 IndexFile,为 Consumer 提供消息查询和投递服务;
- 与 NameServer 保持心跳,定期上报自己的状态;
- 支持主从复制,实现高可用。
Broker 分为主节点(Master)和从节点(Slave):
- Master:负责接收 Producer 的消息写入,同时也可以处理 Consumer 的消息读取;
- Slave:从 Master 同步消息,主要负责处理 Consumer 的消息读取,分担 Master 的读压力。
2.3 Producer
Producer 是消息生产者,负责将业务消息发送到 Broker。RocketMQ 提供了多种发送方式:
- 同步发送:发送后等待 Broker 返回结果,可靠性高但会阻塞线程;
- 异步发送:发送后不等待结果,通过回调函数处理发送结果,性能高;
- Oneway 发送:发送后不等待结果,也不处理回调,性能最高但可靠性最低。
Producer 属于一个生产者组(Producer Group),同一个生产者组的 Producer 可以发送同一个 Topic 的消息,Broker 会根据生产者组进行事务消息的回查。
2.4 Consumer
Consumer 是消息消费者,负责从 Broker 拉取消息并进行业务处理。RocketMQ 提供了两种消费模式:
- Push 模式:Broker 主动将消息推送给 Consumer,实时性高;
- Pull 模式:Consumer 主动从 Broker 拉取消息,灵活性高。
Consumer 也属于一个消费者组(Consumer Group),同一个消费者组的 Consumer 共同消费 Topic 的消息,实现负载均衡和容错。
2.5 Topic
Topic 是消息的主题,用于区分不同类型的消息。Producer 将消息发送到指定的 Topic,Consumer 从指定的 Topic 订阅消息。
Topic 是逻辑概念,一个 Topic 可以包含多个Message Queue(消息队列),Message Queue 是物理概念,用于实现消息的分区存储和并行消费。
2.6 Tag
Tag 是消息的子主题,用于进一步区分 Topic 下的消息。Producer 可以为消息设置 Tag,Consumer 可以根据 Tag 进行消息过滤,只消费感兴趣的消息。
例如,一个 Topic 是“order_topic”,Tag 可以是“create_order”、“pay_order”、“cancel_order”,Consumer 可以只订阅“pay_order”的消息。
2.7 Message Queue
Message Queue(简称 Queue)是消息的物理存储单元,一个 Topic 可以有多个 Queue,分布在不同的 Broker 上。
Queue 的作用是:
- 分区存储:将消息分散存储在多个 Queue 中,提升存储和读写性能;
- 并行消费:同一个消费者组的 Consumer 可以并行消费不同的 Queue,提升消费能力;
- 顺序保证:同一个 Queue 中的消息是有序的,通过将同一组消息发送到同一个 Queue,可以实现分区顺序消息。
2.8 Offset
Offset 是消费进度,用于记录 Consumer 已经消费到 Queue 的哪个位置。RocketMQ 会将 Offset 持久化到 Broker(或本地文件),确保 Consumer 重启后能从上次的位置继续消费。
Offset 分为两种:
- Current Offset:Consumer 当前已经消费到的位置;
- Broker Offset:Broker 中 Queue 的最大位置,即最新消息的位置。
2.9 Message Model
Message Model 是消息消费模式,RocketMQ 支持两种消费模式:
- 集群消费(Clustering):同一个消费者组的 Consumer 共同消费 Topic 的消息,每个消息只被一个 Consumer 消费,实现负载均衡;
- 广播消费(Broadcasting):同一个消费者组的每个 Consumer 都会消费 Topic 的所有消息,即每个消息会被消费多次。
3. 架构设计与组件交互
3.1 整体架构
RocketMQ 的整体架构由以下四个部分组成:
- NameServer 集群:无状态,负责服务发现和路由管理;
- Broker 集群:由 Master 和 Slave 组成,负责消息存储和转发;
- Producer 集群:消息生产者,将消息发送到 Broker;
- Consumer 集群:消息消费者,从 Broker 拉取消息并消费。
架构图如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
┌─────────────────┐
│ NameServer 集群 │
└────────┬────────┘
│ 注册/心跳
│
┌────────▼────────┐ ┌─────────────────┐
│ Broker Master │────▶│ Broker Slave │
└────────┬────────┘ └─────────────────┘
│ 发送/拉取
│
┌────────▼────────┐ ┌─────────────────┐
│ Producer 集群 │ │ Consumer 集群 │
└─────────────────┘ └─────────────────┘
3.2 组件交互流程
3.2.1 Broker 启动与注册
- Broker 启动时,读取配置文件,获取 NameServer 地址列表;
- Broker 向所有 NameServer 发送注册请求,上报自己的信息(包括 Broker 地址、Topic 配置、Queue 信息等);
- Broker 启动定时任务,每隔 30 秒向所有 NameServer 发送心跳,更新自己的状态;
- NameServer 收到心跳后,更新 Broker 的存活时间,如果超过 120 秒没有收到心跳,NameServer 会将该 Broker 从路由信息中移除。
3.2.2 Producer 发送消息
- Producer 启动时,读取配置文件,获取 NameServer 地址列表;
- Producer 从 NameServer 拉取 Topic 的路由信息(包括 Broker 地址、Queue 列表等);
- Producer 启动定时任务,每隔 30 秒从 NameServer 更新一次路由信息;
- Producer 发送消息时,根据负载均衡策略选择一个 Queue,然后将消息发送到该 Queue 所在的 Broker;
- Broker 收到消息后,将消息持久化到 CommitLog,然后更新 ConsumeQueue 和 IndexFile;
- Broker 返回发送结果给 Producer。
3.2.3 Consumer 消费消息
- Consumer 启动时,读取配置文件,获取 NameServer 地址列表;
- Consumer 从 NameServer 拉取 Topic 的路由信息;
- Consumer 启动定时任务,每隔 30 秒从 NameServer 更新一次路由信息;
- Consumer 进行负载均衡,将 Topic 的 Queue 分配给同一个消费者组的各个 Consumer;
- Consumer 启动拉取线程,从分配给自己的 Queue 中拉取消息;
- Broker 收到拉取请求后,从 ConsumeQueue 中读取消息的位置,然后从 CommitLog 中读取消息,返回给 Consumer;
- Consumer 收到消息后,调用业务逻辑处理消息;
- Consumer 处理完消息后,向 Broker 提交 Offset,更新消费进度。
4. 核心特性与原理解析
4.1 高可用性
4.1.1 NameServer 高可用
NameServer 是无状态的,多个 NameServer 之间互不通信,每个 NameServer 都保存完整的路由信息。Producer 和 Consumer 可以配置多个 NameServer 地址,只要有一个 NameServer 可用,就能正常工作。
当某个 NameServer 宕机时,Producer 和 Consumer 会自动切换到其他可用的 NameServer,不会影响正常业务。
4.1.2 Broker 高可用
Broker 采用主从架构,一个 Master 可以配置多个 Slave。Master 负责消息写入,Slave 负责从 Master 同步消息,并分担读压力。
Broker 支持两种复制模式:
- 同步复制(SYNC_MASTER):Master 收到消息后,需要等待所有 Slave 同步完成,才返回成功给 Producer。这种模式可靠性高,但性能低;
- 异步复制(ASYNC_MASTER):Master 收到消息后,立即返回成功给 Producer,然后异步将消息同步给 Slave。这种模式性能高,但可靠性低,如果 Master 宕机,可能会丢失部分消息。
当 Master 宕机时,Slave 可以自动切换为 Master(需要手动配置或使用 RocketMQ 5.x 的 Dledger 模式),实现故障转移。
4.2 高性能
4.2.1 零拷贝技术
RocketMQ 使用零拷贝(Zero-Copy)技术提升消息读写性能。零拷贝是指在数据传输过程中,减少或避免 CPU 拷贝数据的次数,从而提升性能。
传统 IO 的流程:
- 应用程序调用
read()方法,从磁盘读取数据到内核缓冲区(CPU 拷贝:磁盘→内核缓冲区); - 内核将数据从内核缓冲区拷贝到用户缓冲区(CPU 拷贝:内核缓冲区→用户缓冲区);
- 应用程序调用
write()方法,将数据从用户缓冲区拷贝到 Socket 缓冲区(CPU 拷贝:用户缓冲区→Socket 缓冲区); - 内核将数据从 Socket 缓冲区拷贝到网卡(DMA 拷贝:Socket 缓冲区→网卡)。
传统 IO 涉及 4 次上下文切换和 3 次 CPU 拷贝。
RocketMQ 使用 mmap(内存映射) + sendfile 实现零拷贝:
- mmap:将文件直接映射到用户态内存,应用程序可以直接操作内存中的文件数据,减少一次 CPU 拷贝(内核缓冲区→用户缓冲区);
- sendfile:直接在内核态将数据从文件传输到 Socket 缓冲区,减少一次 CPU 拷贝(用户缓冲区→Socket 缓冲区)。
使用零拷贝后,流程变为:
- 应用程序调用
mmap(),将文件映射到用户态内存; - 应用程序调用
sendfile(),内核将数据从磁盘读取到内核缓冲区(DMA 拷贝),然后直接从内核缓冲区传输到 Socket 缓冲区(DMA 拷贝); - 内核将数据从 Socket 缓冲区传输到网卡(DMA 拷贝)。
零拷贝涉及 2 次上下文切换和 0 次 CPU 拷贝,大大提升了性能。
4.2.2 顺序写盘
磁盘的顺序写性能远高于随机写性能(顺序写可以达到几百 MB/s,随机写只有几 MB/s)。RocketMQ 采用顺序写盘技术,将所有消息写入同一个 CommitLog 文件,实现顺序写盘,提升写入性能。
4.2.3 内存映射文件
RocketMQ 使用内存映射文件(MappedByteBuffer)来操作 CommitLog、ConsumeQueue 和 IndexFile。内存映射文件将文件直接映射到内存,应用程序可以像操作内存一样操作文件,减少 IO 次数,提升性能。
4.2.4 刷盘策略
RocketMQ 支持两种刷盘策略:
- 同步刷盘(SYNC_FLUSH):消息写入内存后,立即调用
fsync()刷到磁盘,才返回成功给 Producer。这种模式可靠性高,但性能低; - 异步刷盘(ASYNC_FLUSH):消息写入内存后,立即返回成功给 Producer,然后由后台线程定时(默认 500ms)调用
fsync()刷到磁盘。这种模式性能高,但可靠性低,如果 Broker 宕机,可能会丢失部分消息。
4.3 消息存储原理
RocketMQ 的消息存储由三个部分组成:
- CommitLog:消息的存储文件,所有消息都顺序写入 CommitLog;
- ConsumeQueue:消费队列,存储消息在 CommitLog 中的偏移量,用于 Consumer 拉取消息;
- IndexFile:索引文件,存储消息的 Key 与 CommitLog 偏移量的映射,用于根据 Key 查询消息。
4.3.1 CommitLog
CommitLog 是消息的物理存储文件,默认存储在 $HOME/store/commitlog 目录下。每个 CommitLog 文件的大小默认是 1G,当一个文件写满后,会创建一个新的文件。
CommitLog 的存储格式如下:
| 字段 | 长度(字节) | 说明 |
|---|---|---|
| msgId | 8 | 消息 ID |
| storeSize | 4 | 消息总长度 |
| bodyCRC | 4 | 消息体 CRC 校验码 |
| queueId | 4 | 队列 ID |
| flag | 4 | 消息标志 |
| queueOffset | 8 | 队列偏移量 |
| physicOffset | 8 | CommitLog 物理偏移量 |
| sysFlag | 4 | 系统标志 |
| bornTimestamp | 8 | 消息生成时间戳 |
| bornHost | 8 | 消息生成主机 |
| storeTimestamp | 8 | 消息存储时间戳 |
| storeHost | 8 | 消息存储主机 |
| reconsumeTimes | 4 | 重试次数 |
| preparedTransactionOffset | 8 | 事务消息偏移量 |
| bodyLen | 4 | 消息体长度 |
| body | bodyLen | 消息体 |
| topicLen | 1 | Topic 长度 |
| topic | topicLen | Topic |
| propertiesLen | 2 | 属性长度 |
| properties | propertiesLen | 属性 |
4.3.2 ConsumeQueue
ConsumeQueue 是消费队列,默认存储在 $HOME/store/consumequeue/{topic}/{queueId} 目录下。每个 ConsumeQueue 文件的大小默认是 600W 字节,每个条目是 20 字节,所以每个文件可以存储 30W 条条目。
ConsumeQueue 的存储格式如下:
| 字段 | 长度(字节) | 说明 |
|---|---|---|
| commitLogOffset | 8 | 消息在 CommitLog 中的物理偏移量 |
| size | 4 | 消息长度 |
| tagHashCode | 8 | Tag 的哈希值 |
当 Broker 收到消息并写入 CommitLog 后,会根据消息的 Topic 和 QueueId,将消息的 CommitLog 偏移量、长度和 Tag 哈希值写入对应的 ConsumeQueue。
Consumer 拉取消息时,首先从 ConsumeQueue 中读取消息的位置,然后从 CommitLog 中读取消息,这样可以避免随机读 CommitLog,提升读性能。
4.3.3 IndexFile
IndexFile 是索引文件,默认存储在 $HOME/store/index 目录下。每个 IndexFile 的大小默认是 400M,包含 500W 个哈希槽和 2000W 个索引条目。
IndexFile 的存储格式如下:
| 字段 | 长度(字节) | 说明 |
|---|---|---|
| beginTimestamp | 8 | 索引文件中消息的最小存储时间戳 |
| endTimestamp | 8 | 索引文件中消息的最大存储时间戳 |
| beginPhyOffset | 8 | 索引文件中消息的最小 CommitLog 偏移量 |
| endPhyOffset | 8 | 索引文件中消息的最大 CommitLog 偏移量 |
| hashSlotCount | 4 | 哈希槽数量 |
| indexCount | 4 | 索引条目数量 |
| hashSlots | 4 * hashSlotCount | 哈希槽数组 |
| indexes | 20 * indexCount | 索引条目数组 |
索引条目的存储格式如下:
| 字段 | 长度(字节) | 说明 |
|---|---|---|
| keyHashCode | 4 | Key 的哈希值 |
| phyOffset | 8 | 消息在 CommitLog 中的物理偏移量 |
| timeDiff | 4 | 消息存储时间戳与 beginTimestamp 的差值 |
| preIndexNo | 4 | 前一个索引条目的编号 |
当 Broker 收到消息并写入 CommitLog 后,如果消息设置了 Key,会将 Key 的哈希值、CommitLog 偏移量等信息写入 IndexFile。
用户可以根据 Key 查询消息,Broker 会根据 Key 的哈希值找到对应的哈希槽,然后遍历哈希槽中的索引条目,找到对应的消息。
4.4 消息可靠性
4.4.1 消息发送可靠性
RocketMQ 通过以下机制保证消息发送的可靠性:
- 重试机制:同步发送和异步发送失败时,会自动重试(默认重试 2 次);
- 多种发送方式:根据业务需求选择同步、异步或 Oneway 发送;
- 发送状态确认:Broker 收到消息后,会返回发送状态(SEND_OK、FLUSH_DISK_TIMEOUT、FLUSH_SLAVE_TIMEOUT、SLAVE_NOT_AVAILABLE),Producer 可以根据发送状态判断是否需要重试。
4.4.2 消息存储可靠性
RocketMQ 通过以下机制保证消息存储的可靠性:
- 同步刷盘:消息写入内存后立即刷到磁盘,确保消息不丢失;
- 同步复制:Master 收到消息后,等待 Slave 同步完成,才返回成功,确保消息在多个节点上都有备份;
- 文件校验:CommitLog、ConsumeQueue 和 IndexFile 都有 CRC 校验码,确保文件内容不被篡改。
4.4.3 消息消费可靠性
RocketMQ 通过以下机制保证消息消费的可靠性:
- ACK 机制:Consumer 处理完消息后,需要向 Broker 提交 Offset,确认消息消费成功;
- 重试队列:如果 Consumer 消费失败,消息会被发送到重试队列,等待重新消费;
- 死信队列:如果消息重试超过一定次数(默认 16 次),会被发送到死信队列,等待人工处理。
4.5 顺序消息原理
顺序消息是指消息的消费顺序与发送顺序一致。RocketMQ 支持两种顺序消息:
- 分区顺序消息:同一个 Queue 中的消息是有序的,不同 Queue 之间的消息不保证顺序;
- 全局顺序消息:所有消息都在一个 Queue 中,保证所有消息的顺序,但性能较低。
4.5.1 分区顺序消息实现
分区顺序消息的实现需要 Producer 和 Consumer 配合:
- Producer 端:将同一组顺序消息发送到同一个 Queue 中。可以通过自定义
MessageQueueSelector实现,根据业务 ID(如订单 ID)选择同一个 Queue; - Consumer 端:使用
MessageListenerOrderly监听器,保证同一个 Queue 的消息由同一个线程消费,并且在消费失败时会暂停当前 Queue 的消费,直到重试成功。
4.5.2 全局顺序消息实现
全局顺序消息的实现只需要将 Topic 的 Queue 数量设置为 1,这样所有消息都会发送到同一个 Queue 中,然后使用分区顺序消息的方式消费即可。但全局顺序消息的性能较低,不适合高并发场景。
4.6 事务消息原理
事务消息是指消息的发送与本地事务的执行保持一致,要么都成功,要么都失败。RocketMQ 采用两阶段提交实现事务消息:
4.6.1 两阶段提交流程
- 第一阶段:发送 Half 消息
- Producer 向 Broker 发送 Half 消息(半消息),Half 消息对 Consumer 不可见;
- Broker 将 Half 消息存储在
RMQ_SYS_TRANS_HALF_TOPIC中; - Broker 返回成功给 Producer。
- 第二阶段:执行本地事务并提交/回滚
- Producer 执行本地事务(如操作数据库);
- 如果本地事务成功,Producer 向 Broker 发送 Commit 消息,Broker 将 Half 消息从
RMQ_SYS_TRANS_HALF_TOPIC移动到目标 Topic,供 Consumer 消费; - 如果本地事务失败,Producer 向 Broker 发送 Rollback 消息,Broker 将 Half 消息删除。
4.6.2 事务状态回查
如果 Producer 发送 Commit/Rollback 消息失败,或者 Broker 长时间没有收到 Commit/Rollback 消息,Broker 会发起事务状态回查:
- Broker 定时(默认 60s)扫描
RMQ_SYS_TRANS_HALF_TOPIC中的 Half 消息; - Broker 向 Producer 发送事务状态回查请求;
- Producer 实现
TransactionListener接口,在checkLocalTransaction方法中查询本地事务状态,返回 Commit、Rollback 或 Unknown; - 如果返回 Commit,Broker 将消息移动到目标 Topic;
- 如果返回 Rollback,Broker 将消息删除;
- 如果返回 Unknown,Broker 会继续回查,直到超过最大回查次数(默认 15 次),然后 Rollback。
4.7 延迟消息原理
延迟消息是指消息发送后,等待一定时间后才会被 Consumer 消费。RocketMQ 支持 18 个延迟级别:
1
1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
4.7.1 延迟消息实现流程
- Producer 发送延迟消息:Producer 设置消息的延迟级别,然后发送到 Broker;
- Broker 存储延迟消息:Broker 收到延迟消息后,将消息的 Topic 替换为
SCHEDULE_TOPIC_XXXX,QueueId 根据延迟级别选择(延迟级别 1 对应 QueueId 0,延迟级别 2 对应 QueueId 1,以此类推),然后将消息存储在 CommitLog 中,并更新对应的 ConsumeQueue; - 定时任务扫描延迟消息:Broker 启动
ScheduleMessageService定时任务,每隔 1s 扫描SCHEDULE_TOPIC_XXXX的 ConsumeQueue; - 投递延迟消息:当消息到达投递时间时,
ScheduleMessageService将消息的 Topic 恢复为原来的 Topic,然后投递到目标 Topic 的 ConsumeQueue 中,供 Consumer 消费。
4.8 消息过滤原理
RocketMQ 支持两种消息过滤方式:
- Tag 过滤:根据消息的 Tag 进行过滤,简单高效;
- SQL92 过滤:根据消息的属性进行 SQL 表达式过滤,灵活强大。
4.8.1 Tag 过滤
Tag 过滤的实现原理:
- Producer 发送消息时,设置消息的 Tag;
- Broker 存储消息时,将 Tag 的哈希值写入 ConsumeQueue;
- Consumer 订阅消息时,指定 Tag(或 Tag 表达式,如
TagA || TagB); - Broker 收到 Consumer 的拉取请求时,根据 Tag 的哈希值在 ConsumeQueue 中过滤消息,只返回符合条件的消息。
Tag 过滤是在 Broker 端进行的,减少了网络传输量,提升了性能。
4.8.2 SQL92 过滤
SQL92 过滤的实现原理:
- Producer 发送消息时,设置消息的属性(如
putUserProperty("age", "18")); - Consumer 订阅消息时,指定 SQL92 表达式(如
age > 18); - Broker 收到 Consumer 的拉取请求时,首先根据 Tag 过滤消息,然后将符合 Tag 条件的消息从 CommitLog 中读取出来,再根据 SQL92 表达式过滤消息,只返回符合条件的消息。
SQL92 过滤是在 Broker 端进行的,但需要读取 CommitLog 中的消息属性,性能比 Tag 过滤低,但更灵活。
5. Java 客户端快速入门
5.1 环境准备
5.1.1 引入依赖
在 Maven 项目的 pom.xml 中引入 RocketMQ 客户端依赖:
1
2
3
4
5
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>4.9.5</version>
</dependency>
注意:请使用最新的稳定版本,本文使用 4.9.5 版本。
5.1.2 启动 NameServer 和 Broker
在使用 Java 客户端之前,需要先启动 NameServer 和 Broker:
- 启动 NameServer:
1 2 3 4 5
# Linux/Mac sh bin/mqnamesrv # Windows bin/mqnamesrv.cmd
- 启动 Broker:
1 2 3 4 5
# Linux/Mac sh bin/mqbroker -n localhost:9876 # Windows bin/mqbroker.cmd -n localhost:9876
注意:
-n参数指定 NameServer 的地址。
5.2 普通消息发送
5.2.1 同步发送
同步发送是指 Producer 发送消息后,等待 Broker 返回结果,可靠性高但会阻塞线程。适用于对可靠性要求较高的场景,如订单创建。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class SyncProducer {
public static void main(String[] args) throws Exception {
// 1. 创建生产者实例,指定生产者组名
DefaultMQProducer producer = new DefaultMQProducer("sync_producer_group");
// 2. 设置 NameServer 地址,多个地址用分号分隔
producer.setNamesrvAddr("localhost:9876");
// 3. 启动生产者
producer.start();
for (int i = 0; i < 10; i++) {
// 4. 创建消息对象,指定 Topic、Tag、消息体
Message message = new Message(
"test_topic", // Topic
"TagA", // Tag
("Hello RocketMQ " + i).getBytes() // 消息体
);
// 5. 同步发送消息,等待返回结果
SendResult sendResult = producer.send(message);
// 6. 打印发送结果
System.out.printf("发送结果:%s%n", sendResult);
}
// 7. 关闭生产者
producer.shutdown();
}
}
代码说明:
DefaultMQProducer:生产者的核心类,用于发送消息;producer.setNamesrvAddr():设置 NameServer 的地址,多个地址用分号分隔(如localhost:9876;localhost:9877);Message:消息对象,构造函数参数依次为 Topic、Tag、消息体;producer.send():同步发送消息,返回SendResult,包含发送状态、MsgID、MessageQueue 等信息;producer.shutdown():关闭生产者,释放资源。
5.2.2 异步发送
异步发送是指 Producer 发送消息后,不等待结果,通过回调函数处理发送结果,性能高。适用于对响应时间要求较高的场景,如日志收集。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class AsyncProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("async_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
// 设置异步发送的重试次数
producer.setRetryTimesWhenSendAsyncFailed(3);
for (int i = 0; i < 10; i++) {
final int index = i;
Message message = new Message(
"test_topic",
"TagA",
("Hello RocketMQ Async " + i).getBytes()
);
// 异步发送消息,注册回调函数
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.printf("发送成功,索引:%d,结果:%s%n", index, sendResult);
}
@Override
public void onException(Throwable e) {
System.out.printf("发送失败,索引:%d,异常:%s%n", index, e);
}
});
}
// 等待异步发送完成,避免直接关闭生产者
Thread.sleep(10000);
producer.shutdown();
}
}
代码说明:
producer.send(message, new SendCallback()):异步发送消息,注册SendCallback回调函数;onSuccess():发送成功时调用;onException():发送失败时调用;Thread.sleep(10000):等待异步发送完成,避免直接关闭生产者导致消息没发送完。
5.2.3 Oneway 发送
Oneway 发送是指 Producer 发送消息后,不等待结果,也不处理回调,性能最高但可靠性最低。适用于对可靠性要求较低的场景,如日志收集。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.common.message.Message;
public class OnewayProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("oneway_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
for (int i = 0; i < 10; i++) {
Message message = new Message(
"test_topic",
"TagA",
("Hello RocketMQ Oneway " + i).getBytes()
);
// Oneway 发送消息
producer.sendOneway(message);
System.out.printf("发送消息:%d%n", i);
}
producer.shutdown();
}
}
代码说明:
producer.sendOneway(message):Oneway 发送消息,不返回结果。
5.3 普通消息消费
5.3.1 Push 模式集群消费
Push 模式是指 Broker 主动将消息推送给 Consumer,实时性高。集群消费是指同一个消费者组的 Consumer 共同消费消息,每个消息只被一个 Consumer 消费。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class PushConsumer {
public static void main(String[] args) throws Exception {
// 1. 创建消费者实例,指定消费者组名
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("push_consumer_group");
// 2. 设置 NameServer 地址
consumer.setNamesrvAddr("localhost:9876");
// 3. 订阅 Topic,指定 Tag 过滤(* 表示不过滤)
consumer.subscribe("test_topic", "*");
// 4. 注册消息监听器,并发消费
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
try {
for (MessageExt message : messages) {
// 处理消息
String body = new String(message.getBody());
System.out.printf("收到消息:MsgID=%s,Body=%s%n", message.getMsgId(), body);
}
// 消费成功,返回 SUCCESS
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
// 消费失败,返回 RECONSUME_LATER,会重试
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
// 5. 启动消费者
consumer.start();
System.out.println("消费者已启动");
}
}
代码说明:
DefaultMQPushConsumer:Push 模式消费者的核心类;consumer.subscribe("test_topic", "*"):订阅 Topic,第二个参数是 Tag 表达式,*表示不过滤,TagA表示只消费 TagA 的消息,TagA || TagB表示消费 TagA 或 TagB 的消息;MessageListenerConcurrently:并发消费监听器,多个线程同时消费不同的 Queue;consumeMessage():消费消息的方法,参数messages是消息列表,context是消费上下文;ConsumeConcurrentlyStatus.CONSUME_SUCCESS:消费成功,提交 Offset;ConsumeConcurrentlyStatus.RECONSUME_LATER:消费失败,稍后重试。
5.3.2 Push 模式广播消费
广播消费是指同一个消费者组的每个 Consumer 都会消费所有消息,即每个消息会被消费多次。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
import java.util.List;
public class BroadcastPushConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("broadcast_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
// 设置消费模式为广播消费
consumer.setMessageModel(MessageModel.BROADCASTING);
consumer.subscribe("test_topic", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
try {
for (MessageExt message : messages) {
String body = new String(message.getBody());
System.out.printf("收到广播消息:MsgID=%s,Body=%s%n", message.getMsgId(), body);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
consumer.start();
System.out.println("广播消费者已启动");
}
}
代码说明:
consumer.setMessageModel(MessageModel.BROADCASTING):设置消费模式为广播消费,默认为集群消费。
5.3.3 Pull 模式消费
Pull 模式是指 Consumer 主动从 Broker 拉取消息,灵活性高。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
import org.apache.rocketmq.client.consumer.PullResult;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
import java.util.HashMap;
import java.util.Map;
import java.util.Set;
public class PullConsumer {
// 存储每个 Queue 的 Offset
private static final Map<MessageQueue, Long> OFFSET_TABLE = new HashMap<>();
public static void main(String[] args) throws Exception {
DefaultMQPullConsumer consumer = new DefaultMQPullConsumer("pull_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.start();
// 获取 Topic 的 Queue 列表
Set<MessageQueue> messageQueues = consumer.fetchSubscribeMessageQueues("test_topic");
while (true) {
for (MessageQueue messageQueue : messageQueues) {
// 获取 Queue 的 Offset
long offset = getOffset(messageQueue);
// 拉取消息
PullResult pullResult = consumer.pull(
messageQueue, // Queue
"*", // Tag 表达式
offset, // 起始 Offset
10 // 最大拉取数量
);
// 处理拉取结果
switch (pullResult.getPullStatus()) {
case FOUND:
// 拉取到消息
for (MessageExt message : pullResult.getMsgFoundList()) {
String body = new String(message.getBody());
System.out.printf("收到消息:MsgID=%s,Body=%s%n", message.getMsgId(), body);
}
// 更新 Offset
updateOffset(messageQueue, pullResult.getNextBeginOffset());
break;
case NO_NEW_MSG:
// 没有新消息
break;
case NO_MATCHED_MSG:
// 没有匹配的消息
break;
case OFFSET_ILLEGAL:
// Offset 非法
updateOffset(messageQueue, pullResult.getNextBeginOffset());
break;
}
}
// 等待 1s 再拉取
Thread.sleep(1000);
}
}
// 获取 Queue 的 Offset
private static long getOffset(MessageQueue messageQueue) {
Long offset = OFFSET_TABLE.get(messageQueue);
if (offset == null) {
return 0;
}
return offset;
}
// 更新 Queue 的 Offset
private static void updateOffset(MessageQueue messageQueue, long offset) {
OFFSET_TABLE.put(messageQueue, offset);
}
}
代码说明:
DefaultMQPullConsumer:Pull 模式消费者的核心类;consumer.fetchSubscribeMessageQueues("test_topic"):获取 Topic 的 Queue 列表;consumer.pull():拉取消息,参数依次为 Queue、Tag 表达式、起始 Offset、最大拉取数量;PullStatus:拉取状态,包括 FOUND(拉取到消息)、NO_NEW_MSG(没有新消息)、NO_MATCHED_MSG(没有匹配的消息)、OFFSET_ILLEGAL(Offset 非法);OFFSET_TABLE:存储每个 Queue 的 Offset,实际生产中应该持久化到数据库或 Redis。
6. 高级特性实战
6.1 顺序消息
6.1.1 顺序消息发送
顺序消息发送需要自定义 MessageQueueSelector,将同一组顺序消息发送到同一个 Queue 中。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.MessageQueueSelector;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageQueue;
import java.util.List;
public class OrderlyProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("orderly_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
// 模拟订单数据
String[] orderIds = {"1001", "1002", "1003"};
String[] orderTypes = {"创建", "支付", "发货", "收货"};
for (String orderId : orderIds) {
for (String orderType : orderTypes) {
Message message = new Message(
"orderly_topic",
"TagA",
orderId,
("订单 " + orderId + " " + orderType).getBytes()
);
// 发送顺序消息,使用自定义 MessageQueueSelector
SendResult sendResult = producer.send(message, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> queues, Message msg, Object arg) {
// 根据订单 ID 选择 Queue,确保同一个订单的消息发送到同一个 Queue
String orderId = (String) arg;
int index = Math.abs(orderId.hashCode()) % queues.size();
return queues.get(index);
}
}, orderId);
System.out.printf("发送顺序消息:%s,结果:%s%n", new String(message.getBody()), sendResult);
}
}
producer.shutdown();
}
}
代码说明:
MessageQueueSelector:消息队列选择器,用于选择发送消息的 Queue;select():选择 Queue 的方法,参数queues是 Queue 列表,msg是消息对象,arg是传入的参数(这里是订单 ID);Math.abs(orderId.hashCode()) % queues.size():根据订单 ID 的哈希值选择 Queue,确保同一个订单的消息发送到同一个 Queue。
6.1.2 顺序消息消费
顺序消息消费需要使用 MessageListenerOrderly 监听器,保证同一个 Queue 的消息由同一个线程消费。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class OrderlyPushConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("orderly_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("orderly_topic", "*");
// 注册顺序消息监听器
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> messages, ConsumeOrderlyContext context) {
try {
for (MessageExt message : messages) {
String body = new String(message.getBody());
System.out.printf("收到顺序消息:MsgID=%s,Body=%s,QueueId=%d%n",
message.getMsgId(), body, message.getQueueId());
}
// 消费成功,返回 SUCCESS
return ConsumeOrderlyStatus.SUCCESS;
} catch (Exception e) {
e.printStackTrace();
// 消费失败,暂停当前 Queue 的消费,稍后重试
return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
}
}
});
consumer.start();
System.out.println("顺序消费者已启动");
}
}
代码说明:
MessageListenerOrderly:顺序消费监听器,保证同一个 Queue 的消息由同一个线程消费;ConsumeOrderlyStatus.SUCCESS:消费成功,提交 Offset;ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT:消费失败,暂停当前 Queue 的消费,稍后重试。
6.2 事务消息
6.2.1 事务消息发送
事务消息发送需要使用 TransactionMQProducer,并实现 TransactionListener 接口。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.concurrent.*;
public class TransactionProducer {
public static void main(String[] args) throws Exception {
// 1. 创建事务生产者实例,指定生产者组名
TransactionMQProducer producer = new TransactionMQProducer("transaction_producer_group");
// 2. 设置 NameServer 地址
producer.setNamesrvAddr("localhost:9876");
// 3. 创建线程池,用于执行本地事务和事务状态回查
ExecutorService executorService = new ThreadPoolExecutor(
2, 5, 100, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(2000),
r -> {
Thread thread = new Thread(r);
thread.setName("transaction-check-thread");
return thread;
}
);
producer.setExecutorService(executorService);
// 4. 设置事务监听器
producer.setTransactionListener(new TransactionListener() {
// 执行本地事务
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 模拟本地事务执行,比如操作数据库
String businessId = (String) arg;
System.out.printf("执行本地事务,BusinessID=%s%n", businessId);
// 模拟事务成功或失败,这里随机返回
int random = ThreadLocalRandom.current().nextInt(3);
switch (random) {
case 0:
// 事务成功
System.out.println("本地事务成功");
return LocalTransactionState.COMMIT_MESSAGE;
case 1:
// 事务失败
System.out.println("本地事务失败");
return LocalTransactionState.ROLLBACK_MESSAGE;
case 2:
// 事务状态未知,等待回查
System.out.println("本地事务状态未知,等待回查");
return LocalTransactionState.UNKNOW;
default:
return LocalTransactionState.UNKNOW;
}
} catch (Exception e) {
e.printStackTrace();
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
// 事务状态回查
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
try {
String businessId = msg.getKeys();
System.out.printf("事务状态回查,BusinessID=%s%n", businessId);
// 模拟查询本地事务状态,比如查询数据库
// 这里假设回查时事务成功
System.out.println("回查结果:本地事务成功");
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
e.printStackTrace();
return LocalTransactionState.UNKNOW;
}
}
});
// 5. 启动事务生产者
producer.start();
// 6. 发送事务消息
String businessId = "123456";
Message message = new Message(
"transaction_topic",
"TagA",
businessId, // Key 设置为业务 ID,方便回查
("Hello Transaction Message, BusinessID=" + businessId).getBytes()
);
// 发送事务消息,arg 参数会传递给 executeLocalTransaction 方法
SendResult sendResult = producer.sendMessageInTransaction(message, businessId);
System.out.printf("事务消息发送结果:%s%n", sendResult);
// 等待回查完成
Thread.sleep(100000);
producer.shutdown();
executorService.shutdown();
}
}
代码说明:
TransactionMQProducer:事务生产者的核心类;TransactionListener:事务监听器,包含两个方法:executeLocalTransaction():执行本地事务,返回LocalTransactionState.COMMIT_MESSAGE(提交)、LocalTransactionState.ROLLBACK_MESSAGE(回滚)或LocalTransactionState.UNKNOW(未知);checkLocalTransaction():事务状态回查,返回同样的状态;
producer.sendMessageInTransaction(message, businessId):发送事务消息,第二个参数arg会传递给executeLocalTransaction()方法。
6.2.2 事务消息消费
事务消息的消费与普通消息的消费一样,只需要订阅目标 Topic 即可。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class TransactionPushConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("transaction_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("transaction_topic", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
try {
for (MessageExt message : messages) {
String body = new String(message.getBody());
System.out.printf("收到事务消息:MsgID=%s,Body=%s%n", message.getMsgId(), body);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
consumer.start();
System.out.println("事务消费者已启动");
}
}
6.3 延迟消息
6.3.1 延迟消息发送
延迟消息发送只需要设置消息的延迟级别即可。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class DelayProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("delay_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
for (int i = 0; i < 3; i++) {
Message message = new Message(
"delay_topic",
"TagA",
("Hello Delay Message " + i).getBytes()
);
// 设置延迟级别,3 表示 10s(延迟级别:1s 5s 10s 30s 1m...)
message.setDelayTimeLevel(3);
SendResult sendResult = producer.send(message);
System.out.printf("发送延迟消息:%d,发送时间:%d,结果:%s%n",
i, System.currentTimeMillis(), sendResult);
}
producer.shutdown();
}
}
代码说明:
message.setDelayTimeLevel(3):设置延迟级别,延迟级别对应关系如下:级别 延迟时间 1 1s 2 5s 3 10s 4 30s 5 1m 6 2m 7 3m 8 4m 9 5m 10 6m 11 7m 12 8m 13 9m 14 10m 15 20m 16 30m 17 1h 18 2h
6.3.2 延迟消息消费
延迟消息的消费与普通消息的消费一样,只需要订阅目标 Topic 即可。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class DelayPushConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("delay_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("delay_topic", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
try {
for (MessageExt message : messages) {
String body = new String(message.getBody());
System.out.printf("收到延迟消息:MsgID=%s,Body=%s,接收时间:%d%n",
message.getMsgId(), body, System.currentTimeMillis());
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
consumer.start();
System.out.println("延迟消费者已启动");
}
}
6.4 消息过滤
6.4.1 Tag 过滤
Tag 过滤只需要在订阅 Topic 时指定 Tag 表达式即可。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class TagFilterPushConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("tag_filter_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
// 订阅 Topic,指定 Tag 表达式:只消费 TagA 或 TagB 的消息
consumer.subscribe("test_topic", "TagA || TagB");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
try {
for (MessageExt message : messages) {
String body = new String(message.getBody());
System.out.printf("收到消息:MsgID=%s,Tag=%s,Body=%s%n",
message.getMsgId(), message.getTags(), body);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
consumer.start();
System.out.println("Tag 过滤消费者已启动");
}
}
代码说明:
consumer.subscribe("test_topic", "TagA || TagB"):订阅 Topic,指定 Tag 表达式,TagA || TagB表示消费 TagA 或 TagB 的消息,*表示不过滤。
6.4.2 SQL92 过滤
SQL92 过滤需要在订阅 Topic 时指定 SQL92 表达式,并且 Broker 需要开启 SQL92 过滤功能。
首先,修改 Broker 的配置文件 conf/broker.conf,添加以下配置:
1
2
# 开启 SQL92 过滤
enablePropertyFilter=true
然后,重启 Broker。
代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class SqlFilterPushConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("sql_filter_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
// 订阅 Topic,指定 SQL92 表达式:只消费 age > 18 的消息
consumer.subscribe("test_topic", MessageSelector.bySql("age > 18"));
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
try {
for (MessageExt message : messages) {
String body = new String(message.getBody());
String age = message.getUserProperty("age");
System.out.printf("收到消息:MsgID=%s,age=%s,Body=%s%n",
message.getMsgId(), age, body);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
});
consumer.start();
System.out.println("SQL92 过滤消费者已启动");
}
}
发送带属性的消息代码示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
public class SqlFilterProducer {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("sql_filter_producer_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
for (int i = 0; i < 10; i++) {
Message message = new Message(
"test_topic",
"TagA",
("Hello SQL Filter Message " + i).getBytes()
);
// 设置消息属性
message.putUserProperty("age", String.valueOf(15 + i));
SendResult sendResult = producer.send(message);
System.out.printf("发送消息:%d,age=%d,结果:%s%n", i, 15 + i, sendResult);
}
producer.shutdown();
}
}
代码说明:
message.putUserProperty("age", String.valueOf(15 + i)):设置消息属性;MessageSelector.bySql("age > 18"):指定 SQL92 表达式,只消费 age > 18 的消息;- SQL92 表达式支持的语法:
- 比较运算符:
=、>、<、>=、<=、<>、BETWEEN、IN; - 逻辑运算符:
AND、OR、NOT; - 常量:字符串(用单引号)、数字、布尔值(TRUE、FALSE);
- 函数:
IS NULL、IS NOT NULL。
- 比较运算符:
7. 运维与监控
7.1 部署
7.1.1 NameServer 部署
NameServer 是无状态的,可以部署多个实例,实现高可用。
部署步骤:
- 下载 RocketMQ 安装包:https://rocketmq.apache.org/zh/download
- 解压安装包:
1 2
unzip rocketmq-all-4.9.5-bin-release.zip cd rocketmq-all-4.9.5-bin-release - 修改 NameServer 配置文件
conf/namesrv.properties(可选):1 2 3 4
# NameServer 监听端口 listenPort=9876 # NameServer 日志路径 rocketmq.logs.dir=/usr/local/rocketmq/logs
- 启动 NameServer:
1 2
# 后台启动 nohup sh bin/mqnamesrv > /dev/null 2>&1 &
- 查看 NameServer 日志:
1
tail -f logs/namesrv.log
7.1.2 Broker 部署
Broker 可以部署为主从架构,实现高可用。
部署步骤:
- 修改 Master 配置文件
conf/broker.conf:1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30
# Broker 集群名称 brokerClusterName=DefaultCluster # Broker 名称,主从 Broker 的 brokerName 要相同 brokerName=broker-a # Broker ID,0 表示主节点,非 0 表示从节点 brokerId=0 # NameServer 地址,多个用分号分隔 namesrvAddr=localhost:9876;localhost:9877 # 自动创建 Topic,生产环境建议关闭 autoCreateTopicEnable=true # 自动创建订阅组,生产环境建议关闭 autoCreateSubscriptionGroup=true # CommitLog 存储路径 storePathCommitLog=/usr/local/rocketmq/store/commitlog # ConsumeQueue 存储路径 storePathConsumeQueue=/usr/local/rocketmq/store/consumequeue # IndexFile 存储路径 storePathIndex=/usr/local/rocketmq/store/index # 刷盘策略,ASYNC_FLUSH 表示异步刷盘,SYNC_FLUSH 表示同步刷盘 flushDiskType=ASYNC_FLUSH # 同步刷盘超时时间,仅当 flushDiskType=SYNC_FLUSH 时有效 syncFlushTimeout=5000 # 主从复制策略,ASYNC_MASTER 表示异步复制,SYNC_MASTER 表示同步复制 brokerRole=ASYNC_MASTER # 消息最大大小,默认 4M maxMessageSize=4194304 # 延迟消息的延迟级别 messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h # 开启 SQL92 过滤 enablePropertyFilter=true
- 启动 Master:
1
nohup sh bin/mqbroker -c conf/broker.conf > /dev/null 2>&1 &
- 修改 Slave 配置文件
conf/broker-slave.conf:1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
brokerClusterName=DefaultCluster brokerName=broker-a # Broker ID,1 表示从节点 brokerId=1 namesrvAddr=localhost:9876;localhost:9877 autoCreateTopicEnable=true autoCreateSubscriptionGroup=true storePathCommitLog=/usr/local/rocketmq/store-slave/commitlog storePathConsumeQueue=/usr/local/rocketmq/store-slave/consumequeue storePathIndex=/usr/local/rocketmq/store-slave/index flushDiskType=ASYNC_FLUSH # 从节点的 brokerRole 为 SLAVE brokerRole=SLAVE maxMessageSize=4194304 messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h enablePropertyFilter=true
- 启动 Slave:
1
nohup sh bin/mqbroker -c conf/broker-slave.conf > /dev/null 2>&1 &
- 查看 Broker 日志:
1
tail -f logs/broker.log
7.2 监控
7.2.1 RocketMQ Console 部署
RocketMQ Console 是 RocketMQ 的可视化控制台,用于监控和管理 RocketMQ。
部署步骤:
- 下载 RocketMQ Console 源码:https://github.com/apache/rocketmq-dashboard
- 修改
src/main/resources/application.properties:1 2 3 4
# NameServer 地址 rocketmq.config.namesrvAddr=localhost:9876;localhost:9877 # 控制台端口 server.port=8080
- 打包:
1
mvn clean package -Dmaven.test.skip=true
- 启动:
1
java -jar target/rocketmq-dashboard-1.0.0.jar - 访问:http://localhost:8080
7.2.2 RocketMQ Console 功能
RocketMQ Console 提供以下功能:
- Dashboard:查看集群概览,包括 Broker 数量、Topic 数量、Producer 数量、Consumer 数量等;
- Topic:查看 Topic 列表,创建、删除 Topic,查看 Topic 的 Queue 分布、消息堆积情况等;
- Producer:查看 Producer 列表,查看 Producer 的发送状态等;
- Consumer:查看 Consumer 列表,查看 Consumer 的消费状态、消息堆积情况等;
- Message:根据 MsgID 或 Key 查询消息,查看消息的详细信息;
- Cluster:查看 Broker 列表,查看 Broker 的状态、配置等;
- Tools:提供消息发送、消息消费等工具。
7.3 常见问题排查
7.3.1 消息堆积
消息堆积是指 Consumer 消费速度跟不上 Producer 发送速度,导致大量消息积压在 Broker 中。
排查步骤:
- 查看 Consumer 的消费状态:使用 RocketMQ Console 查看 Consumer 的消息堆积情况;
- 检查 Consumer 的消费逻辑:查看 Consumer 的日志,确认是否有消费异常或消费速度慢的问题;
- 增加 Consumer 数量:如果是消费能力不足,可以增加同一个消费者组的 Consumer 数量;
- 优化消费逻辑:如果是消费逻辑慢,可以优化消费逻辑,如批量消费、异步处理等;
- 临时扩容:如果是突发流量,可以临时扩容 Consumer,等流量下降后再缩容。
7.3.2 消息丢失
消息丢失是指 Producer 发送的消息没有被 Consumer 消费到。
排查步骤:
- 检查 Producer 的发送状态:查看 Producer 的日志,确认消息是否发送成功;
- 检查 Broker 的存储:使用 RocketMQ Console 根据 MsgID 查询消息,确认消息是否存储在 Broker 中;
- 检查 Consumer 的订阅:确认 Consumer 是否订阅了正确的 Topic 和 Tag;
- 检查 Consumer 的消费状态:查看 Consumer 的日志,确认是否有消费异常;
- 检查 Broker 的配置:确认 Broker 的刷盘策略和主从复制策略是否符合要求(如果要求高可靠性,应该使用同步刷盘和同步复制)。
7.3.3 消费失败
消费失败是指 Consumer 消费消息时出现异常,导致消息重试或进入死信队列。
排查步骤:
- 查看 Consumer 的日志:确认消费异常的原因;
- 查看重试队列:使用 RocketMQ Console 查看重试队列中的消息;
- 查看死信队列:如果消息重试超过 16 次,会进入死信队列,使用 RocketMQ Console 查看死信队列中的消息;
- 修复消费逻辑:根据异常原因修复消费逻辑;
- 重新消费死信消息:修复消费逻辑后,可以将死信消息重新发送到目标 Topic,重新消费。
8. 最佳实践
8.1 Producer 最佳实践
8.1.1 消息大小限制
RocketMQ 默认消息最大大小为 4M,超过 4M 的消息会发送失败。如果需要发送大消息,可以:
- 调整 Broker 的
maxMessageSize配置; - 将大消息存储在文件系统或对象存储中,只发送文件路径。
8.1.2 发送方式选择
- 同步发送:适用于对可靠性要求较高的场景,如订单创建;
- 异步发送:适用于对响应时间要求较高的场景,如日志收集;
- Oneway 发送:适用于对可靠性要求较低的场景,如日志收集。
8.1.3 重试策略配置
- 同步发送:默认重试 2 次,可以通过
producer.setRetryTimesWhenSendFailed(3)调整; - 异步发送:默认重试 2 次,可以通过
producer.setRetryTimesWhenSendAsyncFailed(3)调整。
8.1.4 Key 的设置
建议为消息设置 Key,Key 可以是业务 ID(如订单 ID),方便根据 Key 查询消息。
8.1.5 Tag 的使用
建议使用 Tag 对消息进行分类,方便 Consumer 进行消息过滤。
8.2 Consumer 最佳实践
8.2.1 消费模式选择
- 集群消费:适用于需要负载均衡的场景,每个消息只被一个 Consumer 消费;
- 广播消费:适用于需要每个 Consumer 都消费所有消息的场景,如配置更新。
8.2.2 批量消费配置
可以通过 consumer.setConsumeMessageBatchMaxSize(10) 设置批量消费的最大数量,提升消费性能。
8.2.3 重试次数设置
默认重试次数为 16 次,可以通过 consumer.setMaxReconsumeTimes(5) 调整重试次数。
8.2.4 死信队列处理
如果消息重试超过最大次数,会进入死信队列。建议:
- 监控死信队列,及时发现问题;
- 修复消费逻辑后,将死信消息重新发送到目标 Topic,重新消费。
8.2.5 消费幂等性保证
RocketMQ 可能会重复投递消息(如消费失败重试、网络问题导致 ACK 没收到等),所以 Consumer 需要保证幂等性。
实现幂等性的方法:
- 数据库唯一索引:使用业务 ID 作为数据库的唯一索引,避免重复插入;
- Redis 记录:使用 Redis 记录已消费的 MsgID,处理消息前先检查是否已消费;
- 状态机:使用状态机记录业务状态,避免重复处理。
代码示例(使用 Redis 记录已消费的 MsgID):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import redis.clients.jedis.Jedis;
import java.util.List;
public class IdempotentPushConsumer {
private static final String REDIS_KEY_PREFIX = "rocketmq:consumed:";
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("idempotent_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("test_topic", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messages, ConsumeConcurrentlyContext context) {
Jedis jedis = null;
try {
jedis = new Jedis("localhost", 6379);
for (MessageExt message : messages) {
String msgId = message.getMsgId();
String redisKey = REDIS_KEY_PREFIX + msgId;
// 检查是否已消费
if (jedis.exists(redisKey)) {
System.out.printf("消息已消费,MsgID=%s%n", msgId);
continue;
}
// 处理消息
String body = new String(message.getBody());
System.out.printf("处理消息:MsgID=%s,Body=%s%n", msgId, body);
// 标记为已消费,设置过期时间,比如 24 小时
jedis.setex(redisKey, 86400, "1");
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
} finally {
if (jedis != null) {
jedis.close();
}
}
}
});
consumer.start();
System.out.println("幂等消费者已启动");
}
}
8.3 Topic 和 Tag 设计
8.3.1 Topic 划分原则
- 按业务划分:不同的业务使用不同的 Topic,如订单 Topic、支付 Topic、物流 Topic;
- 按优先级划分:不同优先级的消息使用不同的 Topic,如高优先级 Topic、低优先级 Topic;
- 避免过多 Topic:Topic 数量过多会增加 Broker 的管理成本,建议 Topic 数量控制在几百个以内。
8.3.2 Tag 使用场景
- 消息分类:对 Topic 下的消息进行分类,如订单 Topic 下的创建订单、支付订单、取消订单;
- 消息过滤:Consumer 根据 Tag 进行消息过滤,只消费感兴趣的消息。
8.4 性能优化
8.4.1 Broker 配置优化
- 刷盘策略:如果对性能要求较高,使用异步刷盘(ASYNC_FLUSH);如果对可靠性要求较高,使用同步刷盘(SYNC_FLUSH);
- 主从复制策略:如果对性能要求较高,使用异步复制(ASYNC_MASTER);如果对可靠性要求较高,使用同步复制(SYNC_MASTER);
- 内存配置:调整 JVM 内存参数,如
-Xms8g -Xmx8g -Xmn4g; - CommitLog 大小:调整 CommitLog 文件大小,如
mapedFileSizeCommitLog=1073741824(1G)。
8.4.2 Producer 配置优化
- 批量发送:使用批量发送提升性能,如
producer.send(messages); - 异步发送:使用异步发送提升性能;
- 线程池配置:调整发送线程池大小,如
producer.setSendMessageThreadPoolNums(16)。
8.4.3 Consumer 配置优化
- 批量消费:使用批量消费提升性能,如
consumer.setConsumeMessageBatchMaxSize(10); - 并发消费:使用并发消费提升性能,如
consumer.setConsumeThreadMin(10)、consumer.setConsumeThreadMax(20); - 拉取大小:调整拉取大小,如
consumer.setPullBatchSize(32)。
总结
RocketMQ 是一款功能强大、性能优异的分布式消息中间件,具有高并发、高可靠、高可用等特性。本文详细介绍了 RocketMQ 的核心概念、架构设计、核心特性、Java 客户端使用、高级特性实战、运维与监控、最佳实践等内容,希望能帮助读者快速上手 RocketMQ,并在实际项目中用好 RocketMQ。
如果需要进一步学习 RocketMQ,可以参考官方文档:https://rocketmq.apache.org/zh/docs/