下面把 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:高阶功能(顺序/延迟/批量)

  1. 顺序消息(同一订单ID进入同一队列) 发送时指定 HashKey:

    rocketMQTemplate.syncSendOrderly("pay-topic", msg, orderId);
  2. 延迟消息(18级 1s~2h)

    Message<String> msg = MessageBuilder.withPayload(body).build(); 
    msg.setDelayTimeLevel(3); // 3=10s 
    rocketMQTemplate.syncSend("delay-topic", msg);
  3. 批量发送(提高吞吐)

    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:备份 & 升级 & 应急

  1. 消息备份 使用 mqadmin exportMetadata 导出元数据,数据文件直接 rsync /store 目录。

  2. 平滑升级 先升级 Slave → 切换 Master → 升级旧 Master,0 停机。

  3. 消息堆积应急

  • 临时扩容消费者线程

  • 调整 consumeThreadMax=100

  • 业务端做批量消费(List<MessageExt> 一次性处理 1000 条)

保姆级总结(口诀)

“装好 JDK 调内存,外网 IP 必配精; 事务消息半提交,回查状态要记清; 顺序延迟批量发,监控 Lag 别停盯; Slave 先升再切换,备份 store 零宕机!”

照抄 12 阶段,从开发到生产,RocketMQ 集群任你玩转!

Logo

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

更多推荐