1、RocketMQ的基础入门和上手体验
本文主要介绍一下RocketMQ的架构模型和基本概念,并通过基础安装和一个Demo演示的RocketMQ的基本使用
1、RocketMQ架构设计和基本概念
RocketMQ是阿里基于内部孵化并经历过多次双十一考验,捐赠给apache基金会后开源出来的一款分布式、队列模型的消息中间件。RocketMQ的官网地址:https://rocketmq.apache.org/zh/
1.1、RocketMQ架构模型

上图来自RocketMQ官网。可以看出和Kafka、RrabbitMQ等消息中间件类似,也是典型的生产者、broker、消费者之间的发布订阅模型。
1.2、RocketMQ基本概念
主题(Topic)
Apache RocketMQ 中消息传输和存储的顶层容器,用于标识同一类业务逻辑的消息。主题通过TopicName来做唯一标识和区分。
消息类型(MessageType)
Apache RocketMQ 中按照消息传输特性的不同而定义的分类,用于类型管理和安全校验。 Apache RocketMQ 支持的消息类型有普通消息、顺序消息、事务消息和定时/延时消息。
Apache RocketMQ 从5.0版本开始,支持强制校验消息类型,即每个主题Topic只允许发送一种消息类型的消息,这样可以更好的运维和管理生产系统,避免混乱。但同时保证向下兼容4.x版本行为,强制校验功能默认开启。
消息队列(MessageQueue)
队列是 Apache RocketMQ 中消息存储和传输的实际容器,也是消息的最小存储单元。 Apache RocketMQ 的所有主题都是由多个队列组成,以此实现队列数量的水平拆分和队列内部的流式存储。队列通过QueueId来做唯一标识和区分。
消息(Message)
消息是 Apache RocketMQ 中的最小数据传输单元。生产者将业务数据的负载和拓展属性包装成消息发送到服务端,服务端按照相关语义将消息投递到消费端进行消费。
消息视图(MessageView)
消息视图是 Apache RocketMQ 面向开发视角提供的一种消息只读接口。通过消息视图可以读取消息内部的多个属性和负载信息,但是不能对消息本身做任何修改。
消息标签(MessageTag)
消息标签是Apache RocketMQ 提供的细粒度消息分类属性,可以在主题层级之下做消息类型的细分。消费者通过订阅特定的标签来实现细粒度过滤。
消息位点(MessageQueueOffset)
消息是按到达Apache RocketMQ 服务端的先后顺序存储在指定主题的多个队列中,每条消息在队列中都有一个唯一的Long类型坐标,这个坐标被定义为消息位点。
消费位点(ConsumerOffset)
一条消息被某个消费者消费完成后不会立即从队列中删除,Apache RocketMQ 会基于每个消费者分组记录消费过的最新一条消息的位点,即消费位点。这点类似kafka的消费者offset
消息索引(MessageKey)
消息索引是Apache RocketMQ 提供的面向消息的索引属性。通过设置的消息索引可以快速查找到对应的消息内容。
生产者(Producer)
生产者是Apache RocketMQ 系统中用来构建并传输消息到服务端的运行实体。生产者通常被集成在业务系统中,将业务消息按照要求封装成消息并发送至服务端。
生产者支持多种消息发送方式:
- 同步发送
- 异步发送
- 顺序发送
- 单向发送
同步发送和异步发送都需要broker给一个应答ack,参考值:-1(重试166次)、0、1
事务检查器(TransactionChecker)
Apache RocketMQ 中生产者用来执行本地事务检查和异常事务恢复的监听器。事务检查器应该通过业务侧数据的状态来检查和判断事务消息的状态。
事务状态(TransactionResolution)
Apache RocketMQ 中事务消息发送过程中,事务提交的状态标识,服务端通过事务状态控制事务消息是否应该提交和投递。事务状态包括事务提交、事务回滚和事务未决。
消费者分组(ConsumerGroup)
消费者分组是Apache RocketMQ 系统中承载多个消费行为一致的消费者的负载均衡分组。和消费者不同,消费者分组并不是运行实体,而是一个逻辑资源。在 Apache RocketMQ 中,通过消费者分组内初始化多个消费者实现消费性能的水平扩展以及高可用容灾。
消费者(Consumer)
消费者是Apache RocketMQ 中用来接收并处理消息的运行实体。消费者通常被集成在业务系统中,从服务端获取消息,并将消息转化成业务可理解的信息,供业务逻辑处理。
消费者支持两种消费模式:
- 推模式(push)
- 拉模式(pull)
消费结果(ConsumeResult)
Apache RocketMQ 中PushConsumer消费监听器处理消息完成后返回的处理结果,用来标识本次消息是否正确处理。消费结果包含消费成功和消费失败。
订阅关系(Subscription)
订阅关系是Apache RocketMQ 系统中消费者获取消息、处理消息的规则和状态配置。订阅关系由消费者分组动态注册到服务端系统,并在后续的消息传输中按照订阅关系定义的过滤规则进行消息匹配和消费进度维护。
消息过滤
消费者可以通过订阅指定消息标签(Tag)对消息进行过滤,确保最终只接收被过滤后的消息合集。过滤规则的计算和匹配在Apache RocketMQ 的服务端完成。
重置消费位点
以时间轴为坐标,在消息持久化存储的时间范围内,重新设置消费者分组对已订阅主题的消费进度,设置完成后消费者将接收设定时间点之后,由生产者发送到Apache RocketMQ 服务端的消息。
消息轨迹
在一条消息从生产者发出到消费者接收并处理过程中,由各个相关节点的时间、地点等数据汇聚而成的完整链路信息。通过消息轨迹,您能清晰定位消息从生产者发出,经由Apache RocketMQ 服务端,投递给消费者的完整链路,方便定位排查问题。
消息堆积
生产者已经将消息发送到Apache RocketMQ 的服务端,但由于消费者的消费能力有限,未能在短时间内将所有消息正确消费掉,此时在服务端保存着未被消费的消息,该状态即消息堆积。
事务消息
事务消息是Apache RocketMQ 提供的一种高级消息类型,支持在分布式场景下保障消息生产和本地事务的最终一致性。
定时/延时消息
定时/延时消息是Apache RocketMQ 提供的一种高级消息类型,消息被发送至服务端后,在指定时间后才能被消费者消费。通过设置一定的定时时间可以实现分布式场景的延时调度触发效果。
顺序消息
顺序消息是Apache RocketMQ 提供的一种高级消息类型,支持消费者按照发送消息的先后顺序获取消息,从而实现业务场景中的顺序处理。
1.3、和其它常见MQ的对比
| RabbitMQ | Kafka | RockeMQ | |
|---|---|---|---|
| 定位 | 传统消息中间件,保证消息,的可靠性 | 日志消息 | 非日志的可靠性传输 |
| 可用性 | 支持cluster普通模式、镜像队列模式 | 异步刷盘,可能出现数据丢失 | 实现了异步/同步刷盘 |
| 单机吞吐量 | 1w | 10w | 10w |
| 积压消息能力 | 根据内存和磁盘阈值来决定(RabbitMQ告警阈值配置) | 非常好,受磁盘限制 | 非常好,受磁盘限制 |
| 顺序消费 | 支持 | 支持 | 支持 |
| 定时消息 | 支持(依赖插件) | 不支持 | 支持 |
| 事务消息 | 不支持(confirm确认模式) | 不支持(通过事务日志和事务协调者) | 支持(half消息,类似于2pc) |
| 消息重试 | 支持 | 不支持(消费者手动提交offset) | 支持 |
| 死信队列 | 支持 | 不支持 | 支持 |
2、RocketMQ快速启动
2.1、下载 & 解压
进入RocketMQ的官方文档,参考页面:https://rocketmq.apache.org/zh/docs/quickStart/01quickstart
这里提供了RocketMQ的源码包地址和二进制包地址,直接下载二进制包即可。(直接下载可能不是最新版本,可以通过这个地址下载对应的版本:https://dist.apache.org/repos/dist/release/rocketmq/)
解压之后的目录结构如图所示:

所有的命令工具都在bin目录下:

2.2、启动nameserver
rocketmq的运行依赖与nameserver,类似与早期的kafka的启动需要依赖zookeeper一样。如果我们需要单独配置执行的namesever,可以在后续启动命令中直接指定nameserver的地址即可,如果我们什么都没有准备,RocketMQ的二进制包也提供了nameserver,可以直接启动使用
nohup /home/ubuntu/tools/rocketmq/rocketmq-all-5.5.1-bin-release/bin/mqnamesrv &
Nameserver和RocketMQ都是java语言编写的,运行的时候难免会涉及一些参数的配置,比如JVM参数。发布者在发布版本的时候主要是基于线上性能考量,如果是我们本地运行可能并不会用到那么大的资源,也或者是本地并没有那么多资源供我们使用,所以可以酌情进行调整。
查看mqnamesrv文件内容,可以看到最终实际是执行runserver.sh这个脚本:
#!/bin/sh
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
if [ -z "$ROCKETMQ_HOME" ] ; then
## resolve links - $0 may be a link to maven's home
PRG="$0"
# need this for relative symlinks
while [ -h "$PRG" ] ; do
ls=`ls -ld "$PRG"`
link=`expr "$ls" : '.*-> \(.*\)$'`
if expr "$link" : '/.*' > /dev/null; then
PRG="$link"
else
PRG="`dirname "$PRG"`/$link"
fi
done
saveddir=`pwd`
ROCKETMQ_HOME=`dirname "$PRG"`/..
# make it fully qualified
ROCKETMQ_HOME=`cd "$ROCKETMQ_HOME" && pwd`
cd "$saveddir"
fi
export ROCKETMQ_HOME
sh ${ROCKETMQ_HOME}/bin/runserver.sh -Drmq.logback.configurationFile=$ROCKETMQ_HOME/conf/rmq.namesrv.logback.xml org.apache.rocketmq.namesrv.NamesrvStartup $@
继续查看runserver.sh文件:
choose_gc_options()
{
# Example of JAVA_MAJOR_VERSION value : '1', '9', '10', '11', ...
# '1' means releases before Java 9
JAVA_MAJOR_VERSION=$("$JAVA" -version 2>&1 | awk -F '"' '/version/ {print $2}' | awk -F '.' '{print $1}')
if [ -z "$JAVA_MAJOR_VERSION" ] || [ "$JAVA_MAJOR_VERSION" -lt "9" ] ; then
JAVA_OPT="${JAVA_OPT} -server -Xms4g -Xmx4g -Xmn2g -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
}
从choose_gc_options方法可以看到jvm初始内存和最大内存都设置的是4g,可以合理调整。调整之后再执行上面的启动命令
ubuntu@study-server-standalone:~/tools/rocketmq/rocketmq-all-5.5.1-bin-release$ nohup /home/ubuntu/tools/rocketmq/rocketmq-all-5.5.1-bin-release/bin/mqnamesrv &
[1] 55149
ubuntu@study-server-standalone:~/tools/rocketmq/rocketmq-all-5.5.1-bin-release$ nohup: ignoring input and appending output to 'nohup.out'
查看日志输出内容可以看到nameserver已经成功启动了:

#查看进程:
ubuntu@study-server-standalone:~/tools/rocketmq/rocketmq-all-5.5.1-bin-release$ jps
55797 Jps
55178 NamesrvStartup
# 确认nameserver进程
ubuntu@study-server-standalone:~/tools/rocketmq/rocketmq-all-5.5.1-bin-release$
ubuntu@study-server-standalone:~/tools/rocketmq/rocketmq-all-5.5.1-bin-release$ pwdx 55178
55178: /home/ubuntu/tools/rocketmq/rocketmq-all-5.5.1-bin-release
# 查看进程监听的端口号,可以看到nameserver监听在了9876端口
ubuntu@study-server-standalone:~/tools/rocketmq/rocketmq-all-5.5.1-bin-release$ netstat -nltp | grep 55178
(Not all processes could be identified, non-owned process info
will not be shown, you would have to be root to see it all.)
tcp6 0 0 :::9876 :::* LISTEN 55178/java
2.3、启动broker
在确认nameserver启动成功之后可以通过命令直接启动rocketmq broker
nohup /home/ubuntu/tools/rocketmq/rocketmq-all-5.5.1-bin-release/bin/mqbroker -n localhost:9876 &
这里通过-n参数指定了nameserver的地址。
但是同样的,这里也会涉及到jvm内存配置问题,查看mqbroker文件发现最终执行的是./bin/runbroker.sh文件,查看文件内容:
#!/bin/sh
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#===========================================================================================
# Java Environment Setting
#===========================================================================================
error_exit ()
{
echo "ERROR: $1 !!"
exit 1
}
find_java_home()
{
if [ -n "$JAVA_HOME" ]; then
return
fi
case "`uname`" in
Darwin)
JAVA_HOME=$(/usr/libexec/java_home)
;;
*)
JAVA_HOME=$(dirname $(dirname $(readlink -f $(which javac))))
;;
esac
}
find_java_home
[ ! -e "$JAVA_HOME/bin/java" ] && JAVA_HOME=$HOME/jdk/java
[ ! -e "$JAVA_HOME/bin/java" ] && JAVA_HOME=/usr/java
[ ! -e "$JAVA_HOME/bin/java" ] && error_exit "Please set the JAVA_HOME variable in your environment, We need java(x64)!"
export JAVA_HOME
export JAVA="$JAVA_HOME/bin/java"
export BASE_DIR=$(dirname $0)/..
export CLASSPATH=.:${BASE_DIR}/conf:${BASE_DIR}/lib/*:${CLASSPATH}
#===========================================================================================
# JVM Configuration
#===========================================================================================
# The RAMDisk initializing size in MB on Darwin OS for gc-log
DIR_SIZE_IN_MB=600
choose_gc_log_directory()
{
case "`uname`" in
Darwin)
if [ ! -d "/Volumes/RAMDisk" ]; then
# create ram disk on Darwin systems as gc-log directory
DEV=`hdiutil attach -nomount ram://$((2 * 1024 * DIR_SIZE_IN_MB))` > /dev/null
diskutil eraseVolume HFS+ RAMDisk ${DEV} > /dev/null
echo "Create RAMDisk /Volumes/RAMDisk for gc logging on Darwin OS."
fi
GC_LOG_DIR="/Volumes/RAMDisk"
;;
*)
# check if /dev/shm exists on other systems
if [ -d "/dev/shm" ]; then
GC_LOG_DIR="/dev/shm"
else
GC_LOG_DIR=${BASE_DIR}
fi
;;
esac
}
choose_gc_options()
{
JAVA_MAJOR_VERSION=$("$JAVA" -version 2>&1 | head -1 | cut -d'"' -f2 | sed 's/^1\.//' | cut -d'.' -f1)
if [ -z "$JAVA_MAJOR_VERSION" ] || [ "$JAVA_MAJOR_VERSION" -lt "8" ] ; then
JAVA_OPT="${JAVA_OPT} -Xmn4g -XX:+UseConcMarkSweepGC -XX:+UseCMSCompactAtFullCollection -XX:CMSInitiatingOccupancyFraction=70 -XX:+CMSParallelRemarkEnabled -XX:SoftRefLRUPolicyMSPerMB=0 -XX:+CMSClassUnloadingEnabled -XX:SurvivorRatio=8 -XX:-UseParNewGC"
else
JAVA_OPT="${JAVA_OPT} -XX:+UseG1GC -XX:G1HeapRegionSize=16m -XX:G1ReservePercent=25 -XX:InitiatingHeapOccupancyPercent=30 -XX:SoftRefLRUPolicyMSPerMB=0"
fi
if [ -z "$JAVA_MAJOR_VERSION" ] || [ "$JAVA_MAJOR_VERSION" -lt "9" ] ; then
JAVA_OPT="${JAVA_OPT} -verbose:gc -Xloggc:${GC_LOG_DIR}/rmq_srv_gc_%p_%t.log -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:+PrintGCApplicationStoppedTime -XX:+PrintAdaptiveSizePolicy"
JAVA_OPT="${JAVA_OPT} -XX:+UseGCLogFileRotation -XX:NumberOfGCLogFiles=5 -XX:GCLogFileSize=30m"
else
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
}
choose_gc_log_directory
JAVA_OPT="${JAVA_OPT} -server -Xms8g -Xmx8g"
choose_gc_options
JAVA_OPT="${JAVA_OPT} -XX:-OmitStackTraceInFastThrow"
JAVA_OPT="${JAVA_OPT} -XX:+AlwaysPreTouch"
JAVA_OPT="${JAVA_OPT} -XX:MaxDirectMemorySize=15g"
JAVA_OPT="${JAVA_OPT} -XX:-UseLargePages -XX:-UseBiasedLocking -XX:+IgnoreUnrecognizedVMOptions"
#JAVA_OPT="${JAVA_OPT} -Xdebug -Xrunjdwp:transport=dt_socket,address=9555,server=y,suspend=n"
JAVA_OPT="${JAVA_OPT} ${JAVA_OPT_EXT}"
JAVA_OPT="${JAVA_OPT} -cp ${CLASSPATH}"
numactl --interleave=all pwd > /dev/null 2>&1
if [ $? -eq 0 ]
then
if [ -z "$RMQ_NUMA_NODE" ] ; then
numactl --interleave=all $JAVA ${JAVA_OPT} $@
else
numactl --cpunodebind=$RMQ_NUMA_NODE --membind=$RMQ_NUMA_NODE $JAVA ${JAVA_OPT} $@
fi
else
"$JAVA" ${JAVA_OPT} $@
fi
可以看到rocketmq的broker更“离谱”,默认直接分配了8g的内存,适当调整后再次执行上面启动broker的命令,然后查看nohup文件输出:

或者查看对应的日志文件:
- nameserver的日志输出
~/logs/rocketmqlogs/namesrv.log

从日志可以看到有一个新的broker注册到了nameserver
- broker的日志输出
~/logs/rocketmqlogs/broker.log

要确认broker是否正确启动,同样可以用命令行检查一下
# 查看进程id
jps
# 确认进程
pwdx [pid]
# 查看监听的端口
netstat -nltp | grep [pid]
2.4、关闭服务
要停止namesever服务,可以执行命令
/home/ubuntu/tools/rocketmq/rocketmq-all-5.5.1-bin-release/bin/mqshutdown namesrv
要停止broker服务,可以执行命令
/home/ubuntu/tools/rocketmq/rocketmq-all-5.5.1-bin-release/bin/mqshutdown broker
3、RocketMQ上手体验
3.1、通过自带的工具发送、接收消息
rocketmq自带了一个Producer和Consumer,但是要使用这两个工具的话会用到一个主题topic:“TopicTest”,所以需要先创建
3.1.1、创建topic
mqadmin updateTopic -n localhost:9876 -b localhost:10911 -t TopicTest
执行结果:

注意:
1、这里直接执行了
mqadmin命令,是因为我在/etc/profile里将rocketmq的bin目录配置到了环境变量中2、
updateTopic可以用于创建或者修改topic3、
-n参数指定对应的nameserver4、
-b参数指定对应的broker5、
-t参数指定要创建或者修改的主题名称
3.1.2、Producer发送消息
如果我们什么都不做,直接执行Producer的代码,比如:
tools.sh org.apache.rocketmq.example.quickstart.Producer
这个时候会看到程序会报错:

这是因为程序在执行的时候会连接到的nameserver,而这个额地址是从环境变量中获取的,所以在执行之前需要设置一个环境变量
export NAMESRV_ADDR=localhost:9876
配置完成之后再执行上面Producer代码就可以发现消息发送成功了

3.1.3、Consumer消费消息
消费者消费消息就很简单了,直接执行Consumer的代码就行
tools.sh org.apache.rocketmq.example.quickstart.Consumer
从控制台可以看到Consumer Started,并且后续有对应消息的消费记录:

Demo的消费者和生产者的代码不同,生产者在消息生产完之后就会自动退出,而消费者在消费完消息后还会继续监听这个topic,并不会自动退出
3.2、Java客户端发送接收消息
3.2.1、引入maven依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>5.5.1</version>
</dependency>
因为前面我们用的是5.5.1版本,所以这里客户端我们也同样直接用5.5.1版本
3.2.2、消费者
Consumer.java
package com.study.rocketmq;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.io.UnsupportedEncodingException;
import java.util.List;
public class Consumer {
public static final String NAME_SERVER_ADDR = "192.168.2.135:9876";
public static void main(String[] args) throws MQClientException {
// 1. 创建消费者对象,这里采用Push模式的消费者
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GROUP_TEST");
// 2. 设置NameServer的地址,如果设置了环境变量NAMESRV_ADDR,可以省略此步
consumer.setNamesrvAddr(NAME_SERVER_ADDR);
// 消费重试次数 -1代表16次
consumer.setMaxReconsumeTimes(-1);
// 3. 订阅对应的主题和Tag(前面测试主题、匹配所有tag)
consumer.subscribe("TopicTest", "*");
// 4. 注册消息接收到Broker消息后的处理接口
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
try {
MessageExt messageExt = list.get(0);
System.out.printf("线程:%-25s 接收到新消息 %s --- %s %n", Thread.currentThread().getName(), messageExt.getTags(), new String(messageExt.getBody(), RemotingHelper.DEFAULT_CHARSET));
} catch (UnsupportedEncodingException e) {
e.printStackTrace();
}
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
});
// 5. 启动消费者(必须在注册完消息监听器后启动,否则会报错)
consumer.start();
System.out.println("已启动消费者");
}
}
这里消费者用的Push模式的消费者。rocketmq提供了Push和Pull模式的接口和实现类,但是DefaultMQPullConsumer被标记了@Deprecated,所以不建议使用

3.2.3、生产者
Producer.java
package com.study.rocketmq;
import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.apache.rocketmq.remoting.exception.RemotingException;
import java.io.UnsupportedEncodingException;
public class Producer {
public static final String NAME_SERVER_ADDR = "192.168.2.135:9876";
public static void main(String[] args) throws MQClientException, UnsupportedEncodingException, RemotingException, InterruptedException, MQBrokerException {
// 1. 创建生产者对象
DefaultMQProducer producer = new DefaultMQProducer("GROUP_TEST");
// 2. 设置NameServer的地址,如果设置了环境变量NAMESRV_ADDR,可以省略此步
producer.setNamesrvAddr(NAME_SERVER_ADDR);
// 3. 启动生产者
producer.start();
// 4. 生产者发送消息
for (int i = 0; i < 10; i++) {
Message message = new Message("TopicTest", "TagA", ("Hello MQ:" + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
SendResult result = producer.send(message);
System.out.printf("发送结果:%s%n", result);
}
// 5. 停止生产者
producer.shutdown();
}
}
3.2.4、先后执行消费者和生产者
消费者输出:

生产者输出:

4、坑
这里遗留了额一个坑,那就是如果broker有多个IP的时候,有可能会出现broker在nameserever上注册的IP地址并不是我们预期的地址,比如内网地址是:192.168.2.135,可能另外还有一个虚拟网卡的地址是172.16.0.96,在程序执行的时候提示172.16.0.96:10911timeout
这里边的原因是注册的时候会读取BrokerConfig类的配置信息,如果我们注册的时候直接按照前面的方式注册的话,程序从BrokerConfig获取IP的方式是:
brokerIP1 = RemotingUtil.getLocalAddress();
这里边的思路是,遍历本地的所有网卡ip,过滤掉 “127.0” 和“192.168”开头的ip地址然后得到第一个ip,为本机ip,当这台机器有很多别的网卡(如:安装docker后),broker使用的ip,就可能会导致我们的客户端无法连接
这只是默认的情况,如果我们在将broker往nameserver注册的时候指定了-c参数的话,程序会去加载额外配置文件。加载时,通过反射的方式,根据配置文件中的键值对,赋值到BrokerConfig 中对应的属性中。
而这个配置文件在./config目录下:

默认的文件内容:
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
brokerClusterName = DefaultCluster
brokerName = broker-a
brokerId = 0
deleteWhen = 04
fileReservedTime = 48
brokerRole = ASYNC_MASTER
flushDiskType = ASYNC_FLUSH
添加一行配置即可:
brokerIP1=[broker的IP地址]

