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
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