Java 消息中间件 - RocketMQ 全解(2026 详细版)
下面把 RocketMQ 的全流程拆成 12 大阶段,每段都给出:
-
为什么要做 → 怎么做 → 命令/代码 → 常见报错 → 急救方案 跟着敲,零基础也能搭出生产级集群。
阶段1:学前必知——RocketMQ 到底解决啥?
|
关键词 |
一句话 |
|---|---|
|
高吞吐 |
单broker 10w+ QPS,阿里双11核心链路 |
|
低延迟 |
写入 <2ms,适合交易冲正、订单关单 |
|
金融级事务 |
半消息+回查,本地DB与消息原子性 |
|
万亿级堆积 |
文件顺序写+PageCache,磁盘即队列 |
如果你需要 “订单支付后必须发消息,且两者不能丢” → 选 RocketMQ。
阶段2:安装前准备——3 分钟配好环境
1.系统要求
-
JDK 8+(推荐 11)
-
Linux kernel ≥ 3.10(CentOS 7/8、Ubuntu 18+)
-
磁盘 ≥ 100G SSD(顺序写吃磁盘)
-
关闭防火墙 or 开放 9876⁄10911
2.创建专用用户
useradd rocketmq echo "rocketmq ALL=(ALL) NOPASSWD:ALL" >> /etc/sudoers su - rocketmq
3.调大系统参数(否则启动失败)
echo "* soft nofile 65536" >> /etc/security/limits.conf echo "* hard nofile 65536" >> /etc/security/limits.conf
阶段3:一键安装——官方二进制 3 命令
# 1. 下载(国内镜像快)
wget https://mirrors.tuna.tsinghua.edu.cn/apache/rocketmq/5.1.4/rocketmq-all-5.1.4-bin-release.zip
unzip rocketmq-all-5.1.4-bin-release.zip && cd rocketmq-all-5.1.4-bin-release
# 2. 调小 JVM(开发机 8G 内存必做,否则 OOM)
sed -i 's/-Xms4g -Xmx4g/-Xms256m -Xmx256m/g' bin/runserver.sh bin/runbroker.sh
# 3. 启动
nohup sh bin/mqnamesrv > ~/logs/ns.log 2>&1 &
sleep 3
nohup sh bin/mqbroker -n localhost:9876 -c conf/broker.conf > ~/logs/broker.log 2>&1 &
# 4. 验证
export NAMESRV_ADDR=localhost:9876
sh bin/tools.sh org.apache.rocketmq.example.quickstart.Producer
sh bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer
看到 SendResult=[SEND_OK] 说明成功;如果 Connect to <127.0.0.1:9876> failed → 检查 NameServer 是否启动。
阶段4:配置生产级参数——外网IP+同步刷盘
编辑 conf/broker.conf 末尾追加:
# 必须改成你机器外网IP,否则Docker/NAT客户端连不上
brokerIP1=192.168.1.101
# 金融场景防丢消息
flushDiskType=SYNC_FLUSH
# 主从异步,平衡性能
brokerRole=ASYNC_MASTER
# 堆外内存缓冲,降延迟
transientStorePoolEnable=true
# 关闭自动建Topic,生产必关
autoCreateTopicEnable=false
重启 broker 生效:
sh bin/mqshutdown broker
nohup sh bin/mqbroker -n localhost:9876 -c conf/broker.conf &
阶段5:创建 Topic & ConsumerGroup(CLI)
export NAMESRV_ADDR=localhost:9876
# 创建Topic 8队列 1副本
sh bin/mqadmin updateTopic -n localhost:9876 -c DefaultCluster -t order-topic -r 8 -w 8
# 创建消费者组
sh bin/mqadmin updateSubGroup -n localhost:9876 -c DefaultCluster -g order-consumer-group
如果提示 Topic already exist 忽略即可。
阶段6:Spring Boot 3 整合(事务消息模板)
依赖:
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.3.0</version>
</dependency>
yml:
rocketmq:
name-server: localhost:9876
producer:
group: my-producer-group
send-message-timeout: 3000
代码(事务消息 = 订单+消息原子性):
@RestController
class OrderController {
@Autowired private RocketMQTemplate rocketMQTemplate;
@PostMapping("/create")
public String createOrder(@RequestBody OrderDTO dto) {
Message<String> msg = MessageBuilder
.withPayload(JSON.toJSONString(dto))
.setHeader("orderId", dto.getOrderId())
.build();
// 发送事务消息
rocketMQTemplate.sendMessageInTransaction("order-topic", msg, dto.getOrderId());
return "sent";
}
}
// 本地事务监听器
@RocketMQTransactionListener
class OrderTxListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
orderService.createOrder((Long) arg); // DB操作
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 回查事务状态(防止应用重启)
return orderService.exist((Long) msg.getHeaders().get("orderId"))
? RocketMQLocalTransactionState.COMMIT
: RocketMQLocalTransactionState.ROLLBACK;
}
}
启动后调用 /create 出现 SEND_OK 即成功。
阶段7:消费者(集群+广播双模式)
集群模式(默认,负载均衡):
@Component
@RocketMQMessageConsumer(topic = "order-topic", consumerGroup = "order-consumer-group")
public class OrderConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
OrderDTO dto = JSON.parseObject(message, OrderDTO.class);
// 幂等判断
if (redis.exists("order:" + dto.getOrderId())) return;
handleOrder(dto);
redis.setEx("order:" + dto.getOrderId(), 600);
}
}
广播模式(所有实例都收):
@RocketMQMessageConsumer(
topic = "config-topic",
consumerGroup = "config-consumer-group",
messageModel = MessageModel.BROADCASTING)
public class ConfigConsumer implements RocketMQListener<String> { ... }
阶段8:高阶功能(顺序/延迟/批量)
-
顺序消息(同一订单ID进入同一队列) 发送时指定 HashKey:
rocketMQTemplate.syncSendOrderly("pay-topic", msg, orderId); -
延迟消息(18级 1s~2h)
Message<String> msg = MessageBuilder.withPayload(body).build(); msg.setDelayTimeLevel(3); // 3=10s rocketMQTemplate.syncSend("delay-topic", msg); -
批量发送(提高吞吐)
List<Message<String>> list = msgs.stream() .map(m -> MessageBuilder.withPayload(m).build()) .collect(Collectors.toList()); rocketMQTemplate.syncSend("batch-topic", list);
阶段9:监控 & 控制台(Docker 一键)
docker run -d --name rocketmq-dashboard \
-p 8080:8080 \
-e "JAVA_OPTS=-Drocketmq.namesrv.addr=192.168.1.101:9876" \
apacherocketmq/rocketmq-dashboard:latest
浏览器打开 http://localhost:8080 → 实时查看 Topic、ConsumerGroup、TPS、ConsumerLag。 Lag > 5 万 就要扩容或加消费者线程。
阶段10:性能压测(官方工具)
# 生产 100字节消息 200W条 200线程
sh bin/benchmark.sh -t prod -n 2000000 -s 100 -w 200 -c DefaultCluster -k order-topic
# 消费
sh bin/benchmark.sh -t cons -n 2000000 -w 200 -g order-consumer-group -k order-topic
关注 AvgRT 和 MaxRT,RT > 50ms 就要调大 sendThreadPoolNum 或检查磁盘。
阶段11:集群部署(2 主 2 从 最小高可用)
|
机器 |
角色 |
|---|---|
|
node1 |
NameServer + Master-1 |
|
node2 |
NameServer + Master-2 |
|
node3 |
Slave-1 (Replicate Master-1) |
|
node4 |
Slave-2 (Replicate Master-2) |
broker.conf 示例(node1 Master-1):
brokerClusterName=DefaultCluster
brokerName=broker-a
brokerId=0 # 0=Master >0=Slave
brokerIP1=192.168.1.101
namesrvAddr=node1:9876;node2:9876
flushDiskType=SYNC_FLUSH
依次启动 4 台 Broker,Dashboard 里看到 双 Master 双 Slave 即成功。
阶段12:备份 & 升级 & 应急
-
消息备份 使用 mqadmin exportMetadata 导出元数据,数据文件直接 rsync /store 目录。
-
平滑升级 先升级 Slave → 切换 Master → 升级旧 Master,0 停机。
-
消息堆积应急
-
临时扩容消费者线程
-
调整 consumeThreadMax=100
-
业务端做批量消费(List<MessageExt> 一次性处理 1000 条)
保姆级总结(口诀)
“装好 JDK 调内存,外网 IP 必配精; 事务消息半提交,回查状态要记清; 顺序延迟批量发,监控 Lag 别停盯; Slave 先升再切换,备份 store 零宕机!”
照抄 12 阶段,从开发到生产,RocketMQ 集群任你玩转!
更多推荐




所有评论(0)