RocketMQ Dledger 集群搭建
·
RocketMQ简介
Apache RocketMQ 是一款开源的分布式消息中间件,最初由阿里巴巴集团自主研发,用于支撑其高并发、高可靠的核心业务,并于 2016 年捐赠给 Apache 软件基金会,2017 年成为 Apache 顶级项目(TLP)。如今,RocketMQ 已广泛应用于金融、电商、物流、物联网、大数据等场景,是企业级消息传递的行业标准之一。
🧩 一、集群规划
| 主机 | IP 地址 | 服务 |
|---|---|---|
| node1 | 192.168.10.101 | NameServer + Broker (RaftNode00) |
| node2 | 192.168.10.102 | NameServer + Broker (RaftNode00) |
| node3 | 192.168.10.103 | NameServer + Broker (RaftNode00) |
✅ 所有节点角色相同,组成一个 Dledger Raft Group
🛠 二、前置准备(所有节点执行)
1. 安装 Java(OpenJDK 8+)
# CentOS
sudo yum install -y java-1.8.0-openjdk
# Ubuntu
sudo apt update && sudo apt install -y openjdk-8-jdk
2. 下载并解压 RocketMQ
cd /opt
wget https://dist.apache.org/repos/dist/release/rocketmq/5.2.0/rocketmq-all-5.2.0-bin-release.zip
unzip rocketmq-all-5.2.0-bin-release.zip
mv rocketmq-all-5.2.0 rocketmq
3. 创建存储目录
mkdir -p /data/rocketmq/store
4. 关闭防火墙(或开放端口)
# 开放关键端口
sudo firewall-cmd --permanent --add-port={9876,30911,40911}/tcp
sudo firewall-cmd --reload
📄 三、配置文件(每个节点略有不同)
🔹 node1 (192.168.10.101) 的 /opt/rocketmq/conf/dledger.conf
brokerClusterName=RaftCluster
brokerName=RaftNode00
listenPort=30911
dLegerSelfId=n0
dLegerPeers=n0-192.168.10.101:40911;n1-192.168.10.102:40911;n2-192.168.10.103:40911
storePathRootDir=/data/rocketmq/store
storePathCommitLog=/data/rocketmq/store/commitlog
namesrvAddr=192.168.10.101:9876;192.168.10.102:9876;192.168.10.103:9876
sendMessageThreadPoolNums=128
useReentrantLockWhenPutMessage=true
🔹 node2 (192.168.10.102) 的 /opt/rocketmq/conf/dledger.conf
properties编辑
brokerClusterName=RaftCluster
brokerName=RaftNode00
listenPort=30911
dLegerSelfId=n1
dLegerPeers=n0-192.168.10.101:40911;n1-192.168.10.102:40911;n2-192.168.10.103:40911
storePathRootDir=/data/rocketmq/store
storePathCommitLog=/data/rocketmq/store/commitlog
namesrvAddr=192.168.10.101:9876;192.168.10.102:9876;192.168.10.103:9876
sendMessageThreadPoolNums=128
useReentrantLockWhenPutMessage=true
🔹 node3 (192.168.10.103) 的 /opt/rocketmq/conf/dledger.conf
brokerClusterName=RaftCluster
brokerName=RaftNode00
listenPort=30911
dLegerSelfId=n2
dLegerPeers=n0-192.168.10.101:40911;n1-192.168.10.102:40911;n2-192.168.10.103:40911
storePathRootDir=/data/rocketmq/store
storePathCommitLog=/data/rocketmq/store/commitlog
namesrvAddr=192.168.10.101:9876;192.168.10.102:9876;192.168.10.103:9876
sendMessageThreadPoolNums=128
useReentrantLockWhenPutMessage=true
⚠️ 关键点:
dLegerSelfId每个节点唯一(n0/n1/n2)dLegerPeers在所有节点完全一致namesrvAddr包含所有 NameServer 地址
▶️ 四、启动集群(所有节点执行)
1. 启动 NameServer
nohup /opt/rocketmq/bin/mqnamesrv > /dev/null 2>&1 &
echo $ ! > /var/run/rocketmq-namesrv.pid
2. 启动 Broker
nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/dledger.conf > /dev/null 2>&1 &
echo $ ! > /var/run/rocketmq-broker.pid
3. 查看日志(确认启动成功)
tail -f ~/logs/rocketmqlogs/broker.log
✅ 成功标志:
DLedger startup succeed, self=n0, role=Leader
The broker[RaftNode00, 192.168.10.101:30911] boot success.
🔍 五、验证集群
1. 查看集群列表
/opt/rocketmq/bin/mqadmin clusterList -n 192.168.10.101:9876
输出应包含 3 个 Broker 实例。
2. 创建测试 Topic
/opt/rocketmq/bin/mqadmin updateTopic \
-n 192.168.10.101:9876 \
-c RaftCluster \
-t TestTopic \
-r 3 -w 3
3. 发送/消费测试消息
# 生产者(任意节点)
/opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer
# 消费者(任意节点)
/opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer
✅ 如果消费者能收到消息,说明集群工作正常!
🔄 六、高可用测试(模拟故障)
1. 杀死 Leader 节点的 Broker
# 假设 node1 是 Leader
pkill -f mqbroker
2. 观察其他节点日志
# 在 node2 或 node3 上
tail -f ~/logs/rocketmqlogs/broker.log
你会看到:
[DLedger] Try to elect a new leader...
[DLedger] New leader elected: n1
3. 恢复服务
重新启动 node1 的 Broker,它会自动加入集群并成为 Follower。
📦 七、一键启停脚本(可选)
创建 /opt/rocketmq/bin/start-cluster.sh:
#!/bin/bash
# 启动 NameServer 和 Broker
nohup /opt/rocketmq/bin/mqnamesrv > /dev/null 2>&1 &
nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/dledger.conf > /dev/null 2>&1 &
echo "RocketMQ Dledger cluster started."
赋予执行权限:
chmod +x /opt/rocketmq/bin/start-cluster.sh
✅ 八、生产建议
| 项目 | 建议 |
|---|---|
| JVM 参数 | 调整堆内存:export JAVA_OPT="-server -Xms4g -Xmx4g" |
| 监控 | 部署RocketMQ Exporter + Prometheus |
| 日志轮转 | 配置 logback.xml 自动切割日志 |
| 安全 | 开启 ACL(访问控制),避免未授权访问 |
| 备份 | 定期备份 /data/rocketmq/store |
九、Springboot集成
🧩 一、项目依赖(pom.xml)
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.2.0</version>
<relativePath/>
</parent>
<groupId>com.example</groupId>
<artifactId>rocketmq-demo</artifactId>
<version>1.0.0</version>
<properties>
<java.version>17</java.version>
<rocketmq-spring-boot-starter.version>2.2.3</rocketmq-spring-boot-starter.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- RocketMQ Spring Boot Starter -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version> $ {rocketmq-spring-boot-starter.version}</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
</project>
✅ 注意:
- Spring Boot 3.x 需搭配
rocketmq-spring-boot-starter >= 2.2.0- 若用 Spring Boot 2.x,请降级 starter 到
2.1.1
⚙️二、配置文件(application.yml)
server:
port: 8080
rocketmq:
name-server: 192.168.10.101:9876;192.168.10.102:9876;192.168.10.103:9876
producer:
group: demo-producer-group
consumer:
group: demo-consumer-group
# 自定义 Topic(可选)
app:
topic:
normal: NormalTopic
order: OrderTopic
transaction: TransactionTopic
📦 三、消息实体类(Order.java)
package com.example.rocketmqdemo;
import lombok.Data;
import java.io.Serializable;
import java.math.BigDecimal;
import java.time.LocalDateTime;
@Data
public class Order implements Serializable {
private Long orderId;
private String userId;
private String product;
private BigDecimal amount;
private LocalDateTime createTime;
}
📤 四、生产者服务(ProducerService.java)
package com.example.rocketmqdemo;
import jakarta.annotation.PostConstruct;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.concurrent.ThreadLocalRandom;
@Slf4j
@Service
public class ProducerService {
private final RocketMQTemplate rocketMQTemplate;
@Value(" $ {app.topic.normal}")
private String normalTopic;
@Value(" $ {app.topic.order}")
private String orderTopic;
@Value(" $ {app.topic.transaction}")
private String transactionTopic;
public ProducerService(RocketMQTemplate rocketMQTemplate) {
this.rocketMQTemplate = rocketMQTemplate;
}
// 1. 同步发送(可靠,但阻塞)
public void sendNormalMessage(String message) {
SendResult result = rocketMQTemplate.syncSend(normalTopic, message);
log.info("同步发送结果: {}", result.getSendStatus());
}
// 2. 异步发送(非阻塞,回调处理)
public void sendAsyncMessage(Order order) {
rocketMQTemplate.asyncSend(normalTopic, order, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("异步发送成功: orderId={}", order.getOrderId());
}
@Override
public void onException(Throwable e) {
log.error("异步发送失败: orderId={}, error={}", order.getOrderId(), e.getMessage());
}
});
}
// 3. 顺序消息(按订单ID哈希保证同一订单有序)
public void sendOrderlyMessage(Order order) {
rocketMQTemplate.syncSendOrderly(orderTopic, order, String.valueOf(order.getOrderId()));
log.info("顺序消息已发送: orderId={}", order.getOrderId());
}
// 4. 事务消息(模拟下单+扣库存)
public void sendTransactionMessage(Order order) {
String destination = transactionTopic + ":order_create";
rocketMQTemplate.sendMessageInTransaction(destination,
MessageBuilder.withPayload(order).build(), null);
log.info("事务消息已发送: orderId={}", order.getOrderId());
}
}
📥 五、消费者监听器(ConsumerListener.java)
package com.example.rocketmqdemo;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
@RocketMQMessageListener(
topic = " $ {app.topic.normal}",
consumerGroup = "demo-consumer-group"
)
public class NormalConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
log.info("收到普通消息: {}", message);
}
}
// 订单消费者(处理对象)
@Component
@RocketMQMessageListener(
topic = " $ {app.topic.normal}",
consumerGroup = "order-consumer-group"
)
class OrderConsumer implements RocketMQListener<Order> {
@Override
public void onMessage(Order order) {
log.info("处理订单: orderId={}, user={}", order.getOrderId(), order.getUserId());
// 模拟业务逻辑:更新订单状态、发通知等
}
}
// 顺序消息消费者
@Component
@RocketMQMessageListener(
topic = " $ {app.topic.order}",
consumerGroup = "orderly-consumer-group",
consumeMode = org.apache.rocketmq.spring.annotation.ConsumeMode.ORDERLY
)
class OrderlyConsumer implements RocketMQListener<Order> {
@Override
public void onMessage(Order order) {
log.info("顺序处理订单: orderId={}", order.getOrderId());
// 保证同一订单的操作严格有序(如:创建 → 支付 → 发货)
}
}
六、事务消息监听器(TransactionListener.java)
package com.example.rocketmqdemo;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.springframework.messaging.Message;
import org.springframework.stereotype.Component;
@Slf4j
@Component
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
// 1. 执行本地事务(如:插入订单到DB)
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
Order order = (Order) msg.getPayload();
try {
// 模拟数据库操作
log.info("执行本地事务: 创建订单 orderId={}", order.getOrderId());
// 模拟异常(可注释测试回滚)
if (order.getAmount().compareTo(new java.math.BigDecimal("1000")) > 0) {
throw new RuntimeException("金额超限");
}
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
log.error("本地事务失败: {}", e.getMessage());
return RocketMQLocalTransactionState.ROLLBACK;
}
}
// 2. 事务状态回查(Broker 定期调用)
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
Order order = (Order) msg.getPayload();
log.info("回查订单状态: orderId={}", order.getOrderId());
// 实际应查询DB确认订单是否存在
return RocketMQLocalTransactionState.COMMIT;
}
}
🧪 七、测试控制器(TestController.java)
package com.example.rocketmqdemo;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.concurrent.atomic.AtomicLong;
@RestController
@RequestMapping("/test")
public class TestController {
private final ProducerService producerService;
private final AtomicLong orderIdGenerator = new AtomicLong(1);
public TestController(ProducerService producerService) {
this.producerService = producerService;
}
@GetMapping("/normal")
public String sendNormal() {
producerService.sendNormalMessage("Hello RocketMQ!");
return "Normal message sent";
}
@GetMapping("/order")
public String sendOrder() {
Order order = new Order();
order.setOrderId(orderIdGenerator.getAndIncrement());
order.setUserId("user_123");
order.setProduct("iPhone 15");
order.setAmount(new BigDecimal("5999"));
order.setCreateTime(LocalDateTime.now());
producerService.sendOrderlyMessage(order);
return "Order message sent: " + order.getOrderId();
}
@GetMapping("/transaction")
public String sendTransaction() {
Order order = new Order();
order.setOrderId(orderIdGenerator.getAndIncrement());
order.setUserId("user_456");
order.setProduct("MacBook Pro");
order.setAmount(new BigDecimal("15000")); // 超过1000会触发回滚
order.setCreateTime(LocalDateTime.now());
producerService.sendTransactionMessage(order);
return "Transaction message sent: " + order.getOrderId();
}
}
▶️ 八、启动与验证
1.启动应用
mvn spring-boot:run
2.发送测试消息
# 普通消息
curl http://localhost:8080/test/normal
# 顺序消息
curl http://localhost:8080/test/order
# 事务消息(金额>1000会回滚)
curl http://localhost:8080/test/transaction
3.查看日志
# 成功示例
INFO NormalConsumer: 收到普通消息: Hello RocketMQ!
INFO OrderlyConsumer: 顺序处理订单: orderId=1
INFO OrderTransactionListener: 执行本地事务: 创建订单 orderId=2
ERROR OrderTransactionListener: 本地事务失败: 金额超限 # 金额>1000时
📚 附:常用命令速查
| 功能 | 命令 |
|---|---|
| 查看 topic | ./mqadmin topicList -n 192.168.10.101:9876 |
| 查看 broker 状态 | ./mqadmin brokerStatus -n ... -b 192.168.10.101:30911 |
| 发送消息 | ./tools.sh org.apache.rocketmq.example.quickstart.Producer |
| 消费消息 | ./tools.sh org.apache.rocketmq.example.quickstart.Consumer |
更多推荐




所有评论(0)