Article

消息队列 RocketMQ

更新于:2026-07-14

第一章: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 QueueTopic 下的分区,每个 Topic 可分为多个 Message Queue,用于并行生产和消费。Queue 数量影响并发度,建议根据消费能力合理设置。
Consumer Group消费者组,同一组内的消费者共同消费一个 Topic 的消息,实现负载均衡。集群模式下,一个 Queue 只能被组内一个消费者消费;广播模式下,所有消费者都消费全量消息。

1.4 应用场景与优势对比(Kafka、RabbitMQ)

对比维度RocketMQKafkaRabbitMQ说明
开发语言JavaScalaErlang与 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解压路径避免包含空格或中文字符。
启动 NameServernohup sh bin/mqnamesrv &默认监听 9876 端口,可通过 -p 参数指定配置文件。
启动 Brokernohup sh bin/mqbroker -n localhost:9876 autoCreateTopicEnable=true &autoCreateTopicEnable=true 允许自动创建 Topic,生产环境建议关闭并手动创建。
查看日志tail -f logs/rocketmqlogs/broker.lognamesrv.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 示例

方法名称语法用途代码示例注意事项
DefaultMQProducernew DefaultMQProducer(String producerGroup)创建生产者实例DefaultMQProducer producer = new DefaultMQProducer("ProducerGroupTest");producerGroup 必须唯一标识一个生产者组,建议命名规范。
setNamesrvAddrproducer.setNamesrvAddr("127.0.0.1:9876")设置 NameServer 地址producer.setNamesrvAddr("127.0.0.1:9876");必须设置,否则无法连接 Broker。
startproducer.start()启动生产者producer.start();必须在发送消息前调用,建议放在初始化代码块中。
sendSendResult send(Message msg)同步发送消息,阻塞直到收到响应SendResult result = producer.send(msg);抛出异常需捕获处理,如 MQClientException、RemotingException 等。
createTopicproducer.createTopic("ClusterTest", "TopicTest", 4)手动创建 Topic(可选)producer.createTopic("", "MyTopic", 8);第一个参数通常为空字符串,第三个参数为 Queue 数量。
shutdownproducer.shutdown()关闭生产者,释放资源producer.shutdown();应用退出前调用,避免资源泄漏。

2.4 创建第一个 Consumer 示例

方法名称语法用途代码示例注意事项
DefaultMQPushConsumernew DefaultMQPushConsumer(String consumerGroup)创建消费者实例DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroupTest");consumerGroup 必须唯一,同一组内消费者共享消费负载。
setNamesrvAddrconsumer.setNamesrvAddr("127.0.0.1:9876")设置 NameServer 地址consumer.setNamesrvAddr("127.0.0.1:9876");必须设置,否则无法获取路由信息。
setConsumeFromWhereconsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET)设置首次启动时的消费位置consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);CONSUME_FROM_FIRST_OFFSET:从第一条开始;CONSUME_FROM_LAST_OFFSET:从最新开始
subscribeconsumer.subscribe(String topic, String subExpression)订阅指定 Topic,支持 Tag 过滤consumer.subscribe("TopicTest", "TagA");支持 || 分隔多个 Tag,如 "TagA||TagB"
registerMessageListenerconsumer.registerMessageListener(listener)注册消息监听器,处理收到的消息consumer.registerMessageListener(new ConsumeMessageConcurrentlyListener() { ... });可实现 ConsumeConcurrentlyListener 或 ConsumeOrderlyListener。
startconsumer.start()启动消费者consumer.start();必须在注册监听器后调用。
shutdownconsumer.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 同步发送消息

