RocketMQ发展历史及5.0特性介绍:https://zhuanlan.zhihu.com/p/420102329

官网文档: https://rocketmq.apache.org/zh/docs/quickStart/02quickstart

git: https://github.com/apache/rocketmq

阿里云手册: https://help.aliyun.com/document_detail/444755.html

https://github.com/apache/rocketmq/tree/master/docs/cn

消息队列RocketMQ版基于统一消息存储和轻量计算层,主要应用于微服务异步解耦、流式数据处理、事件驱动等场景。

整体的架构设计主要分为四大部分,分别是:Producer、Consumer、Broker、NameServer。

为了更贴合实际,我画的都是集群部署,像 Broker 我还画了主从。

  • Producer:就是消息生产者,可以集群部署。它会先和 NameServer 集群中的随机一台建立长连接,得知当前要发送的 Topic 存在哪台 Broker Master上,然后再与其建立长连接,支持多种负载平衡模式发送消息。

  • Consumer:消息消费者,也可以集群部署。它也会先和 NameServer 集群中的随机一台建立长连接,得知当前要消息的 Topic 存在哪台 Broker Master、Slave上,然后它们建立长连接,支持集群消费和广播消费消息。

  • Broker:主要负责消息的存储、查询消费,支持主从部署,一个 Master 可以对应多个 Slave,Master 支持读写,Slave 只支持读。Broker 会向集群中的每一台 NameServer 注册自己的路由信息。

  • NameServer:是一个很简单的 Topic 路由注册中心,支持 Broker 的动态注册和发现,保存 Topic 和 Borker 之间的关系。通常也是集群部署,但是各 NameServer 之间不会互相通信, 各 NameServer 都有完整的路由信息,即无状态。

我再用一段话来概括它们之间的交互:

先启动 NameServer 集群,各 NameServer 之间无任何数据交互,Broker 启动之后会向所有 NameServer 定期(每 30s)发送心跳包,包括:IP、Port、TopicInfo,NameServer 会定期扫描 Broker 存活列表,如果超过 120s 没有心跳则移除此 Broker 相关信息,代表下线。

这样每个 NameServer 就知道集群所有 Broker 的相关信息,此时 Producer 上线从 NameServer 就可以得知它要发送的某 Topic 消息在哪个 Broker 上,和对应的 Broker (Master 角色的)建立长连接,发送消息。

Consumer 上线也可以从 NameServer 得知它所要接收的 Topic 是哪个 Broker ,和对应的 Master、Slave 建立连接,接收消息。

简单的工作流程如上所述,相信大家对整体数据流转已经有点印象了,我们再来看看每个部分的详细情况。

NameServer

它的特点就是轻量级,无状态。角色类似于 Zookeeper 的情况,从上面描述知道其主要的两个功能就是:Broker 管理、路由信息管理。

总体而言比较简单,我再贴一些字段,让大家有更直观的印象知道它存储了些什么。

Producer

Producer 无非就是消息生产者,那首先它得知道消息要发往哪个 Broker ,于是每 30s 会从某台 NameServer 获取 Topic 和 Broker 的映射关系存在本地内存中,如果发现新的 Broker 就会和其建立长连接,每 30s 会发送心跳至 Broker 维护连接。

并且会轮询当前可以发送的 Broker 来发送消息,达到负载均衡的目的,在同步发送情况下如果发送失败会默认重投两次(retryTimesWhenSendFailed = 2),并且不会选择上次失败的 broker,会向其他 broker 投递。

异步发送失败的情况下也会重试,默认也是两次 (retryTimesWhenSendAsyncFailed = 2),但是仅在同一个 Broker 上重试。

Producer 启动流程

然后我们再来看看 Producer 的启动流程看看都干了些啥。

大致启动流程图中已经表明的很清晰的,但是有些细节可能还不清楚,比如重平衡啊,TBW102 啥玩意啊,有哪些定时任务啊,别急都会提到的。

有人可能会问这生产者为什么要启拉取服务、重平衡?

因为 Producer 和 Consumer 都需要用 MQClientInstance,而同一个 clientId 是共用一个 MQClientInstance 的, clientId 是通过本机 IP 和 instanceName(默认值 default)拼起来的,所以多个 Producer 、Consumer 实际用的是一个MQClientInstance。

至于有哪些定时任务,请看下图:

Producer 发消息流程

我们再来看看发消息的流程,大致也不是很复杂,无非就是找到要发送消息的 Topic 在哪个 Broker 上,然后发送消息。

现在就知道 TBW102 是啥用的,就是接受自动创建主题的 Broker 启动会把这个默认主题登记到 NameServer,这样当 Producer 发送新 Topic 的消息时候就得知哪个 Broker 可以自动创建主题,然后发往那个 Broker。

而 Broker 接受到这个消息的时候发现没找到对应的主题,但是它接受创建新主题,这样就会创建对应的 Topic 路由信息。

Producer 每 30s 会向 NameSrv 拉取路由信息更新本地路由表,有新的 Broker 就和其建立长连接,每隔 30s 发送心跳给 Broker 。

不要在生产环境开启 autoCreateTopicEnable。

Producer 会通过重试和延迟机制提升消息发送的高可用。

Broker

Broker 就比较复杂一些了,但是非常重要。大致分为以下五大模块,我们来看一下官网的图。

  • Remoting 远程模块,处理客户请求。

  • Client Manager 管理客户端,维护订阅的主题。

  • Store Service 提供消息存储查询服务。

  • HA Serivce,主从同步高可用。

  • Index Serivce,通过指定key 建立索引,便于查询。

有几个模块没啥可说的就不分析了,先看看存储的。

下载安装

https://www.apache.org/dyn/closer.cgi?path=rocketmq/5.0.0/rocketmq-all-5.0.0-bin-release.zip

MAC安装

下载zip

上官网下载RocketMQ(Binary:rocketmq-all-4.4.0-bin-release.zip):

http://rocketmq.apache.org/rel

解压

下载下来是一个zip的文件rocketmq-all-4.4.0-bin-release.zip

先mac下解压很简单,双击就可以解压了,命令解压unzip:

#unzip rocketmq-all-4.4.0-bin-release.zip

配置环境变量

export ROCKETMQ_HOME=/app/rocketmq/rocketmq-all-4.9.1-bin-release

export = $ROCKETMQ_HOME\bin:$PATH

需要配置jdk,

默认支持1.8, JDK1.8以上需要单独修改配置:

  1. runserver.sh修改如下:

