Article
第一章:RocketMQ 概述
1.1 消息中间件简介
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 消息中间件 | 是指利用高效可靠的异步消息传递机制在分布式系统之间进行数据交流的中间件系统。 | 选择消息中间件时需考虑吞吐量、延迟、可靠性、运维成本等因素。 |
| 异步通信 | 发送方发送消息后无需等待接收方响应,提升系统响应速度和解耦程度。 | 异步可能导致消息丢失或顺序错乱,需配合持久化和确认机制保障可靠性。 |
| 解耦 | 系统之间通过消息传递进行交互,降低模块间的直接依赖。 | 过度解耦可能增加系统复杂性,需合理设计消息格式与协议。 |
| 削峰填谷 | 在高并发场景下,通过消息队列缓冲请求,防止后端服务被瞬间流量压垮。 | 需设置合理的队列长度和消费速度,避免消息积压导致延迟过高或内存溢出。 |
| 可靠传输 | 支持消息持久化、重试机制、事务消息等,确保消息不丢失。 | 需权衡性能与可靠性,例如同步刷盘比异步刷盘更安全但性能较低。 |
1.2 RocketMQ 简介与核心特性
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| RocketMQ | 阿里巴巴开源的分布式消息中间件,现为 Apache 顶级项目,具备高吞吐、低延迟、高可用等特性。 | 适用于大规模分布式系统,学习曲线略高于 RabbitMQ,但性能更强。 |
| 高吞吐量 | 单机可支持数十万消息/秒的处理能力,适合海量消息场景。 | 实际吞吐受网络、磁盘IO、消息大小等因素影响,需结合压测评估。 |
| 低延迟 | 普通消息延迟可控制在毫秒级,满足实时性要求较高的业务。 | 延迟受 Broker 刷盘策略、主从同步方式等配置影响,需根据业务需求调整。 |
| 高可用性 | 支持主从架构、自动故障切换、NameServer 集群部署,保障服务持续可用。 | 需合理配置主从复制方式(同步双写/异步复制)以平衡性能与数据安全。 |
| 顺序消息 | 支持 FIFO(先进先出)的消息传递,保证同一队列中的消息有序消费。 | 仅保证局部有序,全局有序需谨慎使用,可能影响并发性能。 |
| 事务消息 | 支持分布式事务,通过两阶段提交机制实现本地事务与消息发送的一致性。 | 需实现事务回查逻辑,避免事务状态未知导致消息无法提交或回滚。 |
| 定时/延时消息 | 支持消息在指定时间后被消费,适用于延迟任务调度场景。 | 延时等级固定(如1s, 5s, 10s…2h),不支持任意时间精度,需合理规划业务逻辑。 |
1.3 RocketMQ 架构组成(NameServer、Broker、Producer、Consumer)
| 组件名称 | 说明 | 注意事项 |
|---|---|---|
| NameServer | 轻量级服务发现组件,管理 Broker 的路由信息,Producer 和 Consumer 通过其获取连接地址。 | NameServer 无状态,可集群部署;Broker 定期向所有 NameServer 心跳注册。 |
| Broker | 消息中转节点,负责存储消息、转发消息,支持主从部署以实现高可用。 | 主 Broker 可读写,从 Broker 仅可读;主从间通过复制同步消息数据。 |
| Producer | 消息生产者,负责创建并发送消息到指定的 Topic。 | 生产者需指定 Namesrv 地址列表,启动后从 NameServer 拉取路由信息。 |
| Consumer | 消息消费者,从 Broker 拉取消息并进行处理。 | 消费者需注册消费组(Consumer Group),支持集群或广播模式消费。 |
| Topic | 消息的主题分类,是消息的逻辑分类单元,Producer 发送到 Topic,Consumer 订阅 Topic。 | 同一个 Topic 的消息可分布在多个 Broker 上,提高并发处理能力。 |
| Message Queue | Topic 下的分区,每个 Topic 可分为多个 Message Queue,用于并行生产和消费。 | Queue 数量影响并发度,建议根据消费能力合理设置。 |
| Consumer Group | 消费者组,同一组内的消费者共同消费一个 Topic 的消息,实现负载均衡。 | 集群模式下,一个 Queue 只能被组内一个消费者消费;广播模式下,所有消费者都消费全量消息。 |
1.4 应用场景与优势对比(Kafka、RabbitMQ)
| 对比维度 | RocketMQ | Kafka | RabbitMQ | 说明 |
|---|---|---|---|---|
| 开发语言 | Java | Scala | Erlang | 与 Java 系统集成更友好 |
| 吞吐量 | 高(单机10万+/秒) | 极高(适合大数据场景) | 中等(万级) | Kafka 在日志收集等大数据场景更具优势 |
| 延迟 | 毫秒级 | 毫秒级~秒级 | 毫秒级 | 三者均可满足大多数实时性需求 |
| 顺序消息 | 支持(局部有序) | 支持(Partition 内有序) | 不支持原生顺序 | RocketMQ 和 Kafka 均可保障分区/队列内有序 |
| 事务消息 | 支持(两阶段提交 + 回查) | 支持(幂等 + 事务 API) | 不支持 | RocketMQ 更适合金融级事务场景 |
| 延时消息 | 支持(内置延时等级) | 不支持(需外部调度) | 支持(插件 rabbitmq_delayed_message) | RocketMQ 原生支持,使用更便捷 |
| 运维复杂度 | 中等 | 较高 | 较低 | RabbitMQ 管理界面友好,Kafka 配置复杂 |
| 典型应用场景 | 电商交易、订单系统、金融支付 | 日志收集、流式处理、大数据管道 | 任务调度、内部系统通信 | 根据业务需求选择合适中间件 |
第二章:环境搭建与快速入门
2.1 安装与启动 RocketMQ 服务(单机/集群)
| 操作步骤 | 说明 | 注意事项 |
|---|---|---|
| 下载安装包 | 从 Apache 官网下载 RocketMQ 最新版二进制包(如 rocketmq-all-4.9.4-bin-release.zip) | 建议使用稳定版本,避免使用快照版本用于生产环境。 |
| 解压安装 | tar -zxvf rocketmq-all-x.x.x-bin-release.zip | 解压路径避免包含空格或中文字符。 |
| 启动 NameServer | nohup sh bin/mqnamesrv & | 默认监听 9876 端口,可通过 -p 参数指定配置文件。 |
| 启动 Broker | nohup sh bin/mqbroker -n localhost:9876 autoCreateTopicEnable=true & | autoCreateTopicEnable=true 允许自动创建 Topic,生产环境建议关闭并手动创建。 |
| 查看日志 | tail -f logs/rocketmqlogs/broker.log 和 namesrv.log | 启动失败时通过日志排查端口冲突、内存不足等问题。 |
| 集群部署 | 多台机器分别部署 NameServer 和 Broker,Broker 配置主从关系 | 主从 Broker 需配置不同的 brokerId(0为主,>0为从),并配置同一 brokerName。 |
2.2 安装 RocketMQ 控制台(Console)
| 操作步骤 | 说明 | 注意事项 |
|---|---|---|
| 下载控制台项目 | git clone https://github.com/apache/rocketmq-externals | 控制台已迁移至 rocketmq-externals 仓库。 |
| 编译打包 | mvn clean package -Dmaven.test.skip=true | 需提前安装 Maven 和 JDK 8+。 |
| 修改配置文件 | 修改 target/classes/application.properties 中的 rocketmq.config.namesrvAddr=127.0.0.1:9876 | 若 NameServer 不在本地或为集群,需填写多个地址,如:192.168.0.1:9876;192.168.0.2:9876 |
| 启动控制台 | java -jar rocketmq-console-ng-1.0.0.jar | 默认端口 8080,可通过 --server.port=8081 修改。 |
| 访问 Web 界面 | 浏览器访问 http://localhost:8080 | 可查看 Topic、Consumer Group、消息堆积、Broker 状态等信息。 |
2.3 创建第一个 Producer 示例
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| DefaultMQProducer | new DefaultMQProducer(String producerGroup) | 创建生产者实例 | DefaultMQProducer producer = new DefaultMQProducer("ProducerGroupTest"); | producerGroup 必须唯一标识一个生产者组,建议命名规范。 |
| setNamesrvAddr | producer.setNamesrvAddr("127.0.0.1:9876") | 设置 NameServer 地址 | producer.setNamesrvAddr("127.0.0.1:9876"); | 必须设置,否则无法连接 Broker。 |
| start | producer.start() | 启动生产者 | producer.start(); | 必须在发送消息前调用,建议放在初始化代码块中。 |
| send | SendResult send(Message msg) | 同步发送消息,阻塞直到收到响应 | SendResult result = producer.send(msg); | 抛出异常需捕获处理,如 MQClientException、RemotingException 等。 |
| createTopic | producer.createTopic("ClusterTest", "TopicTest", 4) | 手动创建 Topic(可选) | producer.createTopic("", "MyTopic", 8); | 第一个参数通常为空字符串,第三个参数为 Queue 数量。 |
| shutdown | producer.shutdown() | 关闭生产者,释放资源 | producer.shutdown(); | 应用退出前调用,避免资源泄漏。 |
2.4 创建第一个 Consumer 示例
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| DefaultMQPushConsumer | new DefaultMQPushConsumer(String consumerGroup) | 创建消费者实例 | DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroupTest"); | consumerGroup 必须唯一,同一组内消费者共享消费负载。 |
| setNamesrvAddr | consumer.setNamesrvAddr("127.0.0.1:9876") | 设置 NameServer 地址 | consumer.setNamesrvAddr("127.0.0.1:9876"); | 必须设置,否则无法获取路由信息。 |
| setConsumeFromWhere | consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET) | 设置首次启动时的消费位置 | consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); | CONSUME_FROM_FIRST_OFFSET:从第一条开始;CONSUME_FROM_LAST_OFFSET:从最新开始 |
| subscribe | consumer.subscribe(String topic, String subExpression) | 订阅指定 Topic,支持 Tag 过滤 | consumer.subscribe("TopicTest", "TagA"); | 支持 || 分隔多个 Tag,如 "TagA||TagB" |
| registerMessageListener | consumer.registerMessageListener(listener) | 注册消息监听器,处理收到的消息 | consumer.registerMessageListener(new ConsumeMessageConcurrentlyListener() { ... }); | 可实现 ConsumeConcurrentlyListener 或 ConsumeOrderlyListener。 |
| start | consumer.start() | 启动消费者 | consumer.start(); | 必须在注册监听器后调用。 |
| shutdown | consumer.shutdown() | 关闭消费者,释放资源 | consumer.shutdown(); | 应用退出前调用。 |
2.5 消息发送与接收流程解析
| 流程阶段 | 说明 | 注意事项 |
|---|---|---|
| Producer 发送消息 | 1. Producer 启动时从 NameServer 获取 Topic 路由信息 2. 选择一个 Message Queue 3. 向对应 Broker 发送消息 | 若自动创建 Topic 关闭,且 Topic 不存在会抛异常;网络异常会自动重试(默认2次)。 |
| Broker 存储消息 | 1. Broker 接收消息后追加到 CommitLog 文件 2. 异步构建 ConsumeQueue 和 IndexFile | 消息物理存储在 CommitLog,逻辑队列信息在 ConsumeQueue。 |
| Consumer 拉取消息 | 1. Consumer 启动时从 NameServer 获取路由 2. 向 Broker 发起拉取请求 3. Broker 返回消息数据 | 拉取是长轮询机制,默认等待 15s 若无消息再返回空。 |
| 消费确认(ACK) | Consumer 处理完消息后向 Broker 发送 ACK,Broker 更新消费位点(Offset) | 若未及时 ACK,Broker 会重复推送(根据重试策略)。 |
| 消费位点管理 | Offset 存储在 Broker 端(集群模式)或本地(广播模式) | Offset 提交可同步或异步,异步提交性能高但可能丢失。 |
| 消息过滤 | Broker 根据 Consumer 订阅的 Tag 或 SQL 表达式进行过滤 | Tag 过滤在服务端完成,SQL 过滤需开启属性支持。 |
第三章:消息发送机制
3.1 同步发送消息
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| send | SendResult send(Message msg) | 同步发送消息,阻塞等待结果返回 | SendResult result = producer.send(msg); | 适用于对可靠性要求高、需确认发送结果的场景,如订单创建。 |
| send | SendResult send(Message msg, long timeout) | 同步发送并设置超时时间 | SendResult result = producer.send(msg, 3000); // 3秒超时 | 超时会抛出 RemotingTimeoutException,建议设置合理超时避免线程阻塞过久。 |
| getSendStatus | result.getSendStatus() | 获取发送状态 | if (result.getSendStatus() == SendStatus.SEND_OK) { ... } | SEND_OK 表示成功,其他如 FLUSH_DISK_TIMEOUT、FLUSH_SLAVE_TIMEOUT 为部分成功或失败。 |
| getMsgId | result.getMsgId() | 获取消息唯一ID | System.out.println("MsgId: " + result.getMsgId()); | 消息ID由Broker生成,可用于消息轨迹追踪。 |
| getQueueOffset | result.getQueueOffset() | 获取消息在队列中的偏移量 | System.out.println("Offset: " + result.getQueueOffset()); | 用于定位消息位置,调试和监控使用。 |
3.2 异步发送消息
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| send | void send(Message msg, SendCallback callback) | 异步发送消息,回调通知结果 | producer.send(msg, new SendCallback() { ... }); | 适用于对响应时间敏感、不关心即时结果的场景,如日志上报。 |
| send | void send(Message msg, SendCallback callback, long timeout) | 异步发送并设置超时 | producer.send(msg, callback, 5000); // 5秒超时 | 超时后触发 onException,需在回调中处理超时逻辑。 |
| onSuccess | public void onSuccess(SendResult result) | 发送成功回调方法 | 在 SendCallback 实现中重写该方法 | 回调执行在 Netty 线程池中,避免执行耗时操作,可提交到业务线程池处理。 |
| onException | public void onException(Throwable e) | 发送失败回调方法 | if (e instanceof MQClientException) { ... } | 常见异常包括 MQClientException(配置错误)、RemotingException(网络问题)等。 |
| setRetryTimesWhenSendFailed | producer.setRetryTimesWhenSendFailed(2) | 设置同步/异步发送失败重试次数 | producer.setRetryTimesWhenSendFailed(3); | 默认2次,异步发送仅在内部线程重试,不影响主线程。 |
3.3 单向发送消息(Oneway)
| 方法名称 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| sendOneway | void sendOneway(Message msg) | 单向发送,不等待任何响应 | producer.sendOneway(msg); | 适用于性能要求极高、允许少量消息丢失的场景,如监控数据、心跳包。 |
| sendOneway | void sendOneway(Message msg, MessageQueue mq) | 指定队列单向发送 | producer.sendOneway(msg, queue); | 用于精确控制消息路由,结合选择器使用。 |
| 性能特点 | - | 极低延迟、高吞吐 | 单向发送可达到百万级TPS | 无法得知发送是否成功,无重试机制,需业务容忍丢失。 |
| 使用建议 | - | 仅用于非关键消息 | 不适用于订单、支付等关键业务 | 结合日志记录或采样监控提升可观测性。 |
3.4 发送结果与异常处理
| 类型 | 名称/方法 | 说明 | 注意事项 |
|---|---|---|---|
| SendResult | getSendStatus() | 返回 SendStatus 枚举值:SEND_OK、FLUSH_DISK_TIMEOUT、FLUSH_SLAVE_TIMEOUT、SLAVE_NOT_AVAILABLE | FLUSH_DISK_TIMEOUT 表示主节点刷盘超时,可能丢失;FLUSH_SLAVE_TIMEOUT 表示从节点同步超时,影响高可用。 |
| SendStatus | SEND_OK | 消息发送成功,已写入主节点 | 最理想状态 |
| FLUSH_DISK_TIMEOUT | 主节点未在规定时间内完成磁盘刷盘 | 若为同步刷盘模式,可能丢失消息;建议生产环境使用同步双写+同步刷盘。 | |
| FLUSH_SLAVE_TIMEOUT | 从节点未在规定时间内完成同步 | 主从数据不一致,故障切换时可能丢失数据。 | |
| SLAVE_NOT_AVAILABLE | 无可用从节点 | 主从架构异常,需检查从节点状态。 | |
| 异常类 | MQClientException | 客户端配置错误、找不到Topic路由等 | 检查 producerGroup、namesrvAddr、topic 是否正确。 |
| RemotingException | 网络通信异常 | 检查网络连接、Broker是否正常运行。 | |
| MQBrokerException | Broker处理失败,如权限不足、磁盘满 | 查看Broker日志定位具体原因。 | |
| InterruptedException | 线程中断异常 | 多在线程池或异步回调中发生,需妥善处理中断状态。 | |
| 处理建议 | - | 捕获异常并记录日志,关键业务需重试或降级 | 建议封装统一的消息发送模板,处理重试、告警、日志等共性逻辑。 |
第四章:消息消费机制
4.1 集群模式消费(Clustering)
| 方法/属性 | 语法/说明 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| setMessageModel | consumer.setMessageModel(MessageModel.CLUSTERING) | 设置消费模式为集群模式 | consumer.setMessageModel(MessageModel.CLUSTERING); | 默认模式,同一Consumer Group内多个实例共享消费队列。 |
| 消费分配规则 | 一个Message Queue只能被Group内的一个Consumer消费 | 实现负载均衡 | 例如:4个Queue,2个Consumer,每个Consumer消费2个Queue | 避免重复消费,提高并发处理能力。 |
| 故障转移 | 若某个Consumer宕机,其负责的Queue会自动分配给其他Consumer | 高可用保障 | NameServer检测心跳超时后触发重平衡 | 重平衡期间可能出现短暂消息堆积。 |
| ACK机制 | Consumer处理完消息后自动或手动提交ACK | 确认消费成功,更新Offset | 默认自动提交,可通过 setConsumeMessageBatchMaxSize 控制批量提交大小 | 若未提交ACK,Broker会在一定时间后重新投递(默认重试16次)。 |
| 适用场景 | - | 适用于大多数业务系统 | 如订单处理、用户注册异步通知 | 要求消费逻辑幂等,防止重试导致重复处理。 |
4.2 广播模式消费(Broadcasting)
| 方法/属性 | 语法/说明 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| setMessageModel | consumer.setMessageModel(MessageModel.BROADCASTING) | 设置消费模式为广播模式 | consumer.setMessageModel(MessageModel.BROADCASTING); | 必须显式设置,否则默认为集群模式。 |
| 消费行为 | 同一Consumer Group内的每个Consumer都会收到全量消息 | 实现消息广播 | 所有实例独立消费所有Queue中的消息 | 不进行负载均衡,每条消息被每个实例处理一次。 |
| Offset管理 | Offset存储在本地文件系统(~/.rocketmq_offsets) | 每个Consumer独立维护消费进度 | 路径可配置,重启后从上次位置继续消费 | 不依赖Broker,适合本地状态管理。 |
| 无ACK机制 | 无需向Broker发送ACK | 简化消费流程 | 消息处理失败不会重试 | 消费失败需自行记录日志或告警,无法通过重试弥补。 |
| 适用场景 | - | 配置更新、缓存同步、本地状态刷新 | 如:推送新的风控规则到所有节点 | 不适用于高吞吐或需可靠消费的场景。 |
4.3 负载均衡与消息分配策略
| 策略名称 | 类名 | 说明 | 注意事项 |
|---|---|---|---|
| 平均分配策略 | AllocateMessageQueueAveragely | 将 Queue 平均分配给 Consumers,余数部分由前几个 Consumer 多承担一个 | 最常用策略,分配均匀,适合 Consumer 数量稳定场景。 |
| 环形平均分配策略 | AllocateMessageQueueCircular | 按 Consumer 和 Queue 的环形顺序依次分配 | 分配结果与启动顺序有关,适合动态扩缩容。 |
| 一致性哈希分配策略 | AllocateMessageQueueConsistentHash | 基于一致性哈希算法分配,减少 Consumer 增减时的重平衡影响 | 适合大规模 Consumer 集群,减少数据迁移。 |
| 同机房优先分配策略 | AllocateMessageQueueByMachineRoom | 优先将 Queue 分配给同机房的 Consumer | 需配合机房信息配置(messageModel、broker、consumer 的 machineRoom 指定) |
| 定点分配策略 | AllocateMessageQueueByConfig | 手动指定每个 Consumer 负责的 Queue 列表 | 适用于特殊业务需求,需自行管理分配逻辑。 |
| 随机分配策略 | AllocateMessageQueueRandom | 随机分配 Queue 给 Consumer | 分配不均,一般不推荐用于生产环境。 |
| 设置分配策略 | consumer.setAllocateMessageQueueStrategy(new AllocateMessageQueueAveragely()) | 设置自定义分配策略 | 必须在 start() 之前调用 |
4.4 消费位点(Offset)管理
| 方法/属性 | 语法/说明 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| setConsumeFromWhere | consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET) | 设置首次消费起始位置 | CONSUME_FROM_FIRST_OFFSET:从第一条开始CONSUME_FROM_LAST_OFFSET:从最新开始CONSUME_FROM_TIMESTAMP:从指定时间开始 | 生产环境通常设为 CONSUME_FROM_LAST_OFFSET 避免历史消息积压。 |
| persistConsumerOffset | consumer.persistConsumerOffset() | 手动持久化当前消费位点 | 一般无需手动调用,自动提交机制已封装 | 用于特殊场景下强制保存Offset。 |
| 持久化时机 | - | 自动提交:每隔5秒 手动提交:调用 ack | 集群模式下Offset提交到Broker 广播模式下保存到本地文件 | 自动提交可能丢失少量消息(最多5秒数据),关键业务建议手动提交。 |
| 手动提交Offset | consumer.updateConsumeOffsetToBroker(MessageQueue, offset) | 强制更新Broker上的Offset | 用于跳过消息或回溯消费 | 操作需谨慎,错误设置可能导致消息丢失或重复消费。 |
| 重试队列Offset | %RETRY%{consumerGroup} Topic下的Queue Offset | 管理重试消息的消费进度 | 重试间隔随次数增加(10s, 30s, 1min…2h) | 重试16次后进入死信队列(%DLQ%)。 |
| 死信队列 | %DLQ%{consumerGroup} | 存储最终消费失败的消息 | 需单独启动消费者监听该Topic处理异常消息 | 死信消息需人工介入分析或补偿。 |
第五章:消息类型详解
5.1 普通消息(Normal Message)
| 特性 | 说明 | 注意事项 |
|---|---|---|
| 定义 | 最基本的消息类型,无特殊顺序、延迟或事务要求。 | 适用于大多数异步通信场景,如日志采集、通知推送。 |
| 发送方式 | 支持同步、异步、Oneway 三种发送模式。 | 根据业务对可靠性与性能的需求选择合适模式。 |
| 存储机制 | 消息写入 CommitLog,构建 ConsumeQueue 和 IndexFile 提供快速检索。 | 物理存储统一,逻辑队列分离,提升并发能力。 |
| 可靠性 | 支持持久化、主从复制、重试机制,保障不丢失。 | 建议开启同步刷盘 + 同步双写以实现最高可靠性(牺牲性能)。 |
| 使用示例 | Message msg = new Message("TopicTest", "TagA", "Hello RocketMQ".getBytes()); | Tag 可用于简单过滤,提高消费端匹配效率。 |
5.2 顺序消息(Ordered Message)
| 方法/概念 | 说明 | 注意事项 |
|---|---|---|
| 顺序保障机制 | 将同一业务维度的消息发送到同一个 Message Queue,消费者单线程消费该 Queue | 仅保证局部有序(Queue 内 FIFO),无法保证全局有序。 |
| 消息队列选择器 | 使用 MessageQueueSelector 自定义消息路由逻辑 | 示例:按订单 ID 哈希选择队列,确保同一订单的消息进入同一队列。 |
| 代码示例 | producer.send(msg, new MessageQueueSelector() { ... }, orderId, callback); | arg 通常为业务键(如 orderId),用于一致性哈希。 |
| 消费端实现 | 实现 MessageListenerOrderly 接口 | 该监听器内部加锁,保证单队列内消息串行消费。 |
| 阻塞风险 | 若某条消息处理失败或阻塞,后续消息将被”堵塞” | 建议设置合理的消费超时和重试策略,避免长时间卡顿。 |
| 适用场景 | 订单状态流转、数据库变更日志(binlog)、交易流水处理 | 不适用于高并发无序场景,可能成为性能瓶颈。 |
5.3 延时消息(Delayed Message)
| 属性/方法 | 说明 | 注意事项 |
|---|---|---|
| 延时等级 | 内置 18 个等级:1s, 5s, 10s, 30s, 1m, 2m, 3m, 4m, 5m, 6m, 7m, 8m, 9m, 10m, 20m, 30m, 1h, 2h | 不支持任意时间精度,需映射到最接近的等级。 |
| 设置延时等级 | msg.setDelayTimeLevel(3); // 10秒后投递 | 等级从 1 开始,0 表示不延时。 |
| 存储机制 | 延时消息先写入 SCHEDULE_TOPIC_XXXX 主题下的特定队列,由定时任务扫描后投递到目标 Topic | 延时消息不影响普通消息性能,但大量延时消息可能影响定时线程。 |
| 应用场景 | - 订单超时取消(30分钟) - 支付结果回调重试(10秒、30秒) - 活动开始提醒(提前1小时) | 避免使用延时消息实现高频调度任务,建议结合外部调度系统(如 XXL-JOB)。 |
| 局限性 | 不支持精确到毫秒的任意时间点调度 | 如需更高精度,可使用事务消息 + 外部定时器实现。 |
5.4 批量消息(Batch Message)
| 方法/属性 | 说明 | 注意事项 |
|---|---|---|
| send | producer.send(Collection<Message> msgs) | 将多条消息打包一次性发送,减少网络请求次数,提升吞吐量。 |
| 消息限制 | - 单批次总大小 ≤ 1MB(可配置) - 所有消息必须属于同一 Topic - 不能包含事务消息或延时消息 | 超过大小会抛 MQClientException。 |
| 分批发送 | 若消息集合过大,需手动分片发送 | 示例:每 50 条或 512KB 为一批 |
| 优点 | 减少网络开销,提升发送效率,适合日志、监控等高频小消息场景。 | 批量发送失败则整批重试,需权衡可靠性与性能。 |
| 缺点 | 若一批中某条消息失败,整个批次需重试 | 建议对重要消息单独发送,非关键消息可批量。 |
5.5 事务消息(Transactional Message)
| 阶段 | 说明 | 注意事项 |
|---|---|---|
| 第一阶段:发送半消息 | Producer 发送一条”半消息”(Half Message),Broker 存储但不投递给 Consumer | 使用 sendMessageInTransaction 方法发送。 |
| 第二阶段:执行本地事务 | Producer 执行本地数据库操作,并返回事务状态(COMMIT、ROLLBACK、UNKNOWN) | 必须实现 LocalTransactionExecuter 或 TransactionListener。 |
| 第三阶段:提交或回滚 | Producer 向 Broker 提交事务状态,Broker 决定是否投递消息 | 提交后消息变为”可消费”状态。 |
| 回查机制 | 若 Broker 未收到状态,会定时回调 Producer 查询事务状态(checkLocalTransaction) | 必须实现回查逻辑,防止”悬挂事务”。 |
| 使用限制 | - 不支持批量事务消息 - 仅支持 Producer 发送 - 延时等级必须为 0 | 事务消息最大回查次数默认 15 次,之后进入死信队列。 |
| 适用场景 | 跨系统数据一致性,如:账户扣款成功后发送扣款通知、库存扣减后生成订单 | 避免在事务中执行耗时操作,防止阻塞消息队列。 |
第六章:Producer 高级特性
6.1 生产者组与命名规范
| 规范项 | 建议命名方式 | 说明 | 注意事项 |
|---|---|---|---|
| 生产者组命名 | {业务域}_{环境}_{功能} | 如:order_prod_payment(订单生产环境支付通知)user_test_register(用户测试环境注册) | 统一命名便于运维管理、监控告警和权限控制。 |
| 命名字符限制 | 只能包含字母、数字、下划线、短横线、点号 | 长度建议 ≤ 64 字符 | 避免使用特殊字符或空格。 |
| 唯一性要求 | 同一集群内 Producer Group 名称必须唯一 | 防止路由混乱和资源冲突 | 多个 JVM 实例使用同一 Group 表示逻辑上同一个生产者集群。 |
| 与 Topic 关系 | 一个 Producer 可发送多个 Topic 的消息 | 但建议按业务边界划分 Group,避免耦合 | 例如:订单系统使用 order_producer,用户系统使用 user_producer。 |
6.2 消息重试机制
| 参数名称 | 默认值 | 说明 | 注意事项 |
|---|---|---|---|
| retryTimesWhenSendFailed | 2 | 同步发送失败时的重试次数(不包括首次) | 网络抖动时自动重试,提升可靠性。 |
| retryTimesWhenSendAsyncFailed | 2 | 异步发送失败时的重试次数 | 重试在内部线程执行,不影响业务线程。 |
| retryAnotherBrokerWhenNotStoreOK | false | 若消息未存储成功(如磁盘满),是否尝试发送到其他 Broker | 开启后可提升可用性,但可能导致消息乱序。 |
| 重试间隔 | - | 无固定间隔,立即重试 | 重试期间仍可能失败,需结合监控告警。 |
| 重试策略控制 | - | 可通过实现 MessageQueueSelector 自定义重试路由 | 例如:避开已知故障的 Broker。 |
| 重试与幂等 | - | 重试可能导致消息重复,消费端必须实现幂等处理 | 建议在消费端使用唯一键(如 messageId 或业务 ID)去重。 |
6.3 消息过滤(Tag 与 SQL 表达式)
| 过滤方式 | 语法示例 | 说明 | 注意事项 |
|---|---|---|---|
| Tag 过滤 | consumer.subscribe("TopicA", "TagA||TagB") | 简单字符串匹配,多个 Tag 用 || 分隔 | 可在 Consumer 端订阅时指定,性能高。 |
| SQL 表达式过滤 | consumer.subscribe("TopicB", MessageSelector.bySql("a > 5 AND b = 'abc'")) | 需消息携带属性(msg.putUserProperty("a", "10")),在 Broker 端执行 SQL 判断 | 需开启 enablePropertyFilter=true,性能低于 Tag 过滤。 |
| 属性设置 | msg.putUserProperty("city", "beijing") | 用户自定义属性,用于 SQL 过滤 | 属性值为字符串类型,建议避免过大或敏感信息。 |
| 性能对比 | Tag 过滤 > SQL 过滤 | Tag 是简单字符串匹配,SQL 需解析执行 | 高吞吐场景优先使用 Tag;复杂条件使用 SQL。 |
| 使用建议 | - 简单分类用 Tag(如:order_created, order_paid) - 多维度筛选用 SQL(如:region=‘north’ AND level>3) | 合理设计消息模型,避免过度依赖过滤降低性能。 |
6.4 消息轨迹(Message Trace)
| 组件/方法 | 说明 | 注意事项 |
|---|---|---|
| 开启轨迹 | producer.setSendMsgTrace(true) | 需部署 rocketmq-trace 服务或集成 OpenTelemetry/SkyWalking |
| 轨迹数据内容 | 包含消息发送、存储、消费各阶段的时间戳、IP、Broker、耗时等 | 用于定位延迟瓶颈、分析消息生命周期 |
| 存储方式 | 轨迹消息以特殊 Topic(RMQ_SYS_TRACE_TOPIC)存储,可被单独消费 | 可对接 ELK、Prometheus 等监控系统。 |
| 查询方式 | 通过消息 ID 在控制台或 API 查询完整轨迹 | 需保留足够日志时间(如7天) |
| 性能影响 | 开启后增加约 10%~20% 的网络和存储开销 | 生产环境建议按需开启,或采样开启(如 10% 流量)。 |
| 分布式链路追踪集成 | 支持与 Jaeger、Zipkin 等标准协议对接 | 实现全链路监控,定位跨系统调用问题。 |
6.5 生产者性能调优参数
| 参数名称 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
| sendMsgTimeout | 3000(ms) | 发送消息超时时间 | 高延迟网络可适当增大(如5000),低延迟场景可减小。 |
| compressMsgBodyOverHowmuch | 4096(4KB) | 消息体超过该大小时启用压缩(LZ4) | 大消息场景建议设为 1024(1KB)以节省带宽。 |
| maxMessageSize | 4 * 1024 * 1024(4MB) | 单条消息最大大小 | 可调大至 16MB,但需确保 Broker 配置一致。 |
| flushDiskType | ASYNC_FLUSH | 刷盘方式:ASYNC(异步)或 SYNC_FLUSH(同步) | 可靠性优先设为 SYNC_FLUSH;性能优先保留 ASYNC。 |
| retryTimesWhenSendFailed | 2 | 同步发送失败重试次数 | 弱网络环境可增至 3,避免频繁失败。 |
| useTLS | false | 是否启用 TLS 加密 | 公网传输建议开启,提升安全性。 |
| pollNameServerInterval | 30000(30s) | 拉取 NameServer 路由信息间隔 | 网络稳定可设为 60s 减少请求;频繁变更可设为 10s。 |
| heartbeatBrokerInterval | 30000(30s) | 向 Broker 发送心跳间隔 | 一般无需调整。 |
| clientPooled | true | 是否启用 Netty 连接池 | 高并发场景建议开启,复用连接提升性能。 |
第七章:Consumer 高级特性
7.1 消费者组与订阅关系
| 概念/方法 | 说明 | 注意事项 |
|---|---|---|
| 消费者组(Consumer Group) | 逻辑上的消费者集合,同一组内共享消费队列(集群模式)或广播接收(广播模式) | 组名必须全局唯一,命名建议:{业务}_{环境}_{功能},如 order_prod_processor |
| 订阅关系一致性 | 同一 Consumer Group 内所有实例必须订阅相同的 Topic 和 Tag(或 SQL 表达式) | 若订阅不一致,可能导致消息分配混乱或消费失败,Broker 会打印警告日志 |
| 订阅表达式设置 | consumer.subscribe("TopicA", "TagA||TagB") 或 consumer.subscribe("TopicB", MessageSelector.bySql("age > 18")) | Tag 支持 || 分隔多个 Tag |
| 动态订阅变更 | 不支持运行时动态修改订阅表达式(需重启 Consumer) | 建议在上线前确定订阅规则,避免中途变更导致不可预期行为 |
| 订阅关系存储位置 | 集群模式下,订阅关系由 Consumer 上报至 Broker 维护 | Broker 根据订阅关系进行消息过滤和投递 |
| 多 Topic 订阅 | 单个 Consumer 实例可订阅多个 Topic | 需确保资源(线程、内存)充足,避免相互影响 |
7.2 消费重试与死信队列
| 机制/参数 | 说明 | 注意事项 |
|---|---|---|
| 重试次数 | 默认最大重试 16 次,间隔递增(10s, 30s, 1min, 2min…2h) | 重试间隔不可配置,由系统固定算法决定 |
| 重试 Topic | 重试消息发送至 %RETRY%{consumerGroup} 主题 | 该 Topic 由系统自动创建,无需手动管理 |
| 重试触发条件 | 消费返回 ConsumeConcurrentlyStatus.RECONSUME_LATER 或抛出异常 | 返回 CONSUME_SUCCESS 不会重试 |
| 死信队列(DLQ) | 重试 16 次仍失败的消息进入 %DLQ%{consumerGroup} | 需单独启动 DLQ 消费者进行人工干预或补偿处理 |
| DLQ 消费建议 | 监听 %DLQ%{consumerGroup},记录日志、告警或进入人工审核流程 | 不建议自动重试 DLQ 消息,防止死循环 |
| 跳过重试 | 可通过业务逻辑判断,直接返回成功或记录后跳过 | 对于已知无法处理的”脏数据”,可主动跳过避免阻塞 |
7.3 消费幂等性设计
| 设计方案 | 实现方式 | 适用场景 | 注意事项 |
|---|---|---|---|
| 唯一键去重(推荐) | 使用 msg.getMsgId() 或业务唯一键(如订单ID)作为幂等键,存储于 Redis/DB | 所有场景通用,可靠性高 | 建议设置过期时间(如7天),避免存储无限增长 |
| 数据库唯一索引 | 消费逻辑写入 DB 时,利用唯一约束防止重复插入 | 订单创建、用户注册等写操作 | 适用于写操作明确的场景,异常需捕获 DuplicateKeyException |
| 状态机控制 | 消费前检查业务状态,仅在特定状态执行操作 | 订单状态流转(待支付 → 已支付) | 需保证状态变更原子性,防止并发问题 |
| Token 机制 | 生产者生成一次性 Token,消费者消费后标记已使用 | 高并发场景,防止重复提交 | 需配合缓存实现,Token 生命周期管理复杂 |
| 幂等框架封装 | 自定义注解 + AOP 拦截,自动处理幂等逻辑 | 多个消费端统一治理 | 可提升开发效率,但需考虑通用性和性能开销 |
⚠️ 重要提示: RocketMQ 不保证消息不重复,消费端必须自行实现幂等。
7.4 流量控制与限流策略
| 策略 | 配置方式 | 说明 | 注意事项 |
|---|---|---|---|
| 并发消费线程数 | consumer.setConsumeThreadMin(20)consumer.setConsumeThreadMax(64) | 控制消费线程池大小,间接限流 | 线程过多可能导致 GC 频繁,建议根据 CPU 核数合理设置 |
| 批量消费大小 | consumer.setConsumeMessageBatchMaxSize(1) | 每次最多消费 N 条消息(默认1) | 设置 >1 可提升吞吐,但需处理批量失败回滚 |
| 拉取间隔控制 | consumer.setPullInterval(1000) | 每次拉取后等待 N 毫秒再拉取 | 适用于低频消费场景,避免频繁拉取 |
| 拉取数量限制 | consumer.setPullBatchSize(32) | 每次最多拉取 32 条消息(默认32) | 减少单次拉取量可降低内存压力 |
| 外部限流组件 | 集成 Sentinel、Hystrix | 实现 QPS 控制、熔断降级 | 适用于复杂流量治理场景,需引入额外依赖 |
| 动态调整策略 | 结合监控指标(堆积量、RT)动态调整线程数或拉取参数 | 实现自适应限流 | 需开发定制化逻辑,提升系统弹性 |
7.5 消费者性能调优参数
| 参数名称 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
| consumeThreadMin / consumeThreadMax | 20 / 64 | 消费线程池最小/最大线程数 | 高吞吐场景可增至 128;低配机器可降低至 10/20 |
| pullBatchSize | 32 | 每次拉取消息数量 | 可调大至 64 或 128 提升吞吐,但增加单次处理压力 |
| consumeMessageBatchMaxSize | 1 | 每次提交给监听器的最大消息数 | 批量消费可设为 10~50,需确保监听器支持批量处理 |
| pullInterval | 0 | 拉取间隔(毫秒) | 为 0 表示连续拉取;设为 1000 可降低 CPU 使用率 |
| subscriptionMode | SHARED | 订阅模式(SHARED、EXLUSIVE、SIMPLE) | 默认 SHARED,支持集群消费 |
| messageModel | CLUSTERING | 消费模式 | 根据业务选择 CLUSTERING 或 BROADCASTING |
| rebalanceInterval | 20000(20s) | 负载均衡检查间隔 | 一般无需调整 |
| maxReconsumeTimes | 16 | 最大重试次数 | 可根据业务需求调整(1~16) |
第八章:Broker 与 NameServer 配置
8.1 Broker 核心配置参数详解
| 参数 | 默认值 | 说明 | 建议值 |
|---|---|---|---|
| brokerName | - | Broker 名称,主从部署时需相同 | 如 broker-a |
| brokerId | 0 | 0 表示 Master,>0 表示 Slave | Slave 设置为 1、2 等 |
| listenPort | 10911 | 客户端通信端口 | 保持默认 |
| namesrvAddr | - | NameServer 地址列表 | 192.168.0.1:9876;192.168.0.2:9876 |
| storePathRootDir | $HOME/store | 存储根目录 | 建议挂载高性能 SSD |
| storePathCommitLog | $storePathRootDir/commitlog | CommitLog 存储路径 | 确保磁盘空间充足 |
| mapedFileSizeCommitLog | 1G | 单个 CommitLog 文件大小 | 保持默认 |
| flushDiskType | ASYNC_FLUSH | 刷盘方式:ASYNC / SYNC_FLUSH | 生产环境建议 SYNC_FLUSH |
| brokerRole | ASYNC_MASTER | 角色:ASYNC_MASTER、SYNC_MASTER、SLAVE | 主从部署时 Master 设为 SYNC_MASTER |
| enableDLegerCommitLog | false | 是否启用 Dledger(Raft)模式 | 高可用部署建议开启,替代主从 |
| autoCreateTopicEnable | true | 是否自动创建 Topic | 建议设为 false,由运维统一管理 |
| clusterPermission | 2 | 集群权限:1-只读,2-读写 | 保持默认 2 |
8.2 NameServer 高可用部署
| 部署方式 | 说明 | 注意事项 |
|---|---|---|
| 多节点独立部署 | 部署 2~3 个 NameServer 节点,无主从之分 | 节点间不通信,数据最终一致 |
| 客户端配置 | Producer/Consumer 配置所有 NameServer 地址(分号分隔) | namesrvAddr=ns1:9876;ns2:9876;ns3:9876 |
| 故障容忍 | 客户端内置路由缓存,NameServer 全部宕机仍可发送/消费一段时间 | 建议部署至少 2 个节点防止单点故障 |
| 监控与告警 | 监控 NameServer 进程、端口、JVM 状态 | 无状态服务,重启不影响已有连接 |
| 不支持动态扩缩容 | 新增 NameServer 节点需更新所有客户端配置 | 建议初期规划好节点数 |
| 推荐部署数量 | 生产环境建议 3 个节点 | 提供更高可用性 |
8.3 消息存储机制(CommitLog、ConsumeQueue、IndexFile)
| 组件 | 说明 | 特点 |
|---|---|---|
| CommitLog | 消息物理存储文件,所有 Topic 消息顺序写入 | - 顺序 I/O,高性能写入 - 单文件 1G,文件名是起始偏移量 - mmap 技术提升读写效率 |
| ConsumeQueue | 逻辑队列,存储消息在 CommitLog 的物理偏移量、大小、Tag Hash | - 每个 Topic-Queue 一个文件 - 固定条目大小(20B),随机读取高效 - 是消息投递的基础 |
| IndexFile | 倒排索引文件,支持按 key 或时间范围查询消息 | - 文件名包含起始时间戳 - 包含 Hash 槽和 Index 链表 - 用于消息轨迹和排查 |
| 存储流程 | 1. 消息写入 CommitLog → 2. 异步构建 ConsumeQueue → 3. 异步构建 IndexFile | 三者解耦,写入性能不受索引影响 |
| 文件清理 | 默认保留 72 小时,超过时间的文件被自动删除 | 可通过 fileReservedTime 参数调整 |
8.4 刷盘策略与主从同步
| 策略 | 配置参数 | 说明 | 优缺点 |
|---|---|---|---|
| 异步刷盘(ASYNC_FLUSH) | flushDiskType=ASYNC_FLUSH | 写入 PageCache 后立即返回,后台线程定时刷盘 | ✅ 性能高 ❌ 断电可能丢失数据 |
| 同步刷盘(SYNC_FLUSH) | flushDiskType=SYNC_FLUSH | 必须等待数据写入磁盘后才返回 | ✅ 数据安全 ❌ 性能较低(依赖磁盘速度) |
| 异步复制(ASYNC_MASTER) | brokerRole=ASYNC_MASTER | Master 写成功即返回,Slave 异步拉取同步 | ✅ 性能高 ❌ 主宕机可能丢数据 |
| 同步双写(SYNC_MASTER) | brokerRole=SYNC_MASTER | Master 必须等待 Slave 写成功才返回 | ✅ 高可用,数据不丢 ❌ 延迟高,依赖网络 |
| Dledger 模式(推荐) | enableDLegerCommitLog=true | 基于 Raft 协议的多副本一致性 | ✅ 自动选主,高可用 ✅ 数据强一致 ✅ 支持动态扩缩容 |
✅ 生产环境推荐组合: SYNC_FLUSH + SYNC_MASTER 或直接使用 Dledger 模式。
第九章:消息存储与可靠性
9.1 消息持久化机制
| 机制/组件 | 工作原理 | 特点 | 注意事项 |
|---|---|---|---|
| CommitLog 顺序写入 | 所有消息按到达顺序写入 CommitLog 文件,使用内存映射(mmap)技术 | - 顺序 I/O,写入性能高 - 文件大小固定(1GB),命名以起始偏移量命名(如 00000000000000000000) | 写入失败时返回 SEND_FAILED,生产者需重试 |
| PageCache 缓存 | 消息先写入操作系统 PageCache,由内核异步刷盘 | - 减少直接磁盘 I/O,提升吞吐量 | 建议使用 UPS 电源保障 - 断电可能导致未刷盘数据丢失 |
| 同步刷盘(SYNC_FLUSH) | 调用 fsync() 确保数据落盘后才返回成功 | - 数据可靠性高,宕机不丢 - 延迟增加,依赖磁盘性能 | 适用于金融、交易等高可靠场景 |
| 异步刷盘(ASYNC_FLUSH) | 写入 PageCache 后立即返回,后台线程定时刷盘 | - 性能高,延迟低 - 断电可能丢失秒级数据 | 适用于日志、通知等容忍少量丢失场景 |
| 文件映射(mmap) | 使用 MappedByteBuffer 将文件映射到内存 | - 零拷贝技术,减少用户态与内核态切换 - 支持大文件高效读写 | 注意 JVM 内存设置,避免 OOM |
| 存储路径管理 | 可配置 storePathRootDir 分离 CommitLog、ConsumeQueue、IndexFile | - 提升 I/O 隔离性 - 便于磁盘容量规划 | 建议 CommitLog 使用 SSD,其他可使用普通磁盘 |
9.2 高可用与容灾设计
| 架构模式 | 说明 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| 主从复制(Master-Slave) | Master 接收读写,Slave 异步或同步复制数据,Master 宕机后手动切换 | - 部署简单 - 支持同步双写保障数据不丢 | - 故障转移需人工介入 - Slave 不对外提供服务 | 小型系统,成本敏感 |
| Dledger 模式(Raft) | 基于 Raft 协议的多副本一致性,自动选主,支持动态扩缩容 | - 自动故障转移(秒级) - 数据强一致 - 支持读写分离 | - 配置稍复杂 - 需要奇数个节点(3/5) | 生产环境推荐,高可用要求高 |
| 多 NameServer 部署 | 部署 2~3 个 NameServer 节点,客户端配置多个地址 | - 无单点故障 - 节点间无通信,部署简单 | - 不支持动态感知新节点(需重启客户端) | 所有部署环境必备 |
| 跨机房容灾(同城双活) | 在两个机房部署独立集群,通过消息同步工具(如 MirrorMaker)复制 | - 单机房故障不影响服务 - 支持流量切换 | - 成本高 - 数据同步延迟存在 | 金融、电商等核心业务 |
| 云上高可用部署 | 使用云厂商提供的 RocketMQ 服务(如阿里云ONS) | - 免运维 - 自动扩缩容 - SLA 高达 99.99% | - 成本较高 - 受限于云平台 | 企业级生产环境首选 |
9.3 消息重复与丢失场景分析
| 问题类型 | 场景描述 | 原因分析 | 解决方案 |
|---|---|---|---|
| 消息丢失 | 消息未到达 Broker 或未持久化 | - Producer 发送失败未重试 - Broker 异步刷盘时宕机 - 磁盘损坏 | - 生产者开启同步发送 + 重试 - Broker 使用 SYNC_FLUSH + SYNC_MASTER - 定期备份 CommitLog |
| 消息重复 | 同一条消息被消费多次 | - 网络抖动导致 Producer 重试 - Consumer 消费后未及时提交 Offset - Broker 重启导致 Offset 回滚 | - 消费端实现幂等处理(唯一键、数据库约束等) - 提高提交频率(flushConsumerOffsetInterval) |
| 消息乱序 | 消息未按发送顺序消费 | - 多线程并发消费 - 网络延迟导致消息乱序到达 | - 顺序消息使用 MessageListenerOrderly - 非关键场景可容忍乱序 |
| 消息堆积 | 消费速度远低于生产速度 | - Consumer 性能瓶颈 - 消费逻辑阻塞或异常 - Topic 队列数不足 | - 增加 Consumer 实例 - 优化消费逻辑(异步处理、批处理) - 扩展 Topic 队列数 |
| Offset 丢失 | Consumer 重启后从最早或最晚开始消费 | - 未开启持久化存储 Offset - Group 名称变更 | - 使用集群模式(CLUSTERING) - 避免随意修改 Group 名 |
9.4 可靠性保障最佳实践
| 实践项 | 推荐方案 | 说明 |
|---|---|---|
| 生产者可靠性 | - 使用同步发送(send) - 设置 retryTimesWhenSendFailed=2 - 开启 sendMsgTrace 跟踪 | 确保消息成功发送并可追溯 |
| Broker 持久化 | - flushDiskType=SYNC_FLUSH - brokerRole=SYNC_MASTER - 使用 SSD 存储 CommitLog | 最大程度防止消息丢失 |
| 消费者幂等性 | - 使用 msgId 或业务 ID 去重(Redis/DB) - 数据库唯一索引约束 | 必须实现,防止重复消费导致数据错乱 |
| 主从高可用 | 优先使用 Dledger 模式(3 节点) | 自动选主,无需人工干预,适合生产环境 |
| 监控告警 | - 监控 TPS、延迟、堆积量 - 设置堆积阈值告警(如 > 10000) - 日志关键字监控(ERROR、WARN) | 及时发现异常,快速响应 |
| Topic 与 Group 管理 | - 禁用 autoCreateTopicEnable - 统一命名规范 - 权限控制(ACL) | 防止滥用,提升可维护性 |
| 定期演练 | - 模拟 Broker 宕机 - 网络隔离测试 - 消费者故障恢复 | 验证高可用机制有效性,提升应急能力 |
第十章:监控与运维
10.1 使用 RocketMQ Console 监控
| 功能模块 | 说明 | 使用建议 |
|---|---|---|
| Dashboard | 展示集群概览:TPS、连接数、消息堆积、存储使用率 | 设置大屏监控,实时掌握系统状态 |
| Topic 管理 | 查看 Topic 列表、创建/删除 Topic、设置队列数 | 生产环境建议由运维统一管理 |
| Consumer 管理 | 查看消费者组、订阅关系、消费进度(Offset)、堆积量 | 重点关注堆积增长趋势 |
| Producer 管理 | 查看生产者列表、发送 TPS、连接状态 | 排查发送异常节点 |
| 消息查询 | 按 Message ID、Key 或时间范围搜索消息 | 用于问题排查、数据核对 |
| Broker 状态 | 查看 Broker CPU、内存、磁盘、网络 I/O | 结合系统监控工具(如 Zabbix)使用 |
| 报警配置 | 设置堆积、延迟、异常等告警规则 | 告警通知到钉钉、邮件、企业微信 |
10.2 关键监控指标(TPS、延迟、堆积)
| 指标 | 监控方式 | 告警阈值 | 说明 |
|---|---|---|---|
| 生产 TPS | broker_tps 或 producer_tps | 持续 > 10000 或突增 200% | 反映消息生产压力 |
| 消费 TPS | consumer_tps | 明显低于生产 TPS | 可能出现堆积 |
| 消息延迟 | 消息存储时间与当前时间差 | 平均 > 1s 或 P99 > 5s | 表示系统处理缓慢 |
| 消息堆积量 | consumer_offset_lag | > 10000 或持续增长 | 消费能力不足,需扩容 |
| Broker CPU/Memory | 系统指标采集 | CPU > 80%,内存 > 70% | 资源瓶颈,影响性能 |
| 磁盘使用率 | df -h 或监控平台 | > 80% | 可能导致写入失败,需清理或扩容 |
| 网络 I/O | iftop 或监控工具 | 出口带宽 > 80% | 影响消息传输效率 |
10.3 日志分析与问题排查
| 日志类型 | 路径 | 常见关键字 | 排查方法 |
|---|---|---|---|
| Broker 日志 | logs/rocketmqlogs/broker.log | ERROR, WARN, Store error, Disk full | 搜索错误信息,定位存储、网络问题 |
| NameServer 日志 | logs/rocketmqlogs/namesrv.log | Register broker, Route info | 查看 Broker 注册状态、路由变化 |
| Producer 日志 | 应用日志 | SendResult, MQException, Timeout | 检查发送失败原因(网络、超时) |
| Consumer 日志 | 应用日志 | ConsumeMessage, RECONSUME, Exception | 分析消费失败、重试、异常堆栈 |
| GC 日志 | JVM 参数开启 | Full GC, Pause time | 判断是否因 GC 导致消费延迟 |
| CommitLog 异常 | storeerror.log | MappedFile, Flush error | 重点关注文件写入失败 |
| 排查流程 | 1. 看监控 → 2. 查日志 → 3. 分析堆栈 → 4. 模拟复现 → 5. 修复验证 | 建立标准化排障流程 |
10.4 常见运维命令与工具
| 命令/工具 | 语法示例 | 功能说明 |
|---|---|---|
| mqadmin | sh mqadmin clusterList -n 192.168.0.1:9876 | 查看集群节点、Topic、消费者状态 |
sh mqadmin topicList -n 192.168.0.1:9876 | 列出所有 Topic | |
sh mqadmin consumerProgress -n 192.168.0.1:9876 -g group_name | 查看消费者堆积情况 | |
sh mqadmin updateTopic -n 192.168.0.1:9876 -t TopicA -r 16 | 修改 Topic 队列数为 16 | |
| Broker 控制脚本 | sh mqbroker -c broker.conf & | 启动 Broker |
kill $(jps | grep broker | awk '{print $1}') | 停止 Broker(生产环境慎用) | |
| NameServer 脚本 | sh mqnamesrv & | 启动 NameServer |
kill $(jps | grep namesrv | awk '{print $1}') | 停止 NameServer | |
| 消息查询工具 | sh mqadmin queryMsgById -i <msgId> | 根据 Message ID 查询消息内容 |
sh mqadmin queryMsgByKey -t TopicA -k orderId123 | 根据 Key 查询消息 | |
| 状态检查 | jstat -gc <pid> 1s | 查看 JVM GC 状态 |
top -p <pid> | 查看进程 CPU、内存占用 | |
df -h | 检查磁盘空间 | |
| 自动化脚本 | 自定义 Shell/Python 脚本 | 实现定时备份、健康检查、告警推送 |
第十一章:Spring 集成
11.1 Spring Boot 整合 RocketMQ
| 步骤 | 操作说明 | 示例代码 |
|---|---|---|
| 1. 添加依赖 | 引入 rocketmq-spring-boot-starter | xml<br/><dependency><br/> <groupId>org.apache.rocketmq</groupId><br/> <artifactId>rocketmq-spring-boot-starter</artifactId><br/> <version>2.2.3</version><br/></dependency> |
| 2. 配置文件 | 在 application.yml 中配置 NameServer 地址 | yaml<br/>rocketmq:<br/> name-server: 192.168.0.1:9876;192.168.0.2:9876<br/> producer:<br/> group: order_producer_group |
| 3. 发送消息 | 使用 RocketMQTemplate 发送 | java<br/>@Autowired<br/>private RocketMQTemplate rocketMQTemplate;<br/><br/>public void sendOrderMsg(String msg) {<br/> rocketMQTemplate.convertAndSend("order-topic", msg);<br/>} |
| 4. 启动类注解 | 可选:@EnableRocketMQ(新版本可省略) | java<br/>@SpringBootApplication<br/>@EnableRocketMQ<br/>public class App { ... } |
| 5. 测试验证 | 编写单元测试或接口调用 | 确保消息成功发送并被消费 |
| 注意事项 | - 确保版本兼容(Spring Boot 与 RocketMQ Starter) - 生产者组名必须配置 | 建议使用 2.2.x 稳定版本 |
11.2 使用 @RocketMQTransactionListener 实现事务消息
| 要素 | 说明 | 示例代码 |
|---|---|---|
| 事务监听器 | 实现本地事务逻辑和回查机制 | java<br/>@Component<br/>@RocketMQTransactionListener(txProducerGroup = "trade_tx_group")<br/>public class TradeTransactionListener implements RocketMQLocalTransactionListener {<br/><br/> @Override<br/> public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {<br/> try {<br/> // 执行扣款等本地事务<br/> boolean result = accountService.deduct((String)arg);<br/> return result ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK;<br/> } catch (Exception e) {<br/> return RocketMQLocalTransactionState.UNKNOWN; // 触发回查<br/> }<br/> }<br/><br/> @Override<br/> public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {<br/> String orderId = new String(msg.getBody());<br/> boolean isCommitted = orderService.isOrderValid(orderId);<br/> return isCommitted ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK;<br/> }<br/>} |
| 发送事务消息 | 使用 RocketMQTemplate 发送半消息 | java<br/>TransactionSendResult sendResult = rocketMQTemplate.sendMessageInTransaction(<br/> "trade_tx_topic",<br/> MessageBuilder.withPayload("order123").build(),<br/> "order123" // 业务参数,传给 executeLocalTransaction<br/>); |
| 关键注解 | @RocketMQTransactionListener(txProducerGroup = "xxx") | 必须指定与生产者一致的事务组 |
| 回查机制 | 当返回 UNKNOWN 时,Broker 定时回调 checkLocalTransaction | 确保该方法幂等且能正确查询本地事务状态 |
| 限制 | - 不支持批量事务消息 - 仅支持集群消费模式 | 事务消息延迟略高于普通消息 |
11.3 注解方式实现消息监听
| 注解 | 说明 | 示例代码 |
|---|---|---|
| @RocketMQMessageListener | 核心注解,定义消费者 | java<br/>@Component<br/>@RocketMQMessageListener(<br/> consumerGroup = "order_consumer_group",<br/> topic = "order-topic",<br/> selectorExpression = "CREATE || PAY", // Tag 过滤<br/> consumeMode = ConsumeMode.CONCURRENTLY, // 并发消费<br/> messageModel = MessageModel.CLUSTERING // 集群模式<br/>)<br/>public class OrderMessageListener implements RocketMQListener<String> {<br/><br/> @Override<br/> public void onMessage(String message) {<br/> System.out.println("收到订单消息: " + message);<br/> // 处理业务逻辑<br/> }<br/>} |
| 并发消费(默认) | consumeMode = ConsumeMode.CONCURRENTLY | 多线程并发处理,性能高 |
| 顺序消费 | consumeMode = ConsumeMode.ORDERLY | 单线程顺序处理,保证 FIFO |
| 广播模式 | messageModel = MessageModel.BROADCASTING | 每条消息被所有消费者实例消费 |
| SQL 过滤 | selectorType = SelectorType.SQL92, selectorExpression = “age > 18” | 需开启 enablePropertyFilter=true |
| 监听泛型 | 可监听 String、byte[] 或自定义对象(需序列化) | 建议使用 JSON 字符串传递对象 |
11.4 配置文件参数说明
| 配置项 | 默认值 | 说明 | 示例 |
|---|---|---|---|
| rocketmq.name-server | - | NameServer 地址列表 | 192.168.0.1:9876;192.168.0.2:9876 |
| rocketmq.producer.group | - | 生产者组名 | order_producer |
| rocketmq.producer.send-message-timeout | 3000ms | 发送超时时间 | 5000 |
| rocketmq.producer.retry-times-when-send-failed | 2 | 同步发送失败重试次数 | 3 |
| rocketmq.producer.retry-next-server | false | 失败后切换 Broker | true |
| rocketmq.producer.max-message-size | 4MB | 最大消息大小 | 8388608 (8MB) |
| rocketmq.consumer.group | - | 消费者组名(在 Listener 中配置更常见) | order_consumer |
| rocketmq.consumer.pull-batch-size | 32 | 每次拉取消息数 | 64 |
| rocketmq.consumer.consume-thread-min | 20 | 消费线程最小数 | 30 |
| rocketmq.consumer.consume-thread-max | 64 | 消费线程最大数 | 128 |
| rocketmq.consumer.access-channel | LOCAL | 访问通道:LOCAL(内网)或 CLOUD(云) | CLOUD(阿里云环境) |
| rocketmq.file-size-threshold | 1024 | 超过该大小启用压缩(单位:KB) | 1024(1MB) |
⚠️ 注意: 部分参数在
@RocketMQMessageListener注解中优先级更高。
第十二章:实战案例与最佳实践
12.1 订单系统解耦案例
| 组件 | 角色 | 实现方式 | 优势 |
|---|---|---|---|
| 订单服务 | Producer | 创建订单后发送 order-created 消息到 order-topic | 解耦核心流程,提升响应速度 |
| 库存服务 | Consumer | 订阅 order-topic,监听 CREATE Tag,扣减库存 | 库存不足可返回重试或通知用户 |
| 支付服务 | Consumer | 监听 order-paid 消息,更新订单状态 | 支付结果异步通知,避免阻塞 |
| 通知服务 | Consumer | 监听订单状态变更,发送短信/邮件 | 用户体验提升,系统独立演进 |
| 技术要点 | - 使用 Tag 区分事件类型(CREATE, PAID, CANCEL) - 消费端幂等处理(订单ID去重) - 消息轨迹追踪 | 实现松耦合、高可用的订单中心 |
12.2 秒杀系统削峰填谷
| 阶段 | 策略 | 实现 |
|---|---|---|
| 前端限流 | 按钮置灰、验证码、排队页面 | 减少无效请求到达后端 |
| 网关层过滤 | 令牌桶/漏桶算法 | 限制总请求数,防止系统崩溃 |
| 消息队列削峰 | 用户请求写入 RocketMQ 而非直接扣库存 | java<br/>seckillService.handleRequest(userId, itemId);<br/>rocketMQTemplate.convertAndSend("seckill-queue", requestJson); |
| 异步处理 | 消费者从队列中逐个处理秒杀请求 | - 验证资格 - 扣减 Redis 库存 - 写入 DB - 发送结果 |
| 结果通知 | 处理完成后通过 WebSocket 或短信通知用户 | 用户无需长时间等待 |
| 优势 | 将瞬时高并发转化为平稳的后台处理,保护数据库 | 提升系统吞吐量和稳定性 |
12.3 日志收集与异步处理
| 架构 | 组件说明 | 优势 |
|---|---|---|
| 日志生成端 | 应用通过 RocketMQTemplate 发送日志消息 | 替代直接写文件或同步调用 ELK |
| Topic 设计 | 按日志类型分 Topic:log-access, log-error, log-trace | 便于分类消费和存储 |
| 异步消费 | 消费者将日志写入 Elasticsearch 或 HDFS | - 解耦应用与日志系统 - 提升写入性能 |
| 批量处理 | 消费者设置 consumeMessageBatchMaxSize=100 批量入库 | 减少数据库连接压力 |
| 高可靠保障 | - 生产者同步发送 - Broker 同步刷盘 - 消费者幂等写入 | 防止日志丢失 |
| 扩展性 | 可增加消费者实现日志分析、告警、归档等 | 一套日志,多用途消费 |
12.4 分布式事务最终一致性方案
| 场景 | 流程 | 技术实现 |
|---|---|---|
| 账户扣款 → 发送通知 | 1. 扣款服务开启事务消息 2. 执行本地扣款 3. 发送”扣款成功”半消息 4. 提交事务(COMMIT/ROLLBACK) 5. 通知服务消费消息并发送短信 | - 使用 @RocketMQTransactionListener - 通知服务幂等处理 |
| 库存扣减 → 创建订单 | 1. 订单服务发送事务消息 2. 执行创建订单(本地事务) 3. 消息提交后,库存服务扣减库存 | - 订单与库存服务最终一致 - 失败可补偿或人工干预 |
| 优势 | - 避免强一致性带来的性能瓶颈 - 系统间解耦 - 可靠消息保障不丢失 | 适用于对实时一致性要求不高的场景 |
| 补偿机制 | 对长期未处理的消息启动定时任务补偿 | 结合死信队列(DLQ)人工处理 |
12.5 性能压测与调优建议
| 项目 | 工具/方法 | 调优建议 |
|---|---|---|
| 压测工具 | JMeter、wrk、自定义 Producer 压测脚本 | 模拟真实业务消息大小和频率 |
| 监控指标 | - TPS(生产/消费) - 消息延迟(P99) - 堆积量 - Broker CPU/Memory/Disk | 使用 RocketMQ Console 或 Prometheus + Grafana |
| 生产者调优 | - 增大 sendMsgTimeout - 开启压缩(>4KB) - 批量发送(send(Collection)) | 批量大小 ≤ 1MB,避免超时 |
| 消费者调优 | - 增加 consumeThreadMax - 增大 pullBatchSize - 批量消费(consumeMessageBatchMaxSize) | 线程数建议为 CPU 核数的 2~4 倍 |
| Broker 调优 | - 使用 SSD 存储 CommitLog - 调整 mapedFileSizeCommitLog=1G - flushDiskType=ASYNC_FLUSH(测试环境) | 生产环境建议 SYNC_FLUSH |
| 系统层调优 | - JVM 参数:-Xms4g -Xmx4g -XX:+UseG1GC - Linux 内核:增大文件句柄、网络缓冲区 | 避免 Full GC 频繁 |
| 容量规划 | 根据日均消息量 × 平均大小 × 保留时间计算磁盘需求 | 建议预留 30% 冗余空间 |