RocketMQ安装部署

参考: 集群配置

一、原理

1.1、核心架构

  • RocketMQ 采用分布式架构设计,核心由 4 个角色组成,各角色解耦且支持水平扩展:

    角色 核心职责 部署特性
    NameServer 轻量级注册中心,存储 Topic 与 Broker 路由信息,无状态、可集群部署 单机 / 集群均可,生产建议≥2 台
    Broker 消息存储与转发核心节点,接收 Producer 消息、存储消息、推 / 拉消息给 Consumer 生产需主从部署(Dledger 切换)
    Producer 消息生产者,负责将消息发送到指定的 Topic。每个 Producer 可以连接多个 Broker(支持集群部署、负载均衡),发送消息到不同的队列 无状态,业务侧集群部署
    Consumer 消息消费者,负责从 Broker 拉取消息进行消费。Consumer 可以是 推模式(Push)或 拉模式(Pull),支持异步或同步消费 按 Consumer Group 集群部署
  • 核心流程

    1. Producer/Consumer 启动时向 NameServer 获取 Topic 对应的 Broker 路由信息;

    2. Broker 定期向 NameServer 上报心跳、路由和负载信息,保证路由数据最新;

    3. Producer 根据路由信息将消息发送到指定 Broker 的 Queue;

    4. Consumer 从 Broker 拉取消息,消费完成后提交偏移量(Offset)。

      请添加图片描述

1.2、核心概念与机制

  1. 消息存储模型.

    Broker 采用「CommitLog + ConsumeQueue + IndexFile」三层存储结构:

    • CommitLog:所有 Topic 的消息统一存储到 CommitLog 日志文件(默认 1G / 个),是消息的物理存储载体;
    • ConsumeQueue:消费队列,建立 Topic-Queue 到 CommitLog 的索引(仅存储偏移量、消息长度等元数据),提升消费效率;
    • IndexFile:消息索引文件,支持按消息 Key 快速查询消息。
  2. 高可用核心机制

    • Broker 主从:Master 负责写操作,Slave 同步 CommitLog 数据负责读操作,Master 宕机后基于 Dledger 自动切换到 Slave;
    • 刷盘策略
      • 同步刷盘(SYNC_FLUSH):消息写入后立即刷盘,可靠性最高;
      • 异步刷盘(ASYNC_FLUSH):先写入内存页缓存,批量刷盘,性能更高;
    • 重试与死信:消费失败自动重试(默认 16 次),重试失败的消息进入死信队列(DLQ),需人工干预;
    • 消息持久化:所有消息落地磁盘,支持按时间 / 大小清理过期消息。
  3. 消费模式

    • 集群消费:同一 Consumer Group 的多个 Consumer 分摊消费 Queue 消息(一条消息仅消费一次),默认模式;
    • 广播消费:同一 Consumer Group 的每个 Consumer 都消费全量消息(一条消息被所有实例消费)。

二、单机部署

单机部署为「单 NameServer + 单 Broker(无主从)」

2.1、前置准备

  • 前置条件

    依赖项 版本 / 配置要求 说明
    JDK 8+(推荐 1.8.0_20+) RocketMQ 基于 Java 开发
    操作系统 Linux(CentOS 7/8)/麒麟v10/v11 Linux(CentOS 7/8)/麒麟v10/v11
    内存 单机 ≥1G,集群节点 ≥8G NameServer 建议 2G+,Broker 建议 4G+
    磁盘 生产节点 ≥100G(SSD) 单独挂载存储目录,避免系统盘占满
    端口 9876(NameServer)、
    10911/10909/10912(Broker)
    开放防火墙 / 安全组
  • 安装包准备

    # java略过:  java version "1.8.0_461"
    
    # 下载稳定版(4.9.7 为例)
    wget https://archive.apache.org/dist/rocketmq/4.9.7/rocketmq-all-4.9.7-bin-release.zip
    
    # 解压并创建目录
    unzip rocketmq-all-4.9.7-bin-release.zip -d /data
    mv /data/rocketmq-all-4.9.7-bin-release /data/rocketmq
    mkdir -p /data/rocketmq/{logs,store}  # 日志/数据目录
    chmod 777 /data/rocketmq -R
    
  • JVM 内存调整

    • 修改 NameServer JVM 配置

      vim /data/rocketmq/bin/runserver.sh
      # 调整 Xms/Xmx 为 1G(原配置为 -Xms4g -Xmx4g)
      JAVA_OPT="${JAVA_OPT} -server -Xms1g -Xmx1g -Xmn512m -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m"
      
    • 修改 Broker JVM 配置

      vim /usr/local/rocketmq/bin/runbroker.sh
      # 调整 Xms/Xmx 为 2G(原配置为 -Xms8g -Xmx8g)
      JAVA_OPT="${JAVA_OPT} -server -Xms2g -Xmx2g -Xmn1g"
      
  • 配置文件修改

    • 修改 broker.conf

      mkdir -p /data/rocketmq/conf
      cat > /data/rocketmq/conf/broker.conf << EOF
      # ===================== 基础集群配置 =====================
      # 集群名称(单机默认 DefaultCluster,多Broker集群时需统一命名)
      brokerClusterName = DefaultCluster
      # Broker 节点名称(主从集群中,主从节点需同名,单机仅一个节点默认 broker-a)
      brokerName = broker-a
      # Broker ID(0=主节点,>0=从节点;单机仅主节点,固定为0)
      brokerId = 0
      # NameServer 地址(多个用分号分隔,单机仅 localhost:9876)
      namesrvAddr = localhost:9876
      # Broker 监听端口(默认 10911,避免与其他服务冲突)
      listenPort = 10911
      # 是否自动创建 Topic(测试环境开启,生产环境建议关闭,手动创建Topic)
      autoCreateTopicEnable = true
      
      # ===================== 存储/日志目录配置 =====================
      # Broker 日志根目录(替代启动命令的 --logBaseDir 参数)
      logBaseDir = /data/rocketmq/logs
      # Broker 数据存储根目录(替代启动命令的 --storePathRootDir 参数)
      storePathRootDir = /data/rocketmq/store
      
      # ===================== 高可用/刷盘配置 =====================
      # 刷盘策略(ASYNC_FLUSH=异步刷盘(性能优先),SYNC_FLUSH=同步刷盘(可靠性优先))
      flushDiskType = ASYNC_FLUSH
      # Broker 角色(ASYNC_MASTER=异步复制主节点(单机默认),SYNC_MASTER=同步复制主节点,SLAVE=从节点)
      brokerRole = ASYNC_MASTER
      
      # ===================== 消息清理配置 =====================
      # 过期消息清理触发时间(整点小时,04=凌晨4点,选业务低峰期执行)
      deleteWhen = 04
      # 消息文件保留时长(单位:小时,48=2天;RocketMQ 4.x 单位为小时,5.x 可指定天)
      fileReservedTime = 48
      EOF
      
    • 定义日志目录

      cd /data/rocketmq/conf
          
      [root@node1 conf]# ls logback_*
      logback_broker.xml  logback_namesrv.xml  logback_tools.xml
      
      # 三个配置文件,都是一样的改法, 添加 property
      <configuration scan="true" scanPeriod="30 seconds">
          <property name="user.home" value="/data/rocketmq"/>
          <appender name="DefaultAppender"...ileAppender">
          <file>${user.home}/logs/roc.........log</file>
      
      # 如果不定义,日志文件则在 user.home,启动对应的用户家目录 如 /root/logs/rocketmqlogs
      

      请添加图片描述