#修改CLASSPATH参数为
export CLASSPATH=.:${BASE_DIR}/conf:${JAVA_HOME}/jre/lib/ext:${BASE_DIR}/lib/*
#JDK版本判断由如下
if  [ -z "$JAVA_MAJOR_VERSION" || "$JAVA_MAJOR_VERSION" -lt "9" ]
#JDK版本判断改为:
if  [  "$JAVA_MAJOR_VERSION" -lt "9" ]
  1. runbroker.sh修改如下:

#修改CLASSPATH参数为
export CLASSPATH=.${JAVA_HOME}/jre/lib/ext:${BASE_DIR}/lib/*:${BASE_DIR}/conf:${CLASSPATH}
#JDK版本判断由如下
if  [ -z "$JAVA_MAJOR_VERSION" || "$JAVA_MAJOR_VERSION" -lt "9" ]
#JDK版本判断改为:
if  [  "$JAVA_MAJOR_VERSION" -lt "9" ]
  1. tools.sh修改如下:

#修改CLASSPATH参数为
export CLASSPATH=.${JAVA_HOME}/jre/lib/ext:${BASE_DIR}/lib/*:${BASE_DIR}/conf:${CLASSPATH}

启动nameserver

https://github.com/apache/rocketmq/blob/develop/docs/cn/operation.md

进入解压目录:/rocketmq-all-4.5.1-bin-release/bin,输入指令nohup sh bin/mqnamesrv &

### 首先启动Name Server
$ nohup sh mqnamesrv &
 
### 验证Name Server 是否启动成功
$ tail -f ~/logs/rocketmqlogs/namesrv.log
The Name Server boot success...

启动broker

./mqbroker -n localhost:9876
The broker[node1, 172.17.0.1:10911] boot success. serializeType=JSON and name server is localhost:9876
### 启动Broker
$ nohup sh bin/mqbroker -n localhost:9876 &

### 验证Broker是否启动成功,例如Broker的IP为:192.168.1.2,且名称为broker-a
$ tail -f ~/logs/rocketmqlogs/broker.log 
The broker[broker-a, 192.169.1.2:10911] boot success...

关闭消息队列

通过mqshutdown命令关闭消息队列,依次关闭nameserver和broker

sh mqshutdown broker
The mqbroker(69255) is running...
Send shutdown request to mqbroker(69255) OK
sh mqshutdown namesrv
The mqnamesrv(68711) is running...
Send shutdown request to mqnamesrv(68711) OK

mqadmin管理工具

注意:
1.执行命令方法:./mqadmin {command} {args}
2.几乎所有命令都需要配置-n表示NameServer地址,格式为ip:port
3.几乎所有命令都可以通过-h获取帮助
4.如果既有Broker地址(-b)配置项又有clusterName(-c)配置项,则优先以Broker地址执行命令,如果不配置Broker地址,则对集群中所有主机执行命令,只支持一个Broker地址。-b格式为ip:port,port默认是10911
5.在tools下可以看到很多命令,但并不是所有命令都能使用,只有在MQAdminStartup中初始化的命令才能使用,你也可以修改这个类,增加或自定义命令
6.由于版本更新问题,少部分命令可能未及时更新,遇到错误请直接阅读相关命令源码

Windows安装

说明

这里的环境是window 10、jdk1.8(jdk环境需要提前配置好)

下载zip

上官网下载RocketMQ(Binary:rocketmq-all-4.4.0-bin-release.zip):

http://rocketmq.apache.org/release_notes/release-notes-4.4.0/

解压

下载下来是一个zip的文件rocketmq-all-4.4.0-bin-release.zip使用解压工具解压即可。

系统环境变量配置

这个是必须的,否则会在启动的时候,会提示:

Please set the ROCKETMQ_HOMEvariable in your environment!

环境变量配置:

变量名(固定值):ROCKETMQ_HOME

变量值(和你存放的路径有关):MQ解压路径\MQ文件夹名

举例说明:

启动

进入到rocketMQ的bin目录下,然后执行下面的命令:

#启动RocketMQ的注册中心

#jdk1.8以上,修改runserver.cmd中jvm大小和GC垃圾回收器和日志配置

start mqnamesrv.cmd

执行成功会弹出如下的提示框,不要关闭:

#启动broker

##jdk1.8以上,修改runbroker.cmd中jvm大小和GC垃圾回收器和日志配置

start mqbroker.cmd -n127.0.0.1:9876 autoCreateTopicEnable=true

成功会弹出如下的提示框,不要关闭:

docker 部署

#拉取镜像
docker pull apache/rocketmq:5.1.4
#主机创建日志目录
mkdir -p /docker/rocketmq/nameserver/logs /docker/rocketmq/nameserver/store

#运行nameserver
docker run -d 
    --restart=always #自动启动
    --name rmqnamesrv #名字
    --privileged=true  #私有
    -p 9876:9876  #端口号
    -v /docker/rocketmq/nameserver/logs:/root/logs #挂卷
    -v /docker/rocketmq/nameserver/store:/root/store 
    -e "MAX_POSSIBLE_HEAP=100000000" #环境变量
apache/rocketmq
sh mqnamesrv #启动后命令

#运行broker
docker run -d  \
--restart=always \
--name rmqbroker \
--link rmqnamesrv:namesrv \
-p 10911:10911 \
-p 10909:10909 \
-v  /data/rocketmq/data/broker/logs:/home/rocketmq/logs/rocketmqlogs \
-v  /data/rocketmq/data/broker/store:/home/rocketmq/store \
-v /data/rocketmq/conf/broker.conf:/home/rocketmq/rocketmq-5.1.0/conf/broker.conf \
-e "NAMESRV_ADDR=namesrv:9876" \
-e "MAX_POSSIBLE_HEAP=200000000" \
apache/rocketmq \
sh mqbroker -c /home/rocketmq/rocketmq-5.1.0/conf/broker.conf

#运行控制台
docker run -d \
--restart=always \
--name rocketmq-admin \
-e "JAVA_OPTS=-Drocketmq.namesrv.addr=10.0.102.43:9876 \
-Dcom.rocketmq.sendMessageWithVIPChannel=false" \
-p 9090:8080 \
apacherocketmq/rocketmq-dashboard:latest
services:
  namesrv:
    image: apache/rocketmq:5.0.0
    container_name: rmqnamesrv
    ports:
      - 9876:9876
    networks:
      - rocketmq
    volumes:
      - /etc/localtime:/etc/localtime
    command: sh mqnamesrv

  broker:
    image: apache/rocketmq:5.0.0
    container_name: rmqbroker
    ports:
      - 10909:10909
      - 10911:10911
      - 10912:10912
    environment:
      - NAMESRV_ADDR=rmqnamesrv:9876
    volumes:
      - /home/deploy-docker/volumes/rocketmq/broker/conf/broker.conf:/home/rocketmq/rocketmq-5.0.0/conf/broker.conf
      - /etc/localtime:/etc/localtime
    depends_on:
      - namesrv
    networks:
      - rocketmq
    command: sh mqbroker -c /home/rocketmq/rocketmq-5.0.0/conf/broker.conf

networks:
  rocketmq:
    driver: bridge

docker部署脚本2

#Docker镜像构建

git clone https://github.com/apache/rocketmq-docker.git

cd rocketmq-docker/

cd image-build/

sh build-image.sh 5.2.0 alpine

docker save -o rocketmq-5.2.0.tar  apache/rocketmq:5.2.0-alpine
docker load -i rocketmq-5.2.0.tar

单节点安装

mkdir -p data/broker/logs data/broker/store data/namesrv/logs
chmod -R 777 data
tee broker.conf <<-'EOF'
brokerClusterName = DefaultCluster
brokerName = broker-a
brokerId = 0
deleteWhen = 04
fileReservedTime = 48
brokerRole = ASYNC_MASTER
flushDiskType = ASYNC_FLUSH
autoCreateTopicEnable=true
autoCreateSubscriptionGroup=true
enablePropertyFilter=true
aclEnable=true     #若不开启认证功能则设置false
listenPort=10911
brokerIP1=192.168.0.100   # 替换当前服务器ip
EOF
tee plain_acl.yml <<-'EOF'
accounts:
 - accessKey: 556230f2da1bf9d18af7124b48f99a72ce5a7589
   secretKey: e8837b9049805d99d9f45061218ed65d8d20b755
   admin: true
EOF
tee docker-compose.yml <<-'EOF'
name: rocketmq

services:
  rmqnamesrv:
    image: apache/rocketmq:5.2.0-alpine
    container_name: rmqnamesrv
    hostname: rmqnamesrv
    restart: always
    ulimits:
      nproc: 655350
      nofile:
       soft: 655350
       hard: 655350
    environment:
    - JAVA_OPT_EXT=-server -Xms2g -Xmx2g -Xmn1g
    ports:
    - 9876:9876
    volumes:
    - ./data/namesrv/logs:/home/rocketmq/logs
    command: sh mqnamesrv

  rmqbroker:
    image: apache/rocketmq:5.2.0-alpine
    container_name: rmqbroker
    hostname: rmqbroker
    restart: always
    ulimits:
      nproc: 655350
      nofile:
       soft: 655350
       hard: 655350
    links:
    - rmqnamesrv
    ports:
    - "10908-10914:10908-10914"
    environment:
    - NAMESRV_ADDR=rmqnamesrv:9876
    - JAVA_OPT_EXT=-server -Xms4g -Xmx4g
    volumes:
    - ./data/broker/logs:/home/rocketmq/logs
    - ./data/broker/store:/home/rocketmq/store
    - ./broker.conf:/opt/rocketmq-5.2.0/conf/broker.conf
    - ./plain_acl.yml:/home/rocketmq/rocketmq-5.2.0/conf/plain_acl.yml
    command: sh mqbroker -c /opt/rocketmq-5.2.0/conf/broker.conf
    depends_on:
    - rmqnamesrv
EOF

集群搭建案例查看

各端口说明:

  • 8081 :proxy默认端口

  • 9876 :nameserver 默认端口,listenport=自定义端口

  • 9878 :controller默认端口,

  • 10911:broker默认端口,同一服务器内不冲突。listenport=自定义端口

  • 10909 : broker的fastRemotingServer服务组件用于slave同步master,默认为listenPort - 2. fastListenPort=自定义端口

  • 10912:broker的HAService服务组件使用,主从同步,默认为listenPort + 1,haListenPort=自定义端口

Proxy模块

RocketMQ 5.0 把客户端的部分功能下沉到 Proxy,Proxy 承接了之前 客户端的计算能力,客户端变得更加轻量级。生产者和消费者不再需要连接nameserver和获取计算topic路由,都由proxy进行交互和计算,生产者和消费者只管生产和消费数据。

客户端所有的请求都要经过 Proxy,Proxy 将流量分发到 Broker。这样在 Proxy 可以进行流量控制和流量治理。

Proxy 有两种部署方式,LOCAL 模式和 CLUSTER 模式。

  • 在 Local 模式下,Broker 和 Proxy 是同进程部署,只是在原有 Broker 的配置基础上新增 Proxy 的简易配置就可以运行。

  • 在 Cluster 模式下,Broker 和 Proxy 分别部署,即在原有的集群基础上,额外再部署 Proxy 即可。

Local模式部署

启动 NameServer

NameServer需要先于Broker启动,且如果在生产环境使用,为了保证高可用,建议一般规模的集群启动3个NameServer,各节点的启动命令相同,如下:

### 首先启动Name Server
$ nohup sh mqnamesrv & 
### 验证Name Server 是否启动成功
$ tail -f ~/logs/rocketmqlogs/namesrv.log
The Name Server boot success...

启动Broker+Proxy

单组节点单副本模式

警告

这种方式风险较大,因为 Broker 只有一个节点,一旦Broker重启或者宕机时,会导致整个服务不可用。不建议线上环境使用, 可以用于本地测试。

启动 Broker+Proxy

$ nohup sh bin/mqbroker -n localhost:9876 --enable-proxy &
### 验证Broker 是否启动成功,例如Broker的IP为:192.168.1.2,且名称为broker-a
$ tail -f ~/logs/rocketmqlogs/broker_default.log
The broker[xxx, 192.169.1.2:10911] boot success...

多组节点(集群)单副本模式

一个集群内全部部署 Master 角色,不部署Slave 副本,例如2个Master或者3个Master,这种模式的优缺点如下:

  • 优点:配置简单,单个Master宕机或重启维护对应用无影响,在磁盘配置为RAID10时,即使机器宕机不可恢复情况下,由于RAID10磁盘非常可靠,消息也不会丢(异步刷盘丢失少量消息,同步刷盘一条不丢),性能最高;

  • 缺点:单台机器宕机期间,这台机器上未被消费的消息在机器恢复之前不可订阅,消息实时性会受到影响。

启动Broker+Proxy集群

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-noslave/broker-a.properties --enable-proxy & 

### 在机器B,启动第二个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-noslave/broker-b.properties --enable-proxy &
...

如上启动命令是在单个NameServer情况下使用的。对于多个NameServer的集群,Broker启动命令中-n后面的地址列表用分号隔开即可,例如 192.168.1.1:9876;192.161.2:9876

多节点(集群)多副本模式-异步复制

每个Master配置一个Slave,有多组 Master-Slave,HA采用异步复制方式,主备有短暂消息延迟(毫秒级),这种模式的优缺点如下:

  • 优点:即使磁盘损坏,消息丢失的非常少,且消息实时性不会受影响,同时Master宕机后,消费者仍然可以从Slave消费,而且此过程对应用透明,不需要人工干预,性能同多Master模式几乎一样;

  • 缺点:Master宕机,磁盘损坏情况下会丢失少量消息。

启动Broker+Proxy集群

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-a.properties --enable-proxy & 

### 在机器B,启动第二个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-b.properties --enable-proxy & 

### 在机器C,启动第一个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-a-s.properties --enable-proxy & 

### 在机器D,启动第二个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-b-s.properties --enable-proxy &

多节点(集群)多副本模式-同步双写

每个Master配置一个Slave,有多对 Master-Slave,HA采用同步双写方式,即只有主备都写成功,才向应用返回成功,这种模式的优缺点如下:

  • 优点:数据与服务都无单点故障,Master宕机情况下,消息无延迟,服务可用性与数据可用性都非常高;

  • 缺点:性能比异步复制模式略低(大约低10%左右),发送单个消息的RT会略高,且目前版本在主节点宕机后,备机不能自动切换为主机。

启动 Broker+Proxy 集群

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-a.properties --enable-proxy & 

### 在机器B,启动第二个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-b.properties --enable-proxy & 

### 在机器C,启动第一个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-a-s.properties --enable-proxy & 

### 在机器D,启动第二个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-b-s.properties --enable-proxy &

以上 Broker 与 Slave 配对是通过指定相同的 BrokerName 参数来配对,Master 的 BrokerId 必须是 0,Slave 的 BrokerId 必须是大于 0 的数。另外一个 Master 下面可以挂载多个 Slave,同一 Master 下的多个 Slave 通过指定不同的 BrokerId 来区分。$ROCKETMQ_HOME指的RocketMQ安装目录,需要用户自己设置此环境变量。

Cluster模式部署

在 Cluster模式下,一个 Proxy集群和 Broker集群为一一对应的关系,可以在 Proxy的配置文件 rmq-proxy.json 中使用 rocketMQClusterName 进行配置

启动 NameServer

### 首先启动Name Server
$ nohup sh mqnamesrv & 
### 验证Name Server 是否启动成功
$ tail -f ~/logs/rocketmqlogs/namesrv.log
The Name Server boot success...

启动 Broker

单组节点单副本模式

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 &

多组节点(集群)单副本模式

一个集群内全部部署 Master 角色,不部署Slave 副本,例如2个Master或者3个Master,这种模式的优缺点如下:

  • 优点:配置简单,单个Master宕机或重启维护对应用无影响,在磁盘配置为RAID10时,即使机器宕机不可恢复情况下,由于RAID10磁盘非常可靠,消息也不会丢(异步刷盘丢失少量消息,同步刷盘一条不丢),性能最高;

  • 缺点:单台机器宕机期间,这台机器上未被消费的消息在机器恢复之前不可订阅,消息实时性会受到影响。

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-noslave/broker-a.properties & 

### 在机器B,启动第二个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-noslave/broker-b.properties &
...

备注
如上启动命令是在单个NameServer情况下使用的。对于多个NameServer的集群,Broker启动命令中-n后面的地址列表用分号隔开即可,例如 192.168.1.1:9876;192.161.2:9876

多节点(集群)多副本模式-异步复制

每个Master配置一个Slave,有多组 Master-Slave,HA采用异步复制方式,主备有短暂消息延迟(毫秒级),这种模式的优缺点如下:

  • 优点:即使磁盘损坏,消息丢失的非常少,且消息实时性不会受影响,同时Master宕机后,消费者仍然可以从Slave消费,而且此过程对应用透明,不需要人工干预,性能同多Master模式几乎一样;

  • 缺点:Master宕机,磁盘损坏情况下会丢失少量消息。

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-a.properties & 

### 在机器B,启动第二个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-b.properties & 

### 在机器C,启动第一个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-a-s.properties & 

### 在机器D,启动第二个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-async/broker-b-s.properties &

多节点(集群)多副本模式-同步双写

每个Master配置一个Slave,有多对 Master-Slave,HA采用同步双写方式,即只有主备都写成功,才向应用返回成功,这种模式的优缺点如下:

  • 优点:数据与服务都无单点故障,Master宕机情况下,消息无延迟,服务可用性与数据可用性都非常高;

  • 缺点:性能比异步复制模式略低(大约低10%左右),发送单个消息的RT会略高,且目前版本在主节点宕机后,备机不能自动切换为主机。

### 在机器A,启动第一个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-a.properties & 

### 在机器B,启动第二个Master,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-b.properties & 

### 在机器C,启动第一个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-a-s.properties & 

### 在机器D,启动第二个Slave,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqbroker -n 192.168.1.1:9876 -c $ROCKETMQ_HOME/conf/2m-2s-sync/broker-b-s.properties &

以上 Broker 与 Slave 配对是通过指定相同的 BrokerName 参数来配对,Master 的 BrokerId 必须是 0,Slave 的 BrokerId 必须是大于 0 的数。另外一个 Master 下面可以挂载多个 Slave,同一 Master 下的多个 Slave 通过指定不同的 BrokerId 来区分。$ROCKETMQ_HOME指的RocketMQ安装目录,需要用户自己设置此环境变量。

启动 Proxy

可以在多台机器启动多个Proxy

### 在机器A,启动第一个Proxy,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqproxy -n 192.168.1.1:9876 &
### 在机器B,启动第二个Proxy,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqproxy -n 192.168.1.1:9876 &
### 在机器C,启动第三个Proxy,例如NameServer的IP为:192.168.1.1
$ nohup sh bin/mqproxy -n 192.168.1.1:9876 &

若需要指定配置文件,可以使用 -pc或者 --proxyConfigPath 进行指定

### 自定义配置文件
$ nohup sh bin/mqproxy -n 192.168.1.1:9876 -pc /path/to/proxyConfig.json &

Controller模块

自动主从切换的 Controller 组件,其可以独立部署也可以内嵌在 NameServer 中。4.5版本增加DLedger模式,使用raft分布式一致性算法。5.0升级为controller组件模块。

Controller 组件提供选主能力,若需要保证 Controller 具备容错能力,Controller 部署需要三副本及以上(遵循 Raft 的多数派协议)。

Controller 部署有两种方式。一种是嵌入于 NameServer 进行部署,可以通过配置 enableControllerInNamesrv 打开(可以选择性打开,并不强制要求每一台 NameServer 都打开),在该模式下,NameServer 本身能力仍然是无状态的,也就是内嵌模式下若 NameServer 挂掉多数派,只影响切换能力,不影响原来路由获取等功能。另一种是独立部署,需要单独部署 Controller 组件。

Controller 嵌入 NameServer 部署

嵌入 NameServer 部署时只需要在 NameServer 的配置文件中设置 enableControllerInNamesrv=true,并填上 Controller 的配置即可。

#namesrv.conf配置文件
#Nameserver 中是否开启 controller,默认 false。
enableControllerInNamesrv = true
#DLedger Raft Group 的名字,同一个 DLedger Raft Group 保持一致即可。
controllerDLegerGroup = group1
#DLedger Group 内各节点的端口信息,同一个 Group 内的各个节点配置必须要保证一致。
controllerDLegerPeers = n0-127.0.0.1:9877;n1-127.0.0.1:9878;n2-127.0.0.1:9879
#节点 id,必须属于 controllerDLegerPeers 中的一个;同 Group 内各个节点要唯一。
controllerDLegerSelfId = n0
#controller 日志存储位置。controller 是有状态的,controller 重启或宕机需要依靠日志来恢复数据,该目录非常重要,不可以轻易删除。
controllerStorePath = /home/admin/DledgerController
#是否可以从 SyncStateSet 以外选举 Master,若为 true,可能会选取数据落后的副本作为 Master 而丢失消息,默认为 false。
enableElectUncleanMaster = false
#当 Broker 副本组上角色发生变化时是否主动通知,默认为 true。
notifyBrokerRoleChanged = true
#nameserver指定配置文件启动
$ nohup sh bin/mqnamesrv -c namesrv.conf &

Controller 独立部署

独立部署执行以下脚本即可

$ nohup sh bin/mqcontroller -c controller.conf &

mqcontroller 脚本在源码包 distribution/bin/mqcontroller,配置参数与内嵌模式相同。

独立部署Controller后,仍然需要单独部署NameServer提供路由发现能力

Broker 部署

Broker 启动方法与之前相同,增加以下参数

  • enableControllerMode:Broker controller 模式的总开关,只有该值为 true,自动主从切换模式才会打开。默认为 false。

  • controllerAddr:controller 的地址,多个 controller 中间用分号隔开。例如controllerAddr = 127.0.0.1:9877;127.0.0.1:9878;127.0.0.1:9879

  • syncBrokerMetadataPeriod:向 controller 同步 Broker 副本信息的时间间隔。默认 5000(5s)。

  • checkSyncStateSetPeriod:检查 SyncStateSet 的时间间隔,检查 SyncStateSet 可能会 shrink SyncState。默认5000(5s)。

  • syncControllerMetadataPeriod:同步 controller 元数据的时间间隔,主要是获取 active controller 的地址。默认10000(10s)。

  • haMaxTimeSlaveNotCatchup:表示 Slave 没有跟上 Master 的最大时间间隔,若在 SyncStateSet 中的 slave 超过该时间间隔会将其从 SyncStateSet 移除。默认为 15000(15s)。

  • storePathEpochFile:存储 epoch 文件的位置。epoch 文件非常重要,不可以随意删除。默认在 store 目录下。

  • allAckInSyncStateSet:若该值为 true,则一条消息需要复制到 SyncStateSet 中的每一个副本才会向客户端返回成功,可以保证消息不丢失。默认为 false。

  • syncFromLastFile:若 slave 是空盘启动,是否从最后一个文件进行复制。默认为 false。

  • asyncLearner:若该值为 true,则该副本不会进入 SyncStateSet,也就是不会被选举成 Master,而是一直作为一个 learner 副本进行异步复制。默认为false。

  • inSyncReplicas:需保持同步的副本组数量,默认为1,allAckInSyncStateSet=true 时该参数无效。

  • minInSyncReplicas:最小需保持同步的副本组数量,若 SyncStateSet 中副本个数小于 minInSyncReplicas 则 putMessage 直接返回 PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH,默认为1。

在Controller模式下,Broker配置必须设置 enableControllerMode=true,并填写 controllerAddr,并以下面命令启动:

$ nohup sh bin/mqbroker -c broker.conf &

注意

自动主备切换模式下Broker无需指定brokerId和brokerRole,其由Controller组件进行分配

自定义启动脚本

#自定义sh脚本
#!/bin/bash

#软件安装目录设置,先读环境变量,读不到则用默认值
ROCKETMQ_HOME=$ROCKETMQ_HOME
if [ -z $ROCKETMQ_HOME]; then
    echo "ROCKETMQ_HOME is not set! the program will use default value"
    ROCKETMQ_HOME= /rocketmq-all-4.9.7-bin-release
fi
echo "ROCKETMQ_HOME=$ROCKETMQ_HOME"


# 启动rocketmq nameserver
cd ${ROCKETMQ_HOME}/bin
# 检查Zookeeper是否启动成功
if jps -ml | grep "namesrv.NamesrvStartup"; then
    echo "rocketmq nameserver is running.";
else
    echo "rocketmq nameserver is Starting ...";
    rm -rf nohup.out;
    nohup sh mqnamesrv &
    count=0
    while [ $count -le 10 ]; do
      if tail -100 nohup.out | grep "Name Server boot success"; then
        echo "rocketmq nameserver has been started successfully";
        break
      else
        count=$((count+1))
        sleep 2
      fi
    done
fi

# 启动rocketmq broker
# 检查Zookeeper是否启动成功
if jps -ml | grep "org.apache.rocketmq.broker.BrokerStartup"; then
    echo "rocketmq broker is running.";
else
    echo "rocketmq broker is Starting ...";
    rm -rf nohup.out;
    nohup sh mqbroker -c ../conf/broker.conf -n localhost:9876 autoCreateTopicEnable=true &  
    count=0
    while [ $count -le 10 ]; do
      if tail -100 nohup.out | grep "broker.*success"; then
        echo "rocketmq broker has been started successfully";
        break
      else
        count=$((count+1))
        sleep 2
      fi
    done
fi

# 启动rocketmq dashboard
cd ${ROCKETMQ_HOME}/dashboard
if jps -ml | grep "rocketmq-dashboard"; then
    echo "rocketmq dashboard is running.";
else
    echo "rocketmq dashboard is Starting ...";
    rm -rf nohup.out;
    nohup java -jar -Drocketmq.namesrv.addr=localhost:9876 rocketmq-dashboard-1.0.1.jar & 
    count=0
    while [ $count -le 10 ]; do
      if tail -100 nohup.out | grep "Tomcat started on port(s)"; then
        echo "rocketmq dashboard has been started successfully";
        break
      else
        count=$((count+1))
        sleep 2
      fi
    done
fi

echo "All services started.";

RocketMQ配置参数说明

Nameserver配置

配置项

说明

serverChannelMaxIdleTimeSeconds

120

网络连接最大空闲时间。如果链接空闲时间超过此参数设置的值,连接将被关闭

listenPort

9876

默认监听端口

serverCallbackExecutorThreads

0

netty public任务线程池个数,netty网络设计没根据业务类型会创建不同线程池毛笔如处理发送消息,消息消费心跳检测等。如果业务类型(RequestCode)未注册线程池,则由public线程池执行

serverAsyncSemaphoreValue

64

异步消息发送最大并发度

serverSocketSndBufSize

4096

netty网络socket发送缓存区大小

rocketmqHome

/usr/local/biyao/apache-rocketmq

RockerMQ主目录,默认用户主目录

clusterTest

FALSE

是否开启集群测试,默认为false

serverSelectorThreads

3

IO线程池线程个数,主要是NameServer.broker端解析请求,返回相应的线程个数,这类线程主要是处理网络请求的,解析请求包。然后转发到各个业务线程池完成具体的业务无操作,然后将结果在返回调用方

useEpollNativeSelector

FALSE

是否启用Epoll IO模型。Linux环境建议开启

orderMessageEnable

FALSE

是否支持顺序消息,默认为false

serverPooledByteBufAllocatorEnable

TRUE

ByteBuffer是否开启缓存

kvConfigPath

/home/biyao/namesrv/kvConfig.json

kv配置文件路径,包含顺序消息主题的配置信息

serverWorkerThreads

8

netty业务线程池个数

serverSocketRcvBufSize

4096

netty网络socket接收缓存区大小

productEnvName

center

默认生产环境名称

serverOnewaySemaphoreValue

256

send oneway消息请求并发度

configStorePath

/home/biyao/namesrv/namesrv.properties

NameServer配置文件路径,建议使用-c指定NameServer配置文件路径

Broker配置

#所属集群名字
brokerClusterName=rocketmq-cluster
#broker名字,注意此处不同的配置文件填写的不一样
brokerName=broker-a
#0 表示Master, > 0 表示slave
brokerId=0
#nameServer 地址,分号分割
namesrvAddr=192.168.0.75:9876
#在发送消息时,自动创建服务器不存在的Topic,默认创建的队列数
defaultTopicQueueNums=4
#是否允许Broker 自动创建Topic,建议线下开启,线上关闭
autoCreateTopicEnable=true
#是否允许Broker自动创建订阅组,建议线下开启,线上关闭
autoCreateSubscriptionGroup=true
#Broker 对外服务的监听端口
listenPort=10911
#删除文件时间点,默认是凌晨4点
deleteWhen=04
#文件保留时间,默认48小时
fileReservedTime=120
#commitLog每个文件的大小默认1G
mapedFileSizeCommitLog=1073741824
#ConsumeQueue每个文件默认存30W条,根据业务情况调整
mapedFileSizeConsumeQueue=300000
#destroyMapedFileIntervalForcibly=120000
#redeleteHangedFileInterval=120000
#检测物理文件磁盘空间
diskMaxUsedSpaceRatio=88
#存储路径
storePathRootDir=/root/rocketmq/store
#commitLog存储路径
storePathCommitLog=/root/rocketmq/store/commitlog
#消费队列存储路径
storePathConsumeQueue=/root/rocketmq/store/consumequeue
#消息索引存储路径
storePathIndex=/root/rocketmq/store/index
#checkpoint 文件存储路径
storeCheckpoint=/root/rocketmq/store/checkpoint
#abort 文件存储路径
abortFile=/root/rocketmq/store/abort
#限制的消息大小
maxMessageSize=65536
# flushCommitLogLeastPages=4
# flushConsumeQueueLeastPages=2
# flushCommitLogThoroughInterval=10000
# flushConsumeQueueThoroughInterval=60000
# Broker 的角色
# - ASYNC_MASTER 异步复制Master
# - SYNC_MASTER 同步双写Master
# - SLAVE
brokerRole=ASYNC_MASTER
# 刷盘方式
# - ASYNC_FLUSH 异步刷盘
# - SYNC_FLUSH 同步刷盘
flushDiskType=ASYNC_FLUSH
#checkTransactionMessageEnable=false
#发消息线程池数量
#sendMessageTreadPoolNums=128
#拉消息线程池数量
#pullMessageTreadPoolNums=128

参数名

默认值

说明

listenPort

10911

接受客户端连接的监听端口

namesrvAddr

null

nameServer 地址

brokerIP1

网卡的 InetAddress

当前 broker 监听的 IP

brokerIP2

跟 brokerIP1 一样

存在主从 broker 时,如果在 broker 主节点上配置了 brokerIP2 属性,broker 从 节点会连接主节点配置的 brokerIP2 进行同步

brokerName

null

broker 的名称

brokerClusterName

DefaultCluster

本 broker 所属的 Cluser 名称

brokerId

0

broker id, 0 表示 master, 其他的正整数表示 slave

storePathCommitLog

$HOME/store/commitlog/

存储 commit log 的路径

storePathConsumerQueue

$HOME/store/consumequeue/

存储 consume queue 的路径

mappedFileSizeCommitLog

1024 1024 1024(1G)

commit log 的映射文件大小

deleteWhen

04

在每天的什么时间删除已经超过文件保留时间的 commit log

fileReservedTime

72

fileReservedTime 72 以小时计算的文件保留时间

brokerRole

ASYNC_MASTER

SYNC_MASTER/ASYNC_MASTER/SLAVE

flushDiskType

ASYNC_FLUSH

SYNC_FLUSH/ASYNC_FLUSH SYNC_FLUSH 模式下的 broker 保证在收到确认生产 者之前将消息刷盘。ASYNC_FLUSH 模式下的 broker 则利用刷盘一组消息的模 式,可以取得更好的性能。

useEpollNativeSelector

FALSE

是否启用Epoll IO模型。Linux环境建议开启

haSlaveFallbehindMax

268435456

允许从服务器落户的最大偏移字节数,默认为256M。超过该值则表示该Slave不可用

haTransferBatchSize

32768

一次HA主从同步传输的最大字节长度,默认为32K

autoCreateSubscriptionGroup

FALSE

是否自动创建消费组

haListenPort

10912

Master监听端口,从服务器连接该端口,默认为10911+1

clientManagerThreadPoolQueueCapacity

1000000

客户端管理线程池任务队列初始大小

flushCommitLogThoroughInterval

10000

commitlog两次刷盘的最大间隔,如果超过该间隔,将fushCommitLogLeastPages要求直接执行刷盘操作

flushCommitLogLeastPages

4

一次刷盘至少需要脏页的数量,针对commitlog文件

clientCallbackExecutorThreads

32

通信层异步回调线程数

notifyConsumerIdsChangedEnable

TRUE

消费者数量变化后是否立即通知RebalenceService线程,以便马上进行重新负载

cleanResourceInterval

10000

清除过期文件线程调度频率

channelNotActiveInterval

60000

NettyClientConfig,默认为60秒

diskMaxUsedSpaceRatio

88

commitlog目录所在分区的最大使用比例,如果commitlog目录所在的分区使用比例大于该值,则触发过期文件删除

debugLockEnable

FALSE

是否支持 PutMessage Lock锁打印信息

messageDelayLevel

1s 5s 10s 30s 1m 3m 4m 5m 6m 8m 9m10m 20m 30m 1h 2h

延迟队列等级(s=秒,m=分,h=小时

clusterTopicEnable

TRUE

集群名称是否可用在主题使用

messageIndexEnable

TRUE

是否支持消息索引文件

serverPooledByteBufAllocatorEnable

TRUE

ByteBuffer是否开启缓存

shortPollingTimeMills

1000

短轮询等待时间

commercialEnable

TRUE

redeleteHangedFileInterval

120000

重试删除文件间隔,配合destorymapedfileintervalforcibly

flushConsumerOffsetInterval

5000

持久化消息消费进度 consumerOffse.json文件的频率ms

flushCommitLogTimed

FALSE

表示await方法等待FlushIntervalCommitlog,如果为true表示使用Thread.sleep方法等待

maxMessageSize

65536

默认允许的最大消息体默认4M

syncFlushTimeout

5000

同步刷盘超时时间

flushConsumeQueueThoroughInterval

60000

Consume两次刷盘的最大间隔,如果超过该间隔,将忽略

clientChannelMaxIdleTimeSeconds

120

NettyClientConfig,默认为120秒

flushDelayOffsetInterval

10000

延迟队列拉取进度刷盘间隔。默认10s

serverSocketRcvBufSize

131072

netty网络socket接收缓存区大小16MB

maxTransferBytesOnMessageInMemory

262144

一次服务端消息拉取,消息在内存中传输允许的最大传输字节数默认256kb

clientManageThreadPoolNums

32

服务端处理客户端管理(心跳注册取消注册线程数量)

serverChannelMaxIdleTimeSeconds

120

网络连接最大空闲时间。如果链接空闲时间超过此参数设置的值,连接将被关闭

serverCallbackExecutorThreads

0

netty public任务线程池个数,netty网络设计没根据业务类型会创建不同线程池毛笔如处理发送消息,消息消费心跳检测等。如果业务类型(RequestCode)未注册线程池,则由public线程池执行

maxTransferBytesOnMessageInDisk

65536

一次服务消息端消息拉取,消息在磁盘中传输允许的最大字节

pullMessageThreadPoolNums

128

服务端处理消息拉取线程池线程数量 默认为16加上当前操作系统CPU核数的两倍

sendThreadPoolQueueCapacity

10000

消息发送线程池任务队列初始大小

diskFallRecorded

TRUE

是否统计磁盘的使用情况,默认为true

transientStorePoolEnable

FALSE

Commitlog是否开启 transientStorePool机制,默认为 false

disableConsumeIfConsumerReadSlowly

FALSE

如果消费组消息消费堆积是否禁用该消费组继续消费消息

commitCommitLogThoroughInterval

200

Commitlog两次提交的最大间隔,如果超过该间隔,将忽略commitCommitLogLeastPages直接提交

consumerManagerThreadPoolQueueCapacity

1000000

消费管理线程池任务队列大小

flushIntervalConsumeQueue

1000

consumuQueue文件刷盘频率

slaveReadEnable

FALSE

从节点是否可读

transferMsgByHeap

TRUE

消息传输是否使用堆内存

consumerFallbehindThreshold

17179869184

消息消费堆积阈值默认16GB在disableConsumeifConsumeIfConsumerReadSlowly为true时生效

serverAsyncSemaphoreValue

64

异步消息发送最大并发度

maxTransferCountOnMessageInDisk

8

一次消息服务端消息拉取,消息在磁盘中传输允许的最大条数,默认为8条

deleteCommitLogFilesInterval

100

删除commitlog文件的时间间隔

deleteConsumeQueueFilesInterval

100

删除consumequeue文件时间间隔

serverOnewaySemaphoreValue

256

send oneway消息请求并发度

defaultQueryMaxNum

32

查询消息默认返回条数,默认为32

clientSocketRcvBufSize

131072

客户端socket接收缓冲区大小

connectTimeoutMillis

3000

链接超时时间

clientPooledByteBufAllocatorEnable

FALSE

NettyClientConfig,默认为60秒

serverSocketSndBufSize

131072

netty网络socket发送缓存区大小16MB

regionId

DefaultRegion

消息区域

duplicationEnable

FALSE

是否允许重复复制,默认为 false

cleanFileForciblyEnable

TRUE

是否支持强行删除过期文件

serverSelectorThreads

3

IO线程池线程个数,主要是NameServer.broker端解析请求,返回相应的线程个数,这类县城主要是处理网络请求的,解析请求包。然后转发到各个业务线程池完成具体的业务无操作,然后将结果在返回调用方

consumerManageThreadPoolNums

32

服务端处理消费管理 获取消费者列表 更新消费者进度查询消费进度等

haSendHeartbeatInterval

5000

Master与Slave心跳包发送间隔

mapedFileSizeConsumeQueue

300000

单个consumequeue文件大小默认30W*20表示单个Consumequeue文件中存储30W个ConsumeQueue条目

storeCheckpoint

/usr/local/biyao/apache-rocketmq/store/checkpoint

存储 checkpoint 的路径

commitCommitLogLeastPages

4

一次提交至少需要脏页的数量,默认4页,针对 commitlog文件

longPollingEnable

TRUE

是否开启长轮询

flushConsumeQueueLeastPages

2

一次刷盘至少需要脏页的数量,默认2页,针对 Consume文件

defaultTopicQueueNums

6

主体在一个broker上创建队列数量

autoCreateTopicEnable

FALSE

是否自动创建主题

commitIntervalCommitLog

200

commitlog提交频率

maxMsgsNumBatch

64

一次查询消息最大返回消息条数,默认64条

maxIndexNum

20000000

单个索引文件索引条目的个数,默认为两千万

registerBrokerTimeoutMills

6000

注册broker超时时间

serverWorkerThreads

8

netty业务线程池个数

clientSocketSndBufSize

131072

客户端socket发送缓冲区大小

aHousekeepingInterval

20000

Master与slave长连接空闲时间,超过该时间将关闭连接

brokerPermission

6

Broker权限 默认为6表示可读可写

maxTransferCountOnMessageInMemory

32

一次服务消息拉取,消息在内存中传输运行的最大消息条数,默认为32条

RocketMQ Dashboard可视化控制台

RocketMQ Dashboard 是 RocketMQ 的管控利器,为用户提供客户端和应用程序的各种事件、性能的统计信息,支持以可视化工具代替 Topic 配置、Broker 管理等命令行操作。

快速开始

系统要求:

  1. Linux/Unix/Mac

  2. 64bit JDK 1.8+

  3. Maven 3.2.x

  4. 启动 RocketMQ

网络配置:

  1. 云服务器可远程访问或本地虚拟机可 PING 通外网

  2. rocketmq 配置文件 broker.conf / broker-x.properties 设置 nameserver 地址和端口号

  3. 用配置文件启动 broker

1. docker 镜像安装

① 安装docker,拉取 rocketmq-dashboard 镜像

$ docker pull apacherocketmq/rocketmq-dashboard:latest

② docker 容器中运行 rocketmq-dashboard

$ docker run -d --name rocketmq-dashboard -e "JAVA_OPTS=-Drocketmq.namesrv.addr=127.0.0.1:9876" -p 8080:8080 -t apacherocketmq/rocketmq-dashboard:latest

提示

namesrv.addr:port 替换为 rocketmq 中配置的 nameserver 地址:端口号

开放端口号:8080,9876,10911,11011 端口

  • 云服务器:设置安全组访问规则

  • 本地虚拟机:关闭防火墙,或 -add-port

2. 源码安装

源码地址:apache/rocketmq-dashboard

使用手册: https://github.com/apache/rocketmq-dashboard/blob/master/docs/1_0_0/UserGuide_CN.md

下载并解压,切换至源码目录 rocketmq-dashboard-master/

① 编译 rocketmq-dashboard

$ mvn clean package -Dmaven.test.skip=true

② 运行 rocketmq-dashboard

$ java -jar target/rocketmq-dashboard-1.0.1-SNAPSHOT.jar

提示:Started App in x.xxx seconds (JVM running for x.xxx) 启动成功

或使用springBoot:run

mvn spring-boot:run

浏览器页面访问:namesrv.addr:8080

如果你下载这个包很慢,你可以改变maven的镜像(maven的settings.xml)

<mirrors>
    <mirror>
          <id>alimaven</id>
          <name>aliyun maven</name>
          <url>http://maven.aliyun.com/nexus/content/groups/public/</url>
          <mirrorOf>central</mirrorOf>        
    </mirror>
</mirrors>

修改resource/application.properties中的rocketmq.config.namesrvAddr。(或者你可以在操作页面更改它)

账密登录

在rocketmq-dashboard工程的application.yml中有相关配置

rocketmq:
  config:
    # if this value is empty,use env value rocketmq.config.namesrvAddr  NAMESRV_ADDR | now, default localhost:9876
    # configure multiple namesrv addresses to manage multiple different clusters
    namesrvAddrs:
      - 127.0.0.1:9876
    #      - 127.0.0.2:9876
    #      - 10.151.47.32:9876;10.151.47.33:9876;10.151.47.34:9876
    #      - 10.151.47.30:9876
    # if you use rocketmq version < 3.5.8, rocketmq.config.isVIPChannel should be false.default true
    isVIPChannel:
    # timeout for mqadminExt, default 5000ms
    timeoutMillis:
    # rocketmq-console's data path:dashboard/monitor
    dataPath: /tmp/rocketmq-console/data
    # set it false if you don't want use dashboard.default true
    enableDashBoardCollect: true
    # set the message track trace topic if you don't want use the default one
    msgTrackTopicName:
    ticketKey: ticket
    # 开启账密登录,必须在data下有users.properties文件
    # must create userInfo file: ${rocketmq.config.dataPath}/users.properties if the login is required
    loginRequired: false
    useTLS: false
    # set the accessKey and secretKey if you used acl
#    accessKey: rocketmq2
#    secretKey: 12345678

users.properties文件

# 该文件支持热修改,即添加和修改用户时,不需要重新启动console
# 格式, 每行定义一个用户, username=password[,N]  #N是可选项,可以为0 (普通用户); 1 (管理员)  

#定义管理员 
admin=admin,1

#定义普通用户
user1=user1
user2=user2

启动时指定参数:rocketmq.config.loginRequired=true

nohup java -jar -server -Xmx1g -Xms1g -Drocketmq.config.loginRequired=true -Drocketmq.config.namesrvAddrs=192.168.10.174:9876 -Dserver.port=9888 rocketmq-dashboard-1.0.1-SNAPSHOT.jar > /dev/null  2>&1 &

docker方式

docker run -d --restart=always --name rocketmq-broker-a 
         --network rocketmq -p 10909:10909 -p 10911:10911 
         -v /home/rocketmq-4.9.7/broker/broker-a/logs:/home/rocketmq/logs 
         -v /home/rocketmq-4.9.7/broker/broker-a/store:/home/rocketmq/store 
         -v /home/rocketmq-4.9.7/broker/broker-a/conf:/home/rocketmq/rocketmq-4.9.7/conf 
         -e "MAX_POSSIBLE_HEAP=200000000" apache/rocketmq:4.9.7 sh mqbroker -c /home/rocketmq/rocketmq-4.9.7/conf/broker.conf

账户权限校验

如果用户访问console时开启了登录功能,会按照登录的角色对访问的接口进行权限控制。

  • 1.在Spring配置文件resources/application.properties中修改rocketmq.config.loginRequired=true开启登录功能

# 开启登录功能
rocketmq.config.loginRequired=true

# Dashboard文件目录,登录用户配置文件所在目录
rocketmq.config.dataPath=/tmp/rocketmq-console/data   
  • 2.确保${rocketmq.config.dataPath}定义的目录存在,并且该目录下创建访问权限配置文件"role-permission.yml", 如果该目录下不存在此文件,则默认使用resources/role-permission.yml文件。该文件保存了普通用户角色所有能访问的接口地址。 role-permission.yml文件格式为:

# 该文件支持热修改,即添加和修改用户时,不需要重新启动console
# 格式,如果增加和删除接口权限,直接在列表中增加和删除接口地址即可。
# 接口路径配置支持通配符
# * 表示匹配0或多个不是/的字符
# ** 表示匹配0或多个任意字符
# ? 表示匹配1个任意字符

rolePerms:
  # 普通用户
  ordinary:
    - /rocketmq/nsaddr
    - /ops/*
    - /dashboard/**
    - /topic/*.query
    - /topic/sendTopicMessage.do
    - /producer/*.query
    - /message/*
    - /messageTrace/*
    - /monitor/*
    ....

特点

异步解耦、削峰填谷

  • 异步解耦缩短链路

通过消息队列RocketMQ版将上游业务和下游系统进行解耦,缩短了无服务调用的链路。系统异步化后响应时间更短、上下游轻松耦合,开发效率更高。

  • 削峰填谷提高稳定性

传统消息中间件仅满足业务的异步化需求,而消息队列RocketMQ版诞生于在线互联网和交易业务场景,除了满足异步能力,还将削峰填谷作为消息的基础能力。通过削峰填谷不仅能够提高系统稳定性,还能大幅降低业务成本。

实现削峰填谷的功能,消息中间件需要支持海量的消息堆积能力以及处理好冷热消息混合情况下的流量模型。消息队列RocketMQ版能够提供亿级消息堆积能力,可以在大促等流量峰值场景下抗住流量洪峰,保证下游业务能够在安全水位内平滑稳定的运行。

分布式事务

消息队列RocketMQ版分布式事务消息的方案具备如下优势:

  • 系统性能高

基于最终一致性的事务消息方案,相比传统XA事务,吞吐性能更高,可扩展性更强。

  • 开发成本低

基于事务消息开发逻辑简单,仅需两阶段接口即可完成多个事务分支的协调,无需业务做补偿处理。

下图以创建订单为例对比传统事务和消息队列RocketMQ版事务消息的方案:

分布式定时/延时调度

消息队列RocketMQ版提供精确度到秒级的分布式定时消息能力,可广泛应用于订单超时中心处理、分布式延时调度系统等场景。

使用消息队列RocketMQ版定时消息有如下优势:

  • 定时精度高、开发门槛低

消息定时时间不存在阶梯间隔,可以轻松实现任意精度事件触发,无需业务去重。

  • 高性能、可扩展

传统的定时实现方案较为复杂,需要进行数据库扫描,容易遇到性能瓶颈的问题,消息队列RocketMQ版可以基于定时消息特性完成事件驱动,实现百万级消息TPS能力。

基本概念

本文介绍消息队列RocketMQ版的基本概念,以便您更好地理解和使用消息队列RocketMQ版。

主题(Topic)

消息队列RocketMQ版中消息传输和存储的顶层容器,用于标识同一类业务逻辑的消息。主题通过TopicName来做唯一标识和区分。更多信息,请参见主题(Topic)

消息类型(MessageType)

消息队列RocketMQ版中按照消息传输特性的不同而定义的分类,用于类型管理和安全校验。消息队列RocketMQ版支持的消息类型有普通消息、顺序消息、事务消息和定时/延时消息。

消息队列(MessageQueue)

队列是消息队列RocketMQ版中消息存储和传输的实际容器,也是消息的最小存储单元。消息队列RocketMQ版的所有主题都是由多个队列组成,以此实现队列数量的水平拆分和队列内部的流式存储。队列通过QueueId来做唯一标识和区分。更多信息,请参见队列(MessageQueue)

消息(Message)

消息是消息队列RocketMQ版中的最小数据传输单元。生产者将业务数据的负载和拓展属性包装成消息发送到消息队列RocketMQ版服务端,服务端按照相关语义将消息投递到消费端进行消费。更多信息,请参见消息(Message)

消息视图(MessageView)

消息视图是消息队列RocketMQ版面向开发视角提供的一种消息只读接口。通过消息视图可以读取消息内部的多个属性和负载信息,但是不能对消息本身做任何修改。

消息标签(MessageTag)

消息标签是消息队列RocketMQ版提供的细粒度消息分类属性,可以在主题层级之下做消息类型的细分。消费者通过订阅特定的标签来实现细粒度过滤。更多信息,请参见消息过滤

消息位点(MessageQueueOffset)

消息是按到达消息队列RocketMQ版服务端的先后顺序存储在指定主题的多个队列中,每条消息在队列中都有一个唯一的Long类型坐标,这个坐标被定义为消息位点。更多信息,请参见消费进度管理

消费位点(ConsumerOffset)

一条消息被某个消费者消费完成后不会立即从队列中删除,消息队列RocketMQ版会基于每个消费者分组记录消费过的最新一条消息的位点,即消费位点。更多信息,请参见消费进度管理

消息索引(MessageKey)

消息索引是消息队列RocketMQ版提供的面向消息的索引属性。通过设置的消息索引可以快速查找到对应的消息内容。

生产者(Producer)

生产者是消息队列RocketMQ版系统中用来构建并传输消息到服务端的运行实体。生产者通常被集成在业务系统中,将业务消息按照要求封装成消息队列RocketMQ版的消息并发送至服务端。更多信息,请参见生产者(Producer)

事务检查器(TransactionChecker)

消息队列RocketMQ版中生产者用来执行本地事务检查和异常事务恢复的监听器。事务检查器应该通过业务侧数据的状态来检查和判断事务消息的状态。更多信息,请参见事务消息

事务状态(TransactionResolution)

消息队列RocketMQ版中事务消息发送过程中,事务提交的状态标识,服务端通过事务状态控制事务消息是否应该提交和投递。事务状态包括事务提交、事务回滚和事务未决。更多信息,请参见事务消息

消费者分组(ConsumerGroup)

消费者分组是消息队列RocketMQ版系统中承载多个消费行为一致的消费者的负载均衡分组。和消费者不同,消费者分组并不是运行实体,而是一个逻辑资源。在消息队列RocketMQ版中,通过消费者分组内初始化多个消费者实现消费性能的水平扩展以及高可用容灾。更多信息,请参见消费者分组(ConsumerGroup)

消费者(Consumer)

消费者是消息队列RocketMQ版中用来接收并处理消息的运行实体。消费者通常被集成在业务系统中,从消息队列RocketMQ版服务端获取消息,并将消息转化成业务可理解的信息,供业务逻辑处理。更多信息,请参见消费者(Consumer)

消费结果(ConsumeResult)

消息队列RocketMQ版中PushConsumer消费监听器处理消息完成后返回的处理结果,用来标识本次消息是否正确处理。消费结果包含消费成功和消费失败。

订阅关系(Subscription)

订阅关系是消息队列RocketMQ版系统中消费者获取消息、处理消息的规则和状态配置。订阅关系由消费者分组动态注册到服务端系统,并在后续的消息传输中按照订阅关系定义的过滤规则进行消息匹配和消费进度维护。更多信息,请参见订阅关系(Subscription)

消息过滤

消费者可以通过订阅指定消息标签(Tag)对消息进行过滤,确保最终只接收被过滤后的消息合集。过滤规则的计算和匹配在消息队列RocketMQ版的服务端完成。更多信息,请参见消息过滤

重置消费位点

以时间轴为坐标,在消息持久化存储的时间范围内,重新设置消费者分组对已订阅主题的消费进度,设置完成后消费者将接收设定时间点之后,由生产者发送到消息队列RocketMQ版服务端的消息。更多信息,请参见重置消费位点

消息轨迹

在一条消息从生产者发出到消费者接收并处理过程中,由各个相关节点的时间、地点等数据汇聚而成的完整链路信息。通过消息轨迹,您能清晰定位消息从生产者发出,经由消息队列RocketMQ版服务端,投递给消费者的完整链路,方便定位排查问题。

消息堆积

生产者已经将消息发送到消息队列RocketMQ版的服务端,但由于消费者的消费能力有限,未能在短时间内将所有消息正确消费掉,此时在消息队列RocketMQ版的服务端保存着未被消费的消息,该状态即消息堆积。

事务消息

事务消息是消息队列RocketMQ版提供的一种高级消息类型,支持在分布式场景下保障消息生产和本地事务的最终一致性。

定时/延时消息

定时/延时消息是消息队列RocketMQ版提供的一种高级消息类型,消息被发送至服务端后,在指定时间后才能被消费者消费。通过设置一定的定时时间可以实现分布式场景的延时调度触发效果。

顺序消息

顺序消息是消息队列RocketMQ版提供的一种高级消息类型,支持消费者按照发送消息的先后顺序获取消息,从而实现业务场景中的顺序处理。

参数约束和建议

Apache RocketMQ 系统中存在很多自定义参数和资源命名,您在使用 Apache RocketMQ 时建议参考如下说明规范系统设置,避对某些具体参数设置不合理导致应用出现异常。

参数

建议范围

说明

Topic名称

字符建议:字母a~z或A~Z、数字0~9以及下划线()、短划线(-)和百分号(%)。

长度建议:1~64个字符。

系统保留字符:Topic名称不允许使用以下保留字符或含有特殊前缀的字符命名。

保留字符: TBW102 BenchmarkTest SELF_TEST_TOPIC OFFSET_MOVED_EVENT SCHEDULE_TOPIC_XXXX RMQ_SYS_TRANS_HALF_TOPIC RMQ_SYS_TRACE_TOPIC RMQ_SYS_TRANS_OP_HALF_TOPIC

特殊前缀: rmq_sys %RETRY% %DLQ% rocketmq-broker-

Topic命名应该尽量使用简短、常用的字符,避免使用特殊字符。特殊字符会导致系统解析出现异常,字符过长可能会导致消息收发被拒绝。

ConsumerGroup名称

字符建议:支持字母a~z或A~Z、数字0~9以及下划线()、短划线(-)和百分号(%)。

长度建议:1~64个字符。

系统保留字符:ConsumerGroup不允许使用以下保留字符或含有特殊前缀的字符命名。

保留字符: DEFAULT_CONSUMER DEFAULT_PRODUCER TOOLS_CONSUMER FILTERSRV_CONSUMER __MONITOR_CONSUMER CLIENT_INNER_PRODUCER SELF_TEST_P_GROUP SELF_TEST_C_GROUP CID_ONS-HTTP-PROXY CID_ONSAPI_PERMISSION CID_ONSAPI_OWNER CID_ONSAPI_PULL CID_RMQ_SYS_TRANS 特殊字符 CID_RMQ_SYS CID_HOUSEKEEPING

无。

ACL Credentials

字符建议:AK(AccessKey ID)、SK(AccessKey Secret)和Token仅支持字母a~z或A~Z、数字0~9。

长度建议:不超过1024个字符。

无。

请求超时时间

默认值:3000毫秒。

取值范围:该参数为客户端本地行为,取值范围建议不要超过30000毫秒。

请求超时时间是客户端本地同步调用的等待时间,请根据实际应用设置合理的取值,避免线程阻塞时间过长。

消息大小

默认值:不超过4 MB。不涉及消息压缩,仅计算消息体body的大小。

取值范围:建议不超过4 MB。

消息传输应尽量压缩和控制负载大小,避免超大文件传输。若消息大小不满足限制要求,可以尝试分割消息或使用OSS存储,用消息传输URL。

消息自定义属性

字符限制:所有可见字符。

长度建议:属性的Key和Value总长度不超过16 KB。

系统保留属性:不允许使用以下保留属性作为自定义属性的Key。 保留属性Key

无。

MessageGroup

字符限制:所有可见字符。

长度建议:1~64字节。

MessageGroup是顺序消息的分组标识。一般设置为需要保证顺序的一组消息标识,例如订单ID、用户ID等。

消息发送重试次数

默认值:3次。

取值范围:无限制。

消息发送重试是客户端SDK内置的重试策略,对应用不可见,建议取值不要过大,避免阻塞业务线程。 如果消息达到最大重试次数后还未发送成功,建议业务侧做好兜底处理,保证消息可靠性。

消息消费重试次数

默认值:16次。

消费重试次数应根据实际业务需求设置合理的参数值,避免使用重试进行无限触发。重试次数过大容易造成系统压力过量增加。

messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

事务异常检查间隔

默认值:60秒。

事务异常检查间隔指的是,半事务消息因系统重启或异常情况导致没有提交,生产者客户端会按照该间隔时间进行事务状态回查。 间隔时长不建议设置过短,否则频繁的回查调用会影响系统性能。

半事务消息第一次回查时间

默认值:取值等于[事务异常检查间隔] * 最大限制:不超过1小时。

无。

半事务消息最大超时时长

默认值:4小时。 * 取值范围:不支持自定义修改。

半事务消息因系统重启或异常情况导致没有提交,生产者客户端会按照事务异常检查间隔时间进行回查,若超过半事务消息超时时长后没有返回结果,半事务消息将会被强制回滚。 您可以通过监控该指标避免异常事务。

PushConsumer本地缓存

默认值:

最大缓存数量:1024条。

最大缓存大小:64 M。

取值范围:支持用户自定义设置,无限制。

消费者类型为PushConsumer时,为提高消费者吞吐量和性能,客户端会在SDK本地缓存部分消息。缓存的消息的数量和大小应设置在系统内存允许的范围内。

PushConsumer重试间隔时长

默认值:

非顺序性投递:间隔时间阶梯变化,具体取值,请参见PushConsumer消费重试策略。

顺序性投递:3000毫秒。

无。

PushConsumer消费并发度

默认值:20个线程。

无。

获取消息最大批次

默认值:32条。

消费者从服务端获取消息时,一次获取到最大消息条数。建议按照实际业务设置合理的参数值,一次获取消息数量过大容易在消费失败时造成大批量消息重复。

SimpleConsumer最大不可见时间

默认值:用户必填参数,无默认值。

取值范围建议:最小10秒;最大12小时。

消费不可见时间指的是消息处理+失败后重试间隔的总时长,建议设置时取值比实际需要耗费的时间稍微长一些。

消息生产

生产者(Producer):消息队列RocketMQ版中用于产生消息的运行实体,一般集成于业务调用链路的上游。生产者是轻量级匿名无身份的。

消息存储

  • 主题(Topic):消息队列RocketMQ版消息传输和存储的分组容器,主题内部由多个队列组成,消息的存储和水平扩展实际是通过主题内的队列实现的。

  • 队列(MessageQueue):消息队列RocketMQ版消息传输和存储的实际单元容器,类比于Kafka中的分区。消息队列RocketMQ版通过流式特性的无限队列结构来存储消息,消息在队列内具备顺序性存储特征。

  • 消息(Message):消息队列RocketMQ版的最小传输单元。消息具备不可变性,在初始化发送和完成存储后即不可变。

消息消费

  • 消费者分组(ConsumerGroup):消息队列RocketMQ版发布订阅模型中定义的独立的消费身份分组,用于统一管理底层运行的多个消费者(Consumer)。同一个消费组的多个消费者必须保持消费逻辑和配置一致,共同分担该消费组订阅的消息,实现消费能力的水平扩展。

  • 消费者(Consumer):消息队列RocketMQ版消费消息的运行实体,一般集成在业务调用链路的下游。消费者必须被指定到某一个消费组中。

  • 订阅关系(Subscription):消息队列RocketMQ版发布订阅模型中消息过滤、重试、消费进度的规则配置。订阅关系以消费组粒度进行管理,消费组通过定义订阅关系控制指定消费组下的消费者如何实现消息过滤、消费重试及消费进度恢复等。

消息队列RocketMQ版的订阅关系除过滤表达式之外都是持久化的,即服务端重启或请求断开,订阅关系依然保留。

消息消费重试次数

消息队列 RocketMQ 默认允许每条消息最多重试 16 次,每次重试的间隔时间如下:

第几次重试与上次重试的间隔时间第几次重试与上次重试的间隔时间

第几次重试

与上次重试的间隔时间

第几次重试

与上次重试的间隔时间

1

10 秒

9

7 分钟

2

30 秒

10

8 分钟

3

1 分钟

11

9 分钟

4

2 分钟

12

10 分钟

5

3 分钟

13

20 分钟

6

4 分钟

14

30 分钟

7

5 分钟

15

1 小时

8

6 分钟

16

2 小时

如果消息重试 16 次后仍然失败,消息将不再投递。如果严格按照上述重试时间间隔计算,某条消息在一直消费失败的前提下,将会在接下来的 4 小时 46 分钟之内进行 16 次重试,超过这个时间范围消息将不再重试投递。

注意: 一条消息无论重试多少次,这些重试消息的 Message ID 不会改变。

订阅关系一致性

  • 不同消费者分组对于同一个主题的订阅相互独立如下图所示,消费者分组Group A和消费者分组Group B分别以不同的订阅关系订阅了同一个主题Topic A,这两个订阅关系互相独立,可以各自定义,不受影响。

  • 同一个消费者分组对于不同主题的订阅也相互独立如下图所示,消费者分组Group A订阅了两个主题Topic A和Topic B,对于Group A中的消费者来说,订阅的Topic A为一个订阅关系,订阅的Topic B为另外一个订阅关系,且这两个订阅关系互相独立,可以各自定义,不受影响。

1.1 订阅的Topic一样,且过滤表达式一致

如下图所示,同一 ConsumerGroup 下的三个Consumer实例C1、C2和C3分别都订阅了TopicA,且订阅TopicA的Tag也都是Tag1,符合订阅关系一致原则。

正确示例代码一

C1、C2、C3的订阅关系一致,即C1、C2、C3订阅消息的代码必须完全一致,代码示例如下:

PushConsumer consumer1 = provider.newPushConsumerBuilder().setConsumerGroup("GroupA").build();        
consumer1.subscribe("TopicA", new FilterExpression("TagA", FilterExpressionType.TAG));                
PushConsumer consumer2 = provider.newPushConsumerBuilder().setConsumerGroup("GroupA").build();        
consumer2.subscribe("TopicA", new FilterExpression("TagA", FilterExpressionType.TAG));                
PushConsumer consumer3 = provider.newPushConsumerBuilder().setConsumerGroup("GroupA").build();        
consumer3.subscribe("TopicA", new FilterExpression("TagA", FilterExpressionType.TAG));

RocketMQ 强调订阅关系一致,核心是指相同 ConsumerGroup 的每个 Consumer 之间一致,因为在服务端视角看来一个 Group 下的所有 Consumer 都应该是相同的副本逻辑。
强调订阅关系一致,并不是指一个 Consumer 不能订阅多个Topic,每个 Consumer 仍然可以按照需要订阅多个 Topic,但前提是相同消费者分组下的 Consumer 要一致。

2.1 常见订阅关系不一致问题

同一ConsumerGroup下的Consumer实例订阅的Topic不同(3.x、4.x SDK适用)

在早期3.x/4.x 版本的SDK中,如下图所示,同一 ConsumerGroup 下的三个Consumer实例C1、C2和C3分别订阅了TopicA、TopicB和TopicC,订阅的Topic不一致,不符合订阅关系一致性原则。

备注
5.x版本SDK 已经支持同一个 ConsumerGroup 下的Consumer实例订阅不同的Topic。

2.2 同一 ConsumerGroup 下的 Consumer 实例订阅的Topic相同,但订阅的Tag不一致

如下图所示,同一 ConsumerGroup 下的三个Consumer实例C1、C2和C3分别都订阅了TopicA,但是C1订阅TopicA的Tag为Tag1,C2和C3订阅的TopicA的Tag为Tag2,订阅同一Topic的Tag不一致,不符合订阅关系一致性原则。

开始使用

有两种方式,

#application.yml
spring:
  application:
    name: rocketMQ

server:
  port: 8081
  servlet:
    context-path: /rocketMQ

rocketmq:
  name-server: localhost:9876
  producer:
    group: rocketMQGroup1
    send-message-timeout: 3000
  consumer:
    group: rocketMQConsumer1

使用适配springBoot版本

#对应rocketMQ5.0.0
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.2.3</version>
</dependency>

生产者:

/**
 * 消息生产者
 */