方法名称语法用途代码示例注意事项
sendSendResult send(Message msg)同步发送消息,阻塞等待结果返回SendResult result = producer.send(msg);适用于对可靠性要求高、需确认发送结果的场景,如订单创建。
sendSendResult send(Message msg, long timeout)同步发送并设置超时时间SendResult result = producer.send(msg, 3000); // 3秒超时超时会抛出 RemotingTimeoutException,建议设置合理超时避免线程阻塞过久。
getSendStatusresult.getSendStatus()获取发送状态if (result.getSendStatus() == SendStatus.SEND_OK) { ... }SEND_OK 表示成功,其他如 FLUSH_DISK_TIMEOUT、FLUSH_SLAVE_TIMEOUT 为部分成功或失败。
getMsgIdresult.getMsgId()获取消息唯一IDSystem.out.println("MsgId: " + result.getMsgId());消息ID由Broker生成,可用于消息轨迹追踪。
getQueueOffsetresult.getQueueOffset()获取消息在队列中的偏移量System.out.println("Offset: " + result.getQueueOffset());用于定位消息位置,调试和监控使用。

3.2 异步发送消息

方法名称语法用途代码示例注意事项
sendvoid send(Message msg, SendCallback callback)异步发送消息,回调通知结果producer.send(msg, new SendCallback() { ... });适用于对响应时间敏感、不关心即时结果的场景,如日志上报。
sendvoid send(Message msg, SendCallback callback, long timeout)异步发送并设置超时producer.send(msg, callback, 5000); // 5秒超时超时后触发 onException,需在回调中处理超时逻辑。
onSuccesspublic void onSuccess(SendResult result)发送成功回调方法在 SendCallback 实现中重写该方法回调执行在 Netty 线程池中,避免执行耗时操作,可提交到业务线程池处理。
onExceptionpublic void onException(Throwable e)发送失败回调方法if (e instanceof MQClientException) { ... }常见异常包括 MQClientException(配置错误)、RemotingException(网络问题)等。
setRetryTimesWhenSendFailedproducer.setRetryTimesWhenSendFailed(2)设置同步/异步发送失败重试次数producer.setRetryTimesWhenSendFailed(3);默认2次,异步发送仅在内部线程重试,不影响主线程。

3.3 单向发送消息(Oneway)

方法名称语法用途代码示例注意事项
sendOnewayvoid sendOneway(Message msg)单向发送,不等待任何响应producer.sendOneway(msg);适用于性能要求极高、允许少量消息丢失的场景,如监控数据、心跳包。
sendOnewayvoid sendOneway(Message msg, MessageQueue mq)指定队列单向发送producer.sendOneway(msg, queue);用于精确控制消息路由,结合选择器使用。
性能特点-极低延迟、高吞吐单向发送可达到百万级TPS无法得知发送是否成功,无重试机制,需业务容忍丢失。
使用建议-仅用于非关键消息不适用于订单、支付等关键业务结合日志记录或采样监控提升可观测性。

3.4 发送结果与异常处理

类型名称/方法说明注意事项
SendResultgetSendStatus()返回 SendStatus 枚举值:SEND_OK、FLUSH_DISK_TIMEOUT、FLUSH_SLAVE_TIMEOUT、SLAVE_NOT_AVAILABLEFLUSH_DISK_TIMEOUT 表示主节点刷盘超时,可能丢失;FLUSH_SLAVE_TIMEOUT 表示从节点同步超时,影响高可用。
SendStatusSEND_OK消息发送成功,已写入主节点最理想状态
FLUSH_DISK_TIMEOUT主节点未在规定时间内完成磁盘刷盘若为同步刷盘模式,可能丢失消息;建议生产环境使用同步双写+同步刷盘。
FLUSH_SLAVE_TIMEOUT从节点未在规定时间内完成同步主从数据不一致,故障切换时可能丢失数据。
SLAVE_NOT_AVAILABLE无可用从节点主从架构异常,需检查从节点状态。
异常类MQClientException客户端配置错误、找不到Topic路由等检查 producerGroup、namesrvAddr、topic 是否正确。
RemotingException网络通信异常检查网络连接、Broker是否正常运行。
MQBrokerExceptionBroker处理失败,如权限不足、磁盘满查看Broker日志定位具体原因。
InterruptedException线程中断异常多在线程池或异步回调中发生,需妥善处理中断状态。
处理建议-捕获异常并记录日志,关键业务需重试或降级建议封装统一的消息发送模板,处理重试、告警、日志等共性逻辑。