2.2、服务启动及使用

  • 启动服务 <-- 共两个 【namesrv跟broker】

    • namesrv

      # 手动启动,不推荐
      # 后台启动,日志输出到指定文件
      nohup /data/rocketmq/bin/mqnamesrv > /data/rocketmq/logs/namesrv.log 2>&1 &
      
      ############# 用systemctl ######
      cat > /usr/lib/systemd/system/rocketmq-namesrv.service << EOF
      [Unit]
      Description=RocketMQ NameServer Service
      Documentation=https://rocketmq.apache.org/
      After=network.target
      Wants=network.target
      
      [Service]
      Type=simple
      # 指定 RocketMQ 运行用户(建议创建专用用户,此处以 root 为例,生产可改为 rocketmq)
      User=root
      Group=root
      # 环境变量(指定 JAVA_HOME,根据实际安装路径调整)
      Environment="JAVA_HOME=/data/jdk1.8"
      Environment="ROCKETMQ_HOME=/data/rocketmq"
      # 启动命令(日志重定向 + 自定义 JVM 参数)
      ExecStart=/data/rocketmq/bin/mqnamesrv
      # 日志输出重定向(与手动启动一致)
      StandardOutput=append:/data/rocketmq/logs/namesrv.log
      StandardError=append:/data/rocketmq/logs/namesrv.log
      # 进程停止方式
      ExecStop=/data/rocketmq/bin/mqshutdown namesrv
      # 重启策略(异常退出时自动重启)
      Restart=on-failure
      RestartSec=5
      # 资源限制
      LimitNOFILE=65535
      LimitNPROC=65535
      
      [Install]
      WantedBy=multi-user.target
      EOF
      
      ##### 检查
      bin]# systemctl start rocketmq-namesrv
      bin]# ss -tnlp|grep 987
      LISTEN     0      128  :::9876 :::*  users:(("java"
      
    • Broker

      # 手动启动 指定 NameServer 地址、日志/数据目录,后台启动
      nohup /data/rocketmq/bin/mqbroker -n localhost:9876 \
      --logBaseDir /data/rocketmq/logs \
      --storePathRootDir /data/rocketmq/store > /data/rocketmq/logs/broker.log 2>&1 &
      
      ############# 用systemctl ######
      cat > /usr/lib/systemd/system/rocketmq-broker.service << EOF
      [Unit]
      Description=RocketMQ Broker Service
      Documentation=https://rocketmq.apache.org/
      After=network.target rocketmq-namesrv.service
      Requires=rocketmq-namesrv.service
      
      [Service]
      Type=simple
      User=root
      Group=root
      Environment="JAVA_HOME=/data/jdk1.8"
      Environment="ROCKETMQ_HOME=/data/rocketmq"
      Environment="PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin:${JAVA_HOME}/bin"
      
      # 启动命令:仅用原生 -c(指定配置文件)和 -n(NameServer)参数
      # 所有目录/业务配置都在 broker.conf 中,避免非法参数
      ExecStart=/data/rocketmq/bin/mqbroker -c /data/rocketmq/conf/broker.conf -n localhost:9876
      
      # 日志重定向
      StandardOutput=append:/data/rocketmq/logs/broker.log
      StandardError=append:/data/rocketmq/logs/broker.log
      
      # 停止命令(修复环境变量后可正常执行)
      ExecStop=/data/rocketmq/bin/mqshutdown broker
      
      # 重启策略
      Restart=on-failure
      RestartSec=5
      
      # 资源限制
      LimitNOFILE=65535
      LimitNPROC=65535
      
      # 关键:指定工作目录 + 防止进程被过早杀死
      WorkingDirectory=/data/rocketmq
      KillMode=process
      TimeoutStopSec=60
      
      [Install]
      WantedBy=multi-user.target
      EOF
      
  • 单机运维常用命令

    运行
    # 停止 NameServer
    /data/rocketmq/bin/mqshutdown namesrv
    
    # 停止 Broker
    /data/rocketmq/bin/mqshutdown broker
    
    # 查看集群状态(需配置 NameServer 地址)
    export NAMESRV_ADDR=localhost:9876
    /data/rocketmq/bin/mqadmin clusterList -n localhost:9876
    
    # 查看 Topic 列表
    /data/rocketmq/bin/mqadmin topicList -n localhost:9876
    
  • 单机测试示例

    • 发送测试消息

      # 使用官方工具发送消息
      export NAMESRV_ADDR=localhost:9876
      /data/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer
      
      
    • 消费测试消息

      export NAMESRV_ADDR=localhost:9876
      /data/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer
      

2.3、端口说明

  • 端口说明

    端口号 所属组件 核心功能(单机场景) 配置 / 验证方式
    9876 NameServer 路由中心唯一通信端口:1. Broker 向其上报自身信息;2. 生产者 / 消费者查询 Topic 路由 配置:namesrv.conflistenPort=9876(默认)
    10911 Broker 业务核心端口:1. 接收生产者发消息;2. 接收消费者拉消息 / 确认消费 配置:broker.conflistenPort=10911
  • 端口详解

    • 9876(NameServer 唯一端口)
      • 单机定位:NameServer 是 RocketMQ 的 “导航仪”,9876 是它的唯一对外通信口,单机部署时 NameServer 仅监听这一个端口。
      • 单机核心场景
        1. Broker 启动后,通过 9876 向 NameServer 上报 “我是 broker-a,地址是 127.0.0.1:10911”;
        2. 生产者发消息前,通过 9876 问 NameServer:“TopicTest1 存在哪个 Broker 上?”,NameServer 返回 127.0.0.1:10911;
        3. 消费者启动时,同样通过 9876 查询 Topic 对应的 Broker 地址,建立消费连接。
    • 10911(Broker 唯一业务端口)
      • 单机定位:Broker 是 RocketMQ 的 “消息仓库”,10911 是仓库的 “业务大门”,单机部署时 Broker 仅监听这一个端口(10909/10912 无意义)。
      • 单机核心场景
        1. 生产者拿到 NameServer 给的 10911 地址后,通过该端口向 Broker 发送消息;
        2. 消费者通过 10911 端口从 Broker 拉取消息,消费完成后再通过该端口告诉 Broker“消息已消费,可删除”;
        3. 所有消息的存储、刷盘、重试,都基于 10911 端口的通信完成。

三、集群部署

3.1、前置准备

  • 架构规划

    机器 IP 角色 核心端口 核心配置标识
    10.4.50.130 NameServer + Master Broker 9876(NS)、10911(业务)、10909(主从复制) brokerName=broker-a、brokerId=0
    10.4.50.139 NameServer + Slave1 Broker 9876(NS)、10911(业务)、10909(主从复制) brokerName=broker-a、brokerId=1
    10.4.50.167 NameServer + Slave2 Broker 9876(NS)、10911(业务)、10909(主从复制) brokerName=broker-a、brokerId=2
    • 设计说明
      • NameServer 集群:3 台都部署 NameServer,避免 NS 单点故障;
      • Broker 主从:130=Master,139/167=Slave,同组(brokerName=broker-a),数据实时同步;
      • 端口统一:所有节点端口保持一致,降低运维成本;
      • 高可用:NS 集群 + Broker 主从,任意 1 台 Slave/NS 宕机不影响核心服务。
  • 前期服务器准备

    • 关闭防火墙/SELinux(生产可按需开放端口,测试直接关闭)

      systemctl stop firewalld && systemctl disable firewalld
      setenforce 0 && sed -i 's/^SELINUX=.*/SELINUX=disabled/' /etc/selinux/config
      
    • 下载jdk,用浏览器下载最新的,按架构下载对应包

      https://www.oracle.com/java/technologies/javase/javase8u211-later-archive-downloads.html
      # 将jdk放到 /data/目录下  目录名: /data/jdk1.8
      echo "JAVA_HOME=/data/jdk1.8" >> /etc/profile
      echo "export PATH=\$JAVA_HOME/bin:\$PATH" >> /etc/profile
      
      
    • rocketmq准备

      # 下载稳定版(4.9.7 为例)
      wget https://archive.apache.org/dist/rocketmq/4.9.7/rocketmq-all-4.9.7-bin-release.zip
      
      # 解压并创建目录
      unzip rocketmq-all-4.9.7-bin-release.zip -d /data
      mv /data/rocketmq-all-4.9.7-bin-release /data/rocketmq
      mkdir -p /data/rocketmq/{logs,store}  # 日志/数据目录
      chmod 777 /data/rocketmq -R
      
      # 修改NameServer JVM内存
      sed -i "s@-server -Xms4g -Xmx4g -Xmn2g@-server -Xms2g -Xmx2g -Xmn1g@gi" /data/rocketmq/bin/runserver.sh 
      
      # 修改Broker JVM内存
      sed -i "s@-server -Xms8g -Xmx8g@-server -Xms2g -Xmx2g@gi" /data/rocketmq/bin/runbroker.sh
      
      # 定义日志路径,跟单机部署一样配置
      cd /data/rocketmq/conf
          
      [root@node1 conf]# ls logback_*
      logback_broker.xml  logback_namesrv.xml  logback_tools.xml
      
      # 三个配置文件,都是一样的改法, 添加 property
      <configuration scan="true" scanPeriod="30 seconds">
          <property name="user.home" value="/data/rocketmq"/>
          <appender name="DefaultAppender" .....
      
    • 准备工作完成,可以将rocketmq整个压缩,复制到另外两台上

3.2、集群-手动切主

用3.3章节,这里用namesrc.conf跟两个.service文件

  • namesrc.conf --> 三台都一样

    # 三台都加上
    cat > /data/rocketmq/conf/namesrv.conf << EOF
    # NameServer监听端口
    listenPort=9876
    # 日志目录
    logDir=/data/rocketmq/logs
    EOF
    
    • 自启脚本

      cat > /usr/lib/systemd/system/rocketmq-namesrv.service << EOF
      [Unit]
      Description=RocketMQ NameServer Service
      Documentation=https://rocketmq.apache.org/
      After=network.target
      Wants=network.target
      
      [Service]
      Type=simple
      User=root
      Group=root
      Environment="JAVA_HOME=/data/jdk1.8"
      Environment="ROCKETMQ_HOME=/data/rocketmq"
      # 修正:补充主类名 + -c 指定配置文件
      ExecStart=/data/rocketmq/bin/runserver.sh org.apache.rocketmq.namesrv.NamesrvStartup -c /data/rocketmq/conf/namesrv.conf
      # 修正 ExecStop:避免 kill -9 强制杀死(推荐优雅停止),且语法更安全
      ExecStop=/bin/sh -c 'pid=\$(pgrep -f "org.apache.rocketmq.namesrv.NamesrvStartup"); if [ -n "\$pid" ]; then kill \$pid; fi'
      Restart=on-failure
      RestartSec=5
      LimitNOFILE=65535
      LimitNPROC=65535
      WorkingDirectory=/data/rocketmq
      
      [Install]
      WantedBy=multi-user.target
      EOF
      
  • broker.conf --> 主节点

    # 创建Master Broker配置文件
    cat > /data/rocketmq/conf/broker.conf << EOF
    # 集群名称
    brokerClusterName=MyRocketMQCluster
    # 主从组名(必须一致)
    brokerName=broker-a
    # 0=Master,1/2=Slave
    brokerId=0
    # NameServer集群地址(3台都配置)
    namesrvAddr=10.4.50.130:9876;10.4.50.139:9876;10.4.50.167:9876
    # Broker业务端口
    listenPort=10911
    # 主从复制端口(默认listenPort-2)
    haListenPort=10909
    # 自动创建Topic(测试开启,生产关闭)
    autoCreateTopicEnable=true
    # Master角色:SYNC_MASTER(同步双写,高可用)
    brokerRole=SYNC_MASTER
    # 刷盘策略:ASYNC_FLUSH(性能优先)
    flushDiskType=ASYNC_FLUSH
    # 日志/存储目录
    logBaseDir=/data/rocketmq/logs
    storePathRootDir=/data/rocketmq/store
    # 消息保留时长(小时)
    fileReservedTime=48
    # 过期消息清理时间(凌晨4点)
    deleteWhen=04
    EOF
    
    • 自启脚本

      # 创建Master Broker systemd服务
      cat > /usr/lib/systemd/system/rocketmq-broker.service << EOF
      [Unit]
      Description=RocketMQ Master Broker
      After=network.target rocketmq-namesrv.service
      Requires=rocketmq-namesrv.service
      
      [Service]
      Type=simple
      User=root
      Group=root
      Environment="JAVA_HOME=/data/jdk1.8"
      Environment="ROCKETMQ_HOME=/data/rocketmq"
      # 核心修正:补充Broker主类名 + -c指定配置文件
      ExecStart=/data/rocketmq/bin/runbroker.sh org.apache.rocketmq.broker.BrokerStartup -c /data/rocketmq/conf/broker.conf
      # 优化ExecStop:优雅停止,避免kill -9强制杀死
      ExecStop=/bin/sh -c 'pid=\$(pgrep -f "org.apache.rocketmq.broker.BrokerStartup"); if [ -n "\$pid" ]; then kill \$pid; fi'
      Restart=on-failure
      RestartSec=5
      LimitNOFILE=65535
      LimitNPROC=65535
      WorkingDirectory=/data/rocketmq
      
      [Install]
      WantedBy=multi-user.target
      EOF
      
    • 启动测试一下

      systemctl start rocketmq-namesrv
      systemctl start rocketmq-broker
      
      # 验证 Broker 与 NameServer 连接
      # 进入 RocketMQ 安装目录
      cd /data/rocketmq/bin
      
      # 查看 NameServer 上注册的 Broker 信息
      ./mqadmin clusterList -n 10.4.50.130:9876
      
      [root@node1 bin]# ./mqadmin clusterList -n 10.4.50.130:9876
      #Cluster Name     #Broker Name            #BID  #Addr                  #Version                #InTPS(LOAD)       #OutTPS(LOAD) #PCWait(ms) #Hour #SPACE
      MyRocketMQCluster  broker-a                0     10.4.50.130:10911      V4_9_7                   0.00(0,0ms)         0.00(0,0ms)          0 490326.78 0.1400
      
      字段 数值 / 说明 状态判断
      Cluster Name MyRocketMQCluster broker.conf 中配置一致 ✅
      Broker Name broker-a broker.conf 中配置一致 ✅
      BID 0 0=Master(主节点),符合配置 ✅
      Addr 10.4.50.130:10911 Broker 监听端口正常 ✅
      Version V4_9_7 RocketMQ 版本正常 ✅
      InTPS/OutTPS 0.00 暂无消息收发(新集群正常)✅
      SPACE 0.1400 存储磁盘空间正常 ✅
  • 两台从节点

    • namesrc.conf 三台都一样,原样复制黏贴

    • broker.conf

      # 节点 139
      # Slave1的brokerId=1
      brokerId=1
      brokerRole=SLAVE
      
      # 节点 167 
      # Slave2的brokerId=2
      brokerId=2
      brokerRole=SLAVE
      
    • 启动服务

      systemctl start rocketmq-namesrv
      systemctl start rocketmq-broker
      
      # 启动需要一等会, 当端口都出来时,ok
      data]# ss -tnlp
      State      监控端口   
      LISTEN    [::]:9876     users:(("java",pid=1676,fd=93))
      LISTEN    [::]:10909    users:(("java",pid=1732,fd=128))
      LISTEN    [::]:10911    users:(("java",pid=1732,fd=127))
      LISTEN    [::]:10912    users:(("java",pid=1732,fd=123))
      
    • 验证集群状态

      [root@node1 bin]# cd /data/rocketmq/bin
      [root@node1 bin]# ./mqadmin clusterList -n 集群内任意ip:9876
      #Cluster Name     #Broker Name            #BID  #Addr              
      MyRocketMQCluster  broker-a                0     10.4.50.130:10911 
      MyRocketMQCluster  broker-a                1     10.4.50.139:10911 
      MyRocketMQCluster  broker-a                2     10.4.50.167:10911 
      
  • 测试相关命令

    # 1. 配置NameServer集群地址(临时生效)
    export NAMESRV_ADDR="10.4.50.130:9876;10.4.50.139:9876;10.4.50.167:9876"
    
    # 2. 查看集群状态(能看到1个Master+2个Slave)
    /data/rocketmq/bin/mqadmin clusterList
    
    # 3. 创建测试Topic
    /data/rocketmq/rocketmq/bin/mqadmin updateTopic -n ${NAMESRV_ADDR} -c MyRocketMQCluster -t TestTopic
    create topic to 10.4.50.130:10911 success.
    TopicConfig [topicName=TestTopic, readQueueNums=8, writeQueueNums=8, perm=RW-, topicFilterType=SINGLE_TAG, topicSysFlag=0, order=false]
    
    # 4. 发送测试消息
    /data/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer -n ${NAMESRV_ADDR} -t TestTopic
    
    # 5. 消费测试消息
    /data/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer -n ${NAMESRV_ADDR} -t TestTopic
    
3.2.1、端口说明
  • 详细说明

    端口号 端口名称 / 用途 角色(主 / 从) 核心作用 配置参数
    9876 NameServer 通信端口 所有 NameServer 1. 接收 Broker 节点的注册 / 心跳 / 注销请求;
    2. 接收客户端(生产者 / 消费者)的 Topic 路由查询请求;
    3. 集群内 NameServer 无端口交互(去中心化)。
    listenPort(namesrv.conf)
    10911 Broker 客户端通信端口 主 / 从均启用 1. 接收生产者的消息发送请求;
    2. 接收消费者的消息拉取 / 确认请求;
    3. 主从节点此端口独立,互不影响。
    listenPort(broker.conf)
    10909 HA 主从复制端口(haListenPort) 仅主节点启用 1. 从节点通过此端口连接主节点,同步 commitlog 日志;
    2. 主节点监听此端口,接收从节点的同步请求;
    3. 从节点无此端口监听(仅作为客户端连接主节点)。
    haListenPort(broker.conf,默认 = listenPort-2)
    10912 Broker 内部通信端口(fastPort) 主 / 从均启用 1. 用于 Broker 内部线程间通信、请求快速转发;
    2. 辅助 10911 端口处理高并发场景的轻量请求;
    3. 无需手动配置(默认 = listenPort+1),极少对外暴露。
    无显式配置(自动 = listenPort+1)
  • 补充说明

    • 端口依赖关系

      • 10909 是主节点专属,依赖 10911(默认 = 10911-2);
      • 10912 是 10911 的辅助端口(默认 = 10911+1),无需单独配置;
      • 9876 与 Broker 端口无直接关联,仅通过「路由注册」间接交互。
    • 端口互通要求

      • 客户端 ↔ NameServer(9876):必须互通;
      • 客户端 ↔ Broker 主 / 从(10911):必须互通;
      • Broker 从 ↔ Broker 主(10909):必须互通(主从同步核心);
      • 10912 仅 Broker 本机使用,无需对外放行。
    • 核心配置示例

      # broker.conf(主节点)
      listenPort=10911          # 客户端通信端口
      haListenPort=10909        # HA 主从复制端口(默认=10911-2,可省略)
      brokerId=0                # 主节点标识
      brokerRole=SYNC_MASTER    # 同步主节点
      haEnable=true             # 开启旧版 HA 同步
      
      # broker.conf(从节点)
      listenPort=10911          # 从节点客户端端口(可与主节点相同,IP 不同即可)
      brokerId=1                # 从节点标识
      brokerRole=SLAVE          # 从节点
      haMasterAddress=10.4.50.130:10909  # 主节点 HA 端口(指定同步源)
      

3.3、集群-自动选主(推荐)

DLedger 基于 Raft 协议实现 Broker 集群自动选主

  • namesrc.conf --> 参考手动切主,三台一样

  • broker.conf --> 主节点

    cat > /data/rocketmq/conf/broker.conf << EOF
    # 基础集群配置
    brokerClusterName=MyRocketMQCluster
    brokerName=broker-a
    # 【删除/注释】brokerId=0 (DLedger 自动管理,手动配置会冲突)
    # brokerId=0
    namesrvAddr=10.4.50.130:9876;10.4.50.139:9876;10.4.50.167:9876
    listenPort=10911
    autoCreateTopicEnable=true
    
    # 日志/存储目录
    logBaseDir=/data/rocketmq/logs
    storePathRootDir=/data/rocketmq/store
    deleteWhen=04
    fileReservedTime=48
    
    # ===================== DLedger 核心配置 =====================
    enableDLegerCommitLog=true
    dLegerGroup=broker-a-group
    # 【修正】dLegerPeers 格式:节点ID-IP:DLedger端口(40911),3个节点必须一致
    dLegerPeers=n0-10.4.50.130:40911;n1-10.4.50.139:40911;n2-10.4.50.167:40911
    # 【每个节点单独配置】dLegerSelfId:130节点填n0,139填n1,167填n2
    dLegerSelfId=n0
    storePathDLedgerCommitLog=/data/rocketmq/store/dledger/commitlog
    haEnable=false
    brokerRole=ASYNC_MASTER
    flushDiskType=ASYNC_FLUSH
    
    # 性能调优
    sendMessageThreadPoolNums=16
    consumeMessageThreadPoolNums=32
    EOF
    
    # 自启参考:手动切主
    
  • 另外两台配置

    # 节点 139
    # Slave1的brokerId=1
    dLegerSelfId = n1
    
    # 节点 167 
    # Slave2的brokerId=2
    dLegerSelfId = n2
    
    # 其它的都不需要改动
    
    # 启动重启
    systemctl stop rocketmq-broker
    rm -rf /data/rocketmq/logs/*
    rm -rf /data/rocketmq/store/*
    systemctl start rocketmq-broker
    
  • 启动及日志分析

    • 130

      [root@node1 rocketmqlogs]# tail -300 /data/rocketmq/logs/rocketmqlogs/broker.log | grep -i "dledger\|follower"
      2025-12-08 16:44:01 INFO main - storePathDLedgerCommitLog=/data/rocketmq/store/dledger/commitlog
      2025-12-08 16:44:03 INFO DLegerRoleChangeHandler_1 - Begin handling broker role change term=1 role=FOLLOWER currStoreRole=SLAVE
      2025-12-08 16:44:03 INFO DLegerRoleChangeHandler_1 - Finish handling broker role change succ=true term=1 role=FOLLOWER currStoreRole=SLAVE cost=1
      
    • 139

      [root@node2 rocketmqlogs]# tail -300 /data/rocketmq/logs/rocketmqlogs/broker.log | grep -i "dledger\|follower"
      2025-12-08 16:44:00 INFO main - storePathDLedgerCommitLog=/data/rocketmq/store/dledger/commitlog
      
    • 167

      [root@node3 rocketmqlogs]# tail -300 /data/rocketmq/logs/rocketmqlogs/broker.log | grep -i "dledger\|follower"
      2025-12-08 16:44:00 INFO main - storePathDLedgerCommitLog=/data/rocketmq/store/dledger/commitlog
      2025-12-08 16:44:03 INFO DLegerRoleChangeHandler_1 - Begin handling broker role change term=1 role=FOLLOWER currStoreRole=SLAVE
      2025-12-08 16:44:03 INFO DLegerRoleChangeHandler_1 - Finish handling broker role change succ=true term=1 role=FOLLOWER currStoreRole=SLAVE cost=2
      
    • 核心状态解读

      节点 日志特征 集群角色 状态判断
      130 role=FOLLOWER currStoreRole=SLAVE DLedger Follower 正常加入集群 ✅
      139 storePathDLedgerCommitLog DLedger Leader 无 Follower 日志 → 因为它是 Leader(无需显示 Follower 日志)✅
      167 role=FOLLOWER currStoreRole=SLAVE DLedger Follower 正常加入集群 ✅
3.3.1、端口说明
  • 端口说明

    端口号 端口名称 / 用途 角色(Leader/Follower) 核心作用 配置参数 / 关键特性
    9876 NameServer 通信端口 所有 NameServer 1. 接收 Broker 节点(Leader/Follower)的注册 / 心跳 / 注销请求;2. 接收客户端(生产者 / 消费者)的 Topic 路由查询请求;3. 集群内 NameServer 去中心化,无端口交互。 listenPort(namesrv.conf)✅ 与传统模式功能完全一致
    10911 Broker 客户端通信端口 Leader/Follower 均启用 1. Leader 节点:接收生产者消息发送、消费者消息拉取 / 确认请求;
    2. Follower 节点:仅备用(客户端默认连接 Leader,Leader 宕机后路由自动切换至新 Leader);3. 主从节点此端口独立,IP 不同即可复用端口号。
    listenPort(broker.conf)✅ 与传统模式功能完全一致
    10909 旧版 HA 主从复制端口 所有 Broker 均禁用 传统模式的主从同步端口,DLedger 模式下因 haEnable=false 完全失效,无监听、无数据交互。 haListenPort(默认 = 10911-2)❌ 无需配置 / 放行,DLedger 模式下建议注释该参数
    40911 DLedger 集群同步端口(核心) Leader/Follower 均启用 1. 基于 Raft 协议的集群内数据同步:Leader 监听此端口,Follower 主动连接 Leader 的 40911 端口同步 commitlog;
    2. 参与 Leader 选举:节点间通过此端口交换投票、任期(term)等 Raft 协议核心信息;
    3. 集群内所有节点需互通此端口(缺一不可)。
    无显式配置参数(需在 dLegerPeers 中指定,格式:节点ID-IP:40911)✅ DLedger 模式专属,替代传统模式的 10909
  • 补充说明

    • 端口互通要求

      通信方向 需放行端口 核心目的
      客户端 ↔ 所有 NameServer 9876 获取路由、连接 Broker
      客户端 ↔ Broker Leader(10911) 10911 收发消息(核心业务)
      Broker 集群内(所有节点互访) 40911 Raft 协议同步、Leader 选举
      10909 无需放行 DLedger 模式下已禁用
    • 关键配置说明

      # broker.conf(Leader/Follower 通用配置,仅 dLegerSelfId 不同)
      listenPort=10911           # 客户端通信端口(与传统模式一致)
      # haListenPort=10909       # 建议注释,DLedger 模式下无需配置
      enableDLegerCommitLog=true # 开启 DLedger 模式(核心开关)
      dLegerGroup=broker-a-group # DLedger 集群组名(所有节点一致)
      # 集群节点配置:节点ID-IP:DLedger端口(40911)
      dLegerPeers=n0-10.4.50.130:40911;n1-10.4.50.139:40911;n2-10.4.50.167:40911
      dLegerSelfId=n0           # 节点专属ID(130=n0/139=n1/167=n2)
      haEnable=false            # 强制关闭旧版 HA(避免与 DLedger 冲突)
      
    • 总结:DLedger 模式仅对「主从同步端口」做了替换(40911 替代 10909),业务端口(9876/10911)完全复用传统模式逻辑,核心变化:

      1. 40911 是 DLedger 集群的 “命脉端口”,集群内所有节点必须互通;
      2. 10909 端口彻底失效,无需配置 / 放行;
      3. 10911 端口仅 Leader 提供业务服务,Follower 仅作为故障转移备用。

四、其它补充

4.1、模式对比

  • 手动主从模式与DLedger模式对比

    维度 传统主从模式 DLedger 模式
    核心配置 brokerId + brokerRole + haEnable=true enableDLegerCommitLog=true + dLeger* 配置 + haEnable=false
    同步协议 基于 HA 的主从复制 基于 Raft 协议的 DLedger 同步
    核心端口 10909(haListenPort) 40911(DLedger 端口)
    日志特征 仅输出 MASTER/SLAVE 角色日志 包含 Leader/Follower/term/raft 等 Raft 协议特征日志
    故障转移 手动切换(需修改 brokerId 等配置) 自动选举 Leader(无需人工干预)
    数据一致性 异步 / 同步复制(极端场景可能丢数据) Raft 协议保障数据强一致性
    运维成本 高(主节点宕机需人工介入操作) 低(故障自动恢复,无人工操作成本)
    集群可用性 中(主节点宕机期间服务不可用) 高(秒级自动切换 Leader,服务无感知)
  • 与传统模式的核心端口差异

    维度 传统主从模式 DLedger 模式
    主从同步核心端口 10909(HA 复制) 40911(Raft 协议同步)
    10909 端口状态 主节点监听、从节点连接 所有节点均不监听、无交互
    40911 端口状态 所有节点监听、集群内互通
    10911/9876 端口 功能一致 功能完全复用

4.2、配置说明

  • 基础集群配置
参数名 说明 生产建议值
brokerClusterName 集群名称(所有 Broker/NameServer 统一) MyRocketMQCluster
brokerName Broker 组名(DLedger 集群内所有节点同名) broker-a
brokerId Broker 类型(DLedger 模式需注释,由 Raft 自动管理) 注释 / 删除(传统模式:主 0 / 从 1+)
namesrvAddr NameServer 地址(多个用分号分隔) 地址1:9876;地址2:9876;地址2:9876
listenPort Broker 客户端通信端口(所有节点可复用,IP 不同即可) 10911
autoCreateTopicEnable 是否允许自动创建 Topic false(生产)/true(测试)
  • DLedger 核心配置(专属)

    参数名 说明 生产建议值
    enableDLegerCommitLog 开启 DLedger 模式(核心开关) true
    dLegerGroup DLedger 集群组名(所有节点统一) broker-a-group
    dLegerPeers DLedger 集群节点(节点 ID-IP:40911,分号分隔) n0-地址1:40911;n1-地址2:40911;n2-地址3:40911
    dLegerSelfId 当前节点 DLedger 唯一 ID(与 peers 中 ID 一一对应) 130 节点:n0
    139 节点:n1
    167 节点:n2
    haEnable 关闭传统 HA 主从复制(避免与 DLedger 冲突) false
    brokerRole Broker 角色(DLedger 模式仅支持主模式) ASYNC_MASTER(性能)/SYNC_MASTER(数据安全)
    • dLegerSelfId 绝对不可以设置为相同值
      • dLegerSelfId 是 DLedger 集群中每个节点的唯一标识,必须与 dLegerPeers 中定义的节点 ID(n0/n1/n2)一一对应,且 3 个节点的 dLegerSelfId 必须完全不同:
        • 10.4.50.130 节点:dLegerSelfId=n0
        • 10.4.50.139 节点:dLegerSelfId=n1
        • 10.4.50.167 节点:dLegerSelfId=n2
    • 若设置相同的后果
      1. 集群无法启动:节点启动后会提示「duplicate selfId in peers」,或日志中出现「peer id conflict」,无法加入集群;
      2. Leader 选举失败:即使强制启动,集群无法选举出有效 Leader,表现为所有节点均为 Follower,无可用的 Leader 提供服务;
      3. 数据不一致:若侥幸启动,节点间会因 ID 冲突导致数据同步错乱,出现消息丢失、重复消费等严重问题;
      4. 端口冲突(间接):虽然 40911 端口可复用(IP 不同),但相同 selfId 会导致节点误判为 “同一节点重复注册”,关闭端口监听。
  • 存储 / 日志配置

    参数名 说明 生产建议值
    flushDiskType 刷盘策略 ASYNC_FLUSH(性能)
    SYNC_FLUSH(金融级)
    logBaseDir 运行日志存储目录 /data/rocketmq/logs
    storePathRootDir 核心数据存储根目录 /data/rocketmq/store
    storePathDLedgerCommitLog DLedger 日志存储路径 /data/rocketmq/store/dledger/commitlog
    fileReservedTime 消息文件保留时长(小时) 72/168(根据业务调整)
    deleteWhen 过期消息清理时间(24 小时制) 04(凌晨 4 点,低峰期)
  • 性能调优配置

    参数名 说明 生产建议值
    sendMessageThreadPoolNums 消息发送线程池数量 16-32(高并发可设 64)
    consumeMessageThreadPoolNums 消息消费线程池数量 32-64(高并发可设 128)

4.3、可视化查看

  • 下载RocketMQ Dashboard

  • 我这里用:rocketmq-dashboard-rocketmq-dashboard-1.0.0

  • 本机windows需要用到: mvnjdk1.8 ,用ai查一下怎么安装,还需要配上加速站点,不然会很慢

  • 修改配置文件

    路径: rocketmq-dashboard-rocketmq-dashboard-1.0.0\src\main\resources
    用文本打开:  application.properties
    
    # 集群配置
    rocketmq.config.namesrvAddr=10.4.50.130:9876;10.4.50.139:9876;10.4.50.167:9876
    # 单机配置
    # rocketmq.config.namesrvAddr=10.4.50.163:9876
    
    # 回到  rocketmq-dashboard-rocketmq-dashboard-1.0.0  这一级
    按shift+右键,打开powershell窗口
    
    # 构建包, 默认端口是8080
    mvn clean package
    java -jar target/rocketmq-dashboard-1.0.1-SNAPSHOT.jar
    
    # 直接运行
    mvn spring-boot:run
    
  • 运行查看

    # 输出日志
    C:\Users\eetrust\Downloads\rocketmq-dashboard-rocketmq-dashboard-1.0.0\rocketmq-dashboard-rocketmq-dashboard-1.0.0>set JAVA_HOME=C:\Program Files\Java\jdk1.8.0_202
    [2025-12-08 17:48:35.618]  INFO setNameSrvAddrByProperty nameSrvAddr=10.4.50.130:9876;10.4.50.139:9876;10.4.50.167:9876
    [2025-12-08 17:48:37.857]  INFO Tomcat initialized with port(s): 18080 (http)
    [2025-12-08 17:48:37.873]  INFO Initializing ProtocolHandler ["http-nio-0.0.0.0-18080"]
    [2025-12-08 17:48:37.874]  INFO Starting service [Tomcat]
    [2025-12-08 17:48:37.876]  INFO Starting Servlet engine: [Apache Tomcat/9.0.29]
    

请添加图片描述

Logo

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

更多推荐