public interface RocketMqService {

    /**
     * 同步发送消息
     */
    void send(MqMsg mqMsg);
    /**
     * 异步发送消息,异步返回消息结果
     */
    void asyncSend(MqMsg mqMsg);
    /**
     * 单向发送消息,不关心返回结果,容易消息丢失,适合日志收集、不精确统计等消息发送;
     */
    void syncSendOrderly(MqMsg mqMsg);
}
/**
 * @ClassName RocketMqServiceImpl
 * @Description 消息生产者
 *
 *  进行SpringBoot和RocketMQ整合时,关键使用的是RocketMQTemplate类来进行消息发送,
 *  其中包含有send()、asyncSend()、sendOneWay()、sendMessageInTransaction()等方法,每个方法都至少包含有参数destination和message。
 *  参数说明
 *      destination:消息发送到哪个topic和tag,在SpringBoot中topic和tag使用一个参数发送,其中用英文的冒号(:)进行连接。
 *      message:消息体,需要使用 MessageBuilder.withPayload方法对消息进行封装。
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@Service
@Slf4j
public class ProducerServiceImpl implements RocketMqService {


    private RocketMQTemplate rocketMQTemplate;
    @Autowired
    public void setRocketMQTemplate(RocketMQTemplate rocketMQTemplate) {
        this.rocketMQTemplate = rocketMQTemplate;
    }

    /**
     * 同步发送消息
     */
    @Override
    public void send(MqMsg mqMsg) {
        log.info("send发送消息:{}", mqMsg);
        rocketMQTemplate.send(mqMsg.getTopic()+":"+ mqMsg.getTag(), MessageBuilder.withPayload(mqMsg.getContent()).build());
    }

    /**
     * 异步发送消息,异步返回消息结果
     */
    @Override
    public void asyncSend(MqMsg mqMsg) {
        log.info("asyncSend发送消息:{}", mqMsg);
        rocketMQTemplate.asyncSend(mqMsg.getTopic() + ":" + mqMsg.getTag(), mqMsg.getContent(), new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                log.info("事物消息发送成功:{}", sendResult.getTransactionId());
            }

            @Override
            public void onException(Throwable throwable) {
                log.info("mqMsg={}消息发送失败", mqMsg);
            }
        });
    }

    /**
     * 单向发送消息,不关心返回结果,容易消息丢失,适合日志收集、不精确统计等消息发送;
     */
    @Override
    public void syncSendOrderly(MqMsg mqMsg) {
        log.info("syncSendOrderly发送消息:{}", mqMsg);
        rocketMQTemplate.sendOneWay(mqMsg.getTopic() + ":" + mqMsg.getTag(), mqMsg.getContent());
    }
}
/**
 * @ClassName SpringController
 * @Description   依赖rocketmq-spring-boot-starter,使用@RocketMQAutoConfiguration自动配置方式
 *      自动扩展了很多依赖,根据配置文件,默认注入DefaultMQProducer,DefaultLitePullConsumer
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@RestController
@RequestMapping("/spring")
public class SpringController {

    private RocketMqService producerService;
    @Autowired
    public void setProducerService(RocketMqService producerService) {
        this.producerService = producerService;
    }

    @GetMapping("/send")
    public void send(){
        MqMsg msg = new MqMsg();
        msg.setTopic("topicA");
        msg.setTag("send");
        Map<String, String> map = new HashMap<>();
        map.put("name", "jack");
        map.put("age", "23");
        map.put("tel", "10086");
        msg.setContent(map.toString());
        producerService.send(msg);
    }

    @GetMapping("/asyncSend")
    public void asyncSend(){
        MqMsg msg = new MqMsg();
        msg.setTopic("topicA");
        msg.setTag("asyncSend");
        Map<String, String> userinfo = new HashMap<>();
        userinfo.put("name", "lee");
        userinfo.put("age", "24");
        userinfo.put("tel", "10010");
        msg.setContent(userinfo.toString());
        producerService.asyncSend(msg);
    }
}

消费者:

/**
 * @ClassName RocketMQConsumer
 * @Description  消息消费者
 *  * @RocketMQMessageListener注解参数:
 *  *      topic:表示需要监听哪个topic的消息
 *  *      consumerGroup:表示消费者组
 *  *      selectorExpression:表示需要监听的tag
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@Service
@Slf4j
@RocketMQMessageListener(topic = "topicA",consumerGroup = "consumer1",selectorExpression = "send")
public class SyncMQConsumer implements RocketMQListener<String> {

    /**
     * 接收消息
     * @param s
     */
    @Override
    public void onMessage(String s) {
        log.info("同步消费者成功消费消息:{}", s);
    }


}
/**
 * @ClassName AsyncMQConsumer
 * @Description  异步消费者
 *  * @RocketMQMessageListener注解参数:
 *  *      topic:表示需要监听哪个topic的消息
 *  *      consumerGroup:表示消费者组
 *  *      selectorExpression:表示需要监听的tag
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@Service
@Slf4j
@RocketMQMessageListener(topic = "topicA", consumerGroup = "consumer2",selectorExpression = "asyncSend")
public class AsyncMQConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String s) {
        log.info("异步消费者成功消费消息:{}", s);
    }
}

使用rocketMQ原生客户端

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client-java</artifactId>
    <version>5.0.0</version>
</dependency>

https://github.com/apache/rocketmq/tree/develop/example/src/main/java/org/apache/rocketmq/example/

/**
 * @ClassName ApacheProducerService
 * @Description  原生客户端的生产者
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@Service
public class ApacheProducerService {


    private RocketMQProperties rocketMQPro;
    @Autowired
    public void setRocketMQPro(RocketMQProperties rocketMQPro) {
        this.rocketMQPro = rocketMQPro;
    }

    private DefaultMQProducer producer;


    public DefaultMQProducer getProducer() {
        return producer;
    }

    /**
     * @PostConstruct
     *      初始化producer, 每个生产者在使用前必须start()启动,且只能调用一次.一但关闭不再使用
     *                  producer.start() 进行启动
     *                  producer.shutdown()进行关闭
     * @return
     */
    @PostConstruct
    public void init(){
        producer = new DefaultMQProducer(rocketMQPro.getProducer().getGroup());
        producer.setNamesrvAddr(rocketMQPro.getNameServer());
//        producer.setNamesrvAddr("192,168.0.100:9876;192,168.0.101:9876;192,168.0.102:9876");
        //如果本地多网卡或多通道需要设置false
        producer.setVipChannelEnabled(false);

        try {
            //在使用前,只能启动一次.
            producer.start();
        } catch (MQClientException e) {
            e.printStackTrace();
        }
//        producer.shutdown();   //进行关闭
    }


}
/**
 * @ClassName ApacheController
 * @Description 原生客户端使用方式.
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@RestController
@RequestMapping("/apache")
public class ApacheController {

    private ApacheProducerService service;
    @Autowired
    public void setService(ApacheProducerService service) {
        this.service = service;
    }

    /**
     * 发送消息
     * @return
     */
    @RequestMapping("/send")
    public String send() throws Exception {

        String aaa = "这是一条消息!";
        // 创建一条消息,并指定topic、tag、body等信息,tag可以理解成标签,对消息进行再归类,RocketMQ可以在消费端对tag进行过滤
        //StandardCharsets.UTF_8
        Message msg = new Message("TopicTest","TagA",aaa.getBytes(RemotingHelper.DEFAULT_CHARSET));
        // 利用producer进行发送,并同步等待发送结果
        SendResult sendResult = service.getProducer().send(msg);
        System.out.printf("发送消息ID:%s  发送状态:%n", sendResult.getMsgId(),sendResult.getSendStatus());
        return "ok";
    }
}