第四章:消息消费机制

4.1 集群模式消费(Clustering)

方法/属性语法/说明用途代码示例注意事项
setMessageModelconsumer.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)

方法/属性语法/说明用途代码示例注意事项
setMessageModelconsumer.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)管理

方法/属性语法/说明用途代码示例注意事项
setConsumeFromWhereconsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET)设置首次消费起始位置CONSUME_FROM_FIRST_OFFSET:从第一条开始
CONSUME_FROM_LAST_OFFSET:从最新开始
CONSUME_FROM_TIMESTAMP:从指定时间开始
生产环境通常设为 CONSUME_FROM_LAST_OFFSET 避免历史消息积压。
persistConsumerOffsetconsumer.persistConsumerOffset()手动持久化当前消费位点一般无需手动调用,自动提交机制已封装用于特殊场景下强制保存Offset。
持久化时机-自动提交:每隔5秒
手动提交:调用 ack
集群模式下Offset提交到Broker
广播模式下保存到本地文件
自动提交可能丢失少量消息(最多5秒数据),关键业务建议手动提交。
手动提交Offsetconsumer.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)

方法/属性说明注意事项
sendproducer.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 消息重试机制

参数名称默认值说明注意事项
retryTimesWhenSendFailed2同步发送失败时的重试次数(不包括首次)网络抖动时自动重试,提升可靠性。
retryTimesWhenSendAsyncFailed2异步发送失败时的重试次数重试在内部线程执行,不影响业务线程。
retryAnotherBrokerWhenNotStoreOKfalse若消息未存储成功(如磁盘满),是否尝试发送到其他 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 生产者性能调优参数

参数名称默认值说明调优建议
sendMsgTimeout3000(ms)发送消息超时时间高延迟网络可适当增大(如5000),低延迟场景可减小。
compressMsgBodyOverHowmuch4096(4KB)消息体超过该大小时启用压缩(LZ4)大消息场景建议设为 1024(1KB)以节省带宽。
maxMessageSize4 * 1024 * 1024(4MB)单条消息最大大小可调大至 16MB,但需确保 Broker 配置一致。
flushDiskTypeASYNC_FLUSH刷盘方式:ASYNC(异步)或 SYNC_FLUSH(同步)可靠性优先设为 SYNC_FLUSH;性能优先保留 ASYNC。
retryTimesWhenSendFailed2同步发送失败重试次数弱网络环境可增至 3,避免频繁失败。
useTLSfalse是否启用 TLS 加密公网传输建议开启,提升安全性。
pollNameServerInterval30000(30s)拉取 NameServer 路由信息间隔网络稳定可设为 60s 减少请求;频繁变更可设为 10s。
heartbeatBrokerInterval30000(30s)向 Broker 发送心跳间隔一般无需调整。
clientPooledtrue是否启用 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 / consumeThreadMax20 / 64消费线程池最小/最大线程数高吞吐场景可增至 128;低配机器可降低至 10/20
pullBatchSize32每次拉取消息数量可调大至 64 或 128 提升吞吐,但增加单次处理压力
consumeMessageBatchMaxSize1每次提交给监听器的最大消息数批量消费可设为 10~50,需确保监听器支持批量处理
pullInterval0拉取间隔(毫秒)为 0 表示连续拉取;设为 1000 可降低 CPU 使用率
subscriptionModeSHARED订阅模式(SHARED、EXLUSIVE、SIMPLE)默认 SHARED,支持集群消费
messageModelCLUSTERING消费模式根据业务选择 CLUSTERING 或 BROADCASTING
rebalanceInterval20000(20s)负载均衡检查间隔一般无需调整
maxReconsumeTimes16最大重试次数可根据业务需求调整(1~16)

第八章:Broker 与 NameServer 配置

8.1 Broker 核心配置参数详解