消费者:

/**
 * @ClassName ApacheConsumerService
 * @Description 原生客户端的消费者
 * @Author LiuQian
 * @Date 2022-10-29
 **/
@Service
public class ApacheConsumerService {


    private RocketMQProperties rocketMQPro;
    @Autowired
    public void setRocketMQPro(RocketMQProperties rocketMQPro) {
        this.rocketMQPro = rocketMQPro;
    }

    private DefaultMQPushConsumer consumer;


    public DefaultMQPushConsumer getConsumer() {
        return consumer;
    }

    /**
     * @PostConstruct
     *      初始化producer, 每个生产者在使用前必须start()启动,且只能调用一次.一但关闭不再使用
     *                  consumer.start() 进行启动
     *                  consumer.shutdown()进行关闭
     * @return
     */
    @PostConstruct
    public void init(){
        consumer = new DefaultMQPushConsumer(rocketMQPro.getConsumer().getGroup());
        consumer.setNamesrvAddr(rocketMQPro.getNameServer());
//        producer.setNamesrvAddr("192,168.0.100:9876;192,168.0.101:9876;192,168.0.102:9876");

        //消费模式
        //CONSUME_FROM_FIRST_OFFSET从头开始消费一遍
        //CONSUME_FROM_LAST_OFFSET从尾开始消费一遍
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);

        try {
            //订阅一个或多个topic,并指定tag过滤条件,这里指定*表示接收所有tag的消息
            consumer.subscribe("TopicTest","TagA");
            //注册回调接口来处理从Broker中收到的消息
            //MessageListenerConcurrently类是无序队列获取
            //MessageListenerOrderly是有序队列获取
            consumer.registerMessageListener(new MessageListenerConcurrently() {
                @Override
                public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
                    try {
                        System.out.printf("pull消费者收到消息: %s /n/n  ", JSON.toJSONString(consumeConcurrentlyContext));
                        list.forEach( x -> {
                            try {
                                System.out.printf("消费消息ID:%s  内容:%s", x.getMsgId(),new String(x.getBody(), RemotingHelper.DEFAULT_CHARSET));
                            } catch (UnsupportedEncodingException e) {
                                e.printStackTrace();
                            }
                        });
                    } catch (Exception e) {
                        e.printStackTrace();
                        //消费失败,稍后再试
                        return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                    }
                    // 返回消息消费状态,ConsumeConcurrentlyStatus.CONSUME_SUCCESS为消费成功
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                }
            });

            //在使用前,只能启动一次.
            consumer.start();
        } catch (Exception e) {
            e.printStackTrace();
        }
//        consumer.shutdown();   //进行关闭
    }


}