参数默认值说明建议值
brokerName-Broker 名称,主从部署时需相同如 broker-a
brokerId00 表示 Master,>0 表示 SlaveSlave 设置为 1、2 等
listenPort10911客户端通信端口保持默认
namesrvAddr-NameServer 地址列表192.168.0.1:9876;192.168.0.2:9876
storePathRootDir$HOME/store存储根目录建议挂载高性能 SSD
storePathCommitLog$storePathRootDir/commitlogCommitLog 存储路径确保磁盘空间充足
mapedFileSizeCommitLog1G单个 CommitLog 文件大小保持默认
flushDiskTypeASYNC_FLUSH刷盘方式:ASYNC / SYNC_FLUSH生产环境建议 SYNC_FLUSH
brokerRoleASYNC_MASTER角色:ASYNC_MASTER、SYNC_MASTER、SLAVE主从部署时 Master 设为 SYNC_MASTER
enableDLegerCommitLogfalse是否启用 Dledger(Raft)模式高可用部署建议开启,替代主从
autoCreateTopicEnabletrue是否自动创建 Topic建议设为 false,由运维统一管理
clusterPermission2集群权限: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_MASTERMaster 写成功即返回,Slave 异步拉取同步✅ 性能高
❌ 主宕机可能丢数据
同步双写(SYNC_MASTER)brokerRole=SYNC_MASTERMaster 必须等待 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、延迟、堆积)

指标监控方式告警阈值说明
生产 TPSbroker_tps 或 producer_tps持续 > 10000 或突增 200%反映消息生产压力
消费 TPSconsumer_tps明显低于生产 TPS可能出现堆积
消息延迟消息存储时间与当前时间差平均 > 1s 或 P99 > 5s表示系统处理缓慢
消息堆积量consumer_offset_lag> 10000 或持续增长消费能力不足,需扩容
Broker CPU/Memory系统指标采集CPU > 80%,内存 > 70%资源瓶颈,影响性能
磁盘使用率df -h 或监控平台> 80%可能导致写入失败,需清理或扩容
网络 I/Oiftop 或监控工具出口带宽 > 80%影响消息传输效率

10.3 日志分析与问题排查

日志类型路径常见关键字排查方法
Broker 日志logs/rocketmqlogs/broker.logERROR, WARN, Store error, Disk full搜索错误信息,定位存储、网络问题
NameServer 日志logs/rocketmqlogs/namesrv.logRegister broker, Route info查看 Broker 注册状态、路由变化
Producer 日志应用日志SendResult, MQException, Timeout检查发送失败原因(网络、超时)
Consumer 日志应用日志ConsumeMessage, RECONSUME, Exception分析消费失败、重试、异常堆栈
GC 日志JVM 参数开启Full GC, Pause time判断是否因 GC 导致消费延迟
CommitLog 异常storeerror.logMappedFile, Flush error重点关注文件写入失败
排查流程1. 看监控 → 2. 查日志 → 3. 分析堆栈 → 4. 模拟复现 → 5. 修复验证建立标准化排障流程

10.4 常见运维命令与工具

命令/工具语法示例功能说明
mqadminsh 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-starterxml<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-timeout3000ms发送超时时间5000
rocketmq.producer.retry-times-when-send-failed2同步发送失败重试次数3
rocketmq.producer.retry-next-serverfalse失败后切换 Brokertrue
rocketmq.producer.max-message-size4MB最大消息大小8388608 (8MB)
rocketmq.consumer.group-消费者组名(在 Listener 中配置更常见)order_consumer
rocketmq.consumer.pull-batch-size32每次拉取消息数64
rocketmq.consumer.consume-thread-min20消费线程最小数30
rocketmq.consumer.consume-thread-max64消费线程最大数128
rocketmq.consumer.access-channelLOCAL访问通道:LOCAL(内网)或 CLOUD(云)CLOUD(阿里云环境)
rocketmq.file-size-threshold1024超过该大小启用压缩(单位: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% 冗余空间