RocketMQ生产部署

  1. 下载并上传二进制执行文件

  2. 环境变量配置ROCKETMQ_HOME和系统参数优化

  3. 修改namesrv的jvm参数和broker的jvm参数

修改namesrv的参数

cd rocketmq-4.8.0
vim bin/runserver.sh

choose_gc_options()
{
    # Example of JAVA_MAJOR_VERSION value : '1', '9', '10', '11', ...
    # '1' means releases befor Java 9
    JAVA_MAJOR_VERSION=$("$JAVA" -version 2>&1 | sed -r -n 's/.* version "([0-9]*).*$/\1/p')
    if [ -z "$JAVA_MAJOR_VERSION" ] || [ "$JAVA_MAJOR_VERSION" -lt "9" ] ; then
      JAVA_OPT="${JAVA_OPT} -server -Xms1g -Xmx1g -Xmn1g -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m"
      JAVA_OPT="${JAVA_OPT} -XX:+UseConcMarkSweepGC -XX:+UseCMSCompactAtFullCollection -XX:CMSInitiatingOccupancyFraction=70 -XX:+CMSParallelRemarkEnabled -XX:SoftRefLRUPolicyMSPerMB=0 -XX:+CMSClassUnloadingEnabled -XX:SurvivorRatio=8 -XX:-UseParNewGC"
      JAVA_OPT="${JAVA_OPT} -verbose:gc -Xloggc:${GC_LOG_DIR}/rmq_srv_gc_%p_%t.log -XX:+PrintGCDetails -XX:+PrintGCDateStamps"
      JAVA_OPT="${JAVA_OPT} -XX:+UseGCLogFileRotation -XX:NumberOfGCLogFiles=5 -XX:GCLogFileSize=30m"
    else
      JAVA_OPT="${JAVA_OPT} -server -Xms4g -Xmx4g -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m"
      JAVA_OPT="${JAVA_OPT} -XX:+UseG1GC -XX:G1HeapRegionSize=16m -XX:G1ReservePercent=25 -XX:InitiatingHeapOccupancyPercent=30 -XX:SoftRefLRUPolicyMSPerMB=0"
      JAVA_OPT="${JAVA_OPT} -Xlog:gc*:file=${GC_LOG_DIR}/rmq_srv_gc_%p_%t.log:time,tags:filecount=5,filesize=30M"
    fi
}

修改broker的参数

vim bin/runbroker.sh

choose_gc_log_directory

JAVA_OPT="${JAVA_OPT} -server -Xms3g -Xmx3g"
  1. 制作启动脚本

#启动namesrv
nohup sh bin/mqnamesrv > /data/logs/rocketmq/nameserver.log 2>&1 &
#启动broker
nohup sh bin/mqbroker -c conf/broker-a.properties -n localhost:9876 > /data/logs/rocketmq/broker.log 2>&1 &
nohup sh mqnamesrv &
nohup sh mqbroker -n localhost:9876 -c /usr/local/rocketmq/rocketmq-5.0.0/conf/2m-2s-sync/broker-a.properties > /usr/local/rocketmq/rocketmq-5.0.0/logs/broker-a.log &
  1. 开启自启动

#创建脚本
vi /etc/systemd/system/namesrv.service

[Unit]
Description=rocketmq nameserver
Documentation=namesrv
After=network.target

[Service]
Type=sample
User=root
ExecStart=/data/rocketmq-4.8.0/bin/mqnamesrv > /data/logs/rocketmq/nameserver.log 
ExecReload=/bin/kill -s HUP $MAINPID
ExecStop=/bin/kill -s QUIT $MAINPID
Restart=0
LimitNOFILE=655350

[Install]
WantedBy=multi-user.target

#加载脚本,并添加到开机项
systemctl daemon-reload
systemctl enable namesrv.service
#配置broker开机自启
vi /etc/systemd/system/broker.service

[Unit]
Description=rocketmq broker
Documentation=broker
After=network.target

[Service]
Type=sample
User=root
ExecStart=/data/rocketmq-4.8.0/bin/mqbroker -c /data/rocketmq-4.8.0/conf/broker.conf > /data/logs/rocketmq/broker.log 
ExecReload=/bin/kill -s HUP $MAINPID
ExecStop=/bin/kill -s QUIT $MAINPID
Restart=0
LimitNOFILE=655350

[Install]
WantedBy=multi-user.target

systemctl daemon-reload
systemctl enable broker.service

编写管理脚本命令,vi /etc/init.d/rocketmq
#!/bin/sh
#
# rocketmq - this script starts and stops the rocketmq daemon
#
# chkconfig: - 85 15

ROCKETMQ_HOME=/root/rocketmq
ROCKETMQ_BIN=${ROCKETMQ_HOME}/bin
ADDR=192.168.0.75:9876
LOG_DIR=${ROCKETMQ_HOME}/logs

start() {
if [ ! -d ${LOG_DIR} ];then
mkdir ${LOG_DIR}
fi
cd ${ROCKETMQ_HOME}
nohup sh bin/mqnamesrv &
echo -n "The Name Server boot success..."
nohup sh bin/mqbroker -c ${ROCKETMQ_HOME}/conf/2m-2s-sync/broker-b-s.properties > /dev/null 2>&1 &
echo -n "The broker[%s, ${ADDR}] boot success..."
}
stop() {
cd ${ROCKETMQ_HOME}
sh bin/mqshutdown broker
sleep 1
sh bin/mqshutdown namesrv
}
restart() {
stop
sleep 5
start
}


case "$1" in
start)
start
;;
stop)
stop
;;
restart)
restart
;;
*)
echo $"Usage: $0 {start|stop|restart}"
exit 2
esac
注意:

(1.1)注: 复制这个脚本到Linux可能会丢失最上面的一部分,不知道为什么,通过xshell打开Linux,然后复制这些内容到文件中,会丢失一部分,自己复制过去以后检查清楚是否丢失一些内容。不然就算注册成一个服务也没法用。

(1.2)ROCKETMQ_HOME的地址、利用nohup启动broker是指定启动的配置文件都要写对

2、将rocketmq服务添加为开机启动服务
chmod +x /etc/init.d/rocketmq

chkconfig --add rocketmq

  1. 开放端口

firewall-cmd --add-port=10911/tcp --permanent
firewall-cmd --add-port=10909/tcp --permanent
firewall-cmd --add-port=9876/tcp --permanent
# 开启后重新加载
firewall-cmd --reload

文章作者: 刘同学
本文链接:
版权声明: 本站所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。转载请注明来自 刘同学的小站
中间件 rocketMQ
喜欢就支持一下吧