微服务架构中,消息队列和远程服务调用已经是两大必不可少的组件,而RocketMQ和Dubbo正是阿里系贡献的对应的两大精品开源,做为两个已经获得普遍应用的框架,好好学习研究是必需的。html
官方文档:https://github.com/alibaba/RocketMQ/wiki/quick-start
根据文档说明,须要如下软件来完成这个快速开始示例:
⑴ 64bit OS, best to have Linux/Unix/Mac;
⑵ 64bit JDK 1.6+;
⑶ Maven 3.x
⑷ Git
⑸ Screenjava
做为纯Java程序,RocketMQ在Windows下也是能够运行的,官方还准备了exe执行文件方便Windows环境下进行开发部署。
Windows下的编译部署大同小异,有兴趣能够参考下面这个网址:
http://blog.csdn.net/ruishenh/article/details/22390809linux
若是Linux已经自带JDK,能够使用命令查看JDK版本,若是版本不符合64位1.6+,须要先卸载旧版本而后安装新版本。
我安装的是jdk-8u111-linux-x64.rpm,过程略过,有问题能够参考下面的网址。
下载地址:http://www.oracle.com/technetwork/java/javase/downloads/index.htm
安装教程:http://www.cnblogs.com/benio/archive/2010/09/14/1825909.html
按理32位JDK也是能够运行的,只是须要调整内存配置,但可能不适合生产环境。未经测试,不妄言。git
1.3.1 下载github
cd /usr/javawork wget http://mirrors.hust.edu.cn/apache/maven/maven-3/3.3.9/binaries/apache-maven-3.3.9-bin.tar.gz
1.3.2 解压算法
tar -zxvf apache-maven-3.3.9-bin.tar.gz
1.3.3 设置环境变量shell
vi /etc/profile
文件末尾添加两行配置:apache
export M2_HOME=/usr/javawork/apache-maven-3.3.9 export PATH=$PATH:$M2_HOME/bin
退出vi执行命令使其生效:windows
source /etc/profile
1.3.4 添加alibaba的Maven仓库镜像(下载速度飞快)bash
vi /usr/javawork/apache-maven-3.3.9/conf/settings.xml
在<mirrors>项下添加镜像信息:
<mirror> <id>alimaven</id> <name>aliyun maven</name> <mirrorOf>central</mirrorOf> <url>http://maven.aliyun.com/nexus/content/groups/public/</url> </mirror>
yum install git
若是已下载RocketMQ源码包,Git能够无需安装。shell安装脚本中有git pull命令,若是未安装git,会提示command not found,但不影响后面的编译。
若是嫌烦,vi install.sh打开文件删掉 git pull这条命令便可。
yum install screen
screen 非必需,但安装后切换会话很是方便。官方文档中使用了这条命令,因此仍是装上较好。
命令介绍:http://www.cnblogs.com/mchina/archive/2013/01/30/2880680.html
git clone https://github.com/alibaba/RocketMQ.git cd RocketMQ bash install.sh
若是编译成功,最终获得的目录以下图。devenv是软连接,源文件在target目录。
设置环境变量
cd devenv echo "ROCKETMQ_HOME=`pwd`" >> ~/.bash_profile
使环境变量生效:
source ~/.bash_profile
screen bash mqnamesrv
若是未安装screen,能够使用下面这条命令:
nohup sh mqnamesrv &(加 & 能够后台运行,不然Ctrl+c命令退出当前会话,服务会中止)
启动成功的信息:
Java HotSpot(TM) 64-Bit Server VM warning: ignoring option PermSize=128m; support was removed in 8.0
Java HotSpot(TM) 64-Bit Server VM warning: ignoring option MaxPermSize=320m; support was removed in 8.0
Java HotSpot(TM) 64-Bit Server VM warning: UseCMSCompactAtFullCollection is deprecated and will likely be removed in a future release.
The Name Server boot success. serializeType=JSON
备注:出现三条警告信息是由于JDK1.8已经不支持设置方法区大小和指定CMS垃圾收集算法进行FullGC,后面会讲如何去掉这些信息。
先按ctrl+a,而后再按d挂起当前会话,而后再执行如下命令:
screen bash mqbroker -n localhost:9876
启动成功后除了警告信息外,会出现如下信息:
3.3.1 修改xml文件
cd bin vi mqadmin.xml vi mqbroker.xml vi mqnamesrv.xml vi mqfiltersrv.xml
依次打开这些xml文件,并删除下图红色框中的标记的配置信息。
3.3.2 修改启动文件
vi runserver.sh vi runbroker.sh
依次打开这两个文件,去除下图红色标记的配置信息。
如图, 黄色部分是内存设置, 由于当前是虚拟机搭建的开发环境, 因此内存调整成以下:
runserver.sh: -Xms512m -Xmx512m -Xmn256m
runbroker.sh: -Xms2g -Xmx2g -Xmn512m
关闭Name Server
sh mqshutdown namesrv
关闭Broker
sh mqshutdown broker
cd ~/logs/rocketmqlogs/
3.6.1 设置地址
export NAMESRV_ADDR=localhost:9876
3.6.2 测试命令
生产者:
bash tools.sh com.alibaba.rocketmq.example.quickstart.Producer
消费者:
bash tools.sh com.alibaba.rocketmq.example.quickstart.Consumer
/**生产者*/ public class Producer { public static void main(String[] args) throws MQClientException, InterruptedException { /** * 一个应用建立一个Producer,由应用来维护此对象,能够设置为全局对象或者单例<br> * 注意:ProducerGroupName须要由应用来保证惟一<br> * ProducerGroup这个概念发送普通的消息时,做用不大,可是发送分布式事务消息时,比较关键, * 由于服务器会回查这个Group下的任意一个Producer */ final DefaultMQProducer producer = new DefaultMQProducer("ProducerGroupName"); producer.setNamesrvAddr("192.168.0.11:9876"); producer.setInstanceName("Producer"); /** * Producer对象在使用以前必需要调用start初始化,初始化一次便可<br> * 注意:切记不能够在每次发送消息时,都调用start方法 */ producer.start(); /** * 下面这段代码代表一个Producer对象能够发送多个topic,多个tag的消息。 * 注意:send方法是同步调用,只要不抛异常就标识成功。可是发送成功也可会有多种状态,<br> * 例如消息写入Master成功,可是Slave不成功,这种状况消息属于成功,可是对于个别应用若是对消息可靠性要求极高,<br> * 须要对这种状况作处理。另外,消息可能会存在发送失败的状况,失败重试由应用来处理。 */ for (int i = 0; i < 10; i++) { try { { Message msg = new Message("TopicTest1", // topic "TagA", // tag "OrderID001", // key ("Hello MetaQA").getBytes());// body SendResult sendResult = producer.send(msg); System.out.println(sendResult); } { Message msg = new Message("TopicTest2", // topic "TagB", // tag "OrderID0034", // key ("Hello MetaQB").getBytes());// body SendResult sendResult = producer.send(msg); System.out.println(sendResult); } { Message msg = new Message("TopicTest3", // topic "TagC", // tag "OrderID061", // key ("Hello MetaQC").getBytes());// body SendResult sendResult = producer.send(msg); System.out.println(sendResult); } } catch (Exception e) { e.printStackTrace(); } TimeUnit.MILLISECONDS.sleep(1000); } /** * 应用退出时,要调用shutdown来清理资源,关闭网络链接,从MetaQ服务器上注销本身 * 注意:咱们建议应用在JBOSS、Tomcat等容器的退出钩子里调用shutdown方法 */ // producer.shutdown(); Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() { public void run() { producer.shutdown(); } })); System.exit(0); } }
/** 消费者 */ public class PushConsumer { /** * 当前例子是PushConsumer用法,使用方式给用户感受是消息从RocketMQ服务器推到了应用客户端。<br> * 可是实际PushConsumer内部是使用长轮询Pull方式从MetaQ服务器拉消息,而后再回调用户Listener方法<br> */ public static void main(String[] args) throws InterruptedException, MQClientException { /** * 一个应用建立一个Consumer,由应用来维护此对象,能够设置为全局对象或者单例<br> * 注意:ConsumerGroupName须要由应用来保证惟一 */ DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroupName"); consumer.setNamesrvAddr("192.168.0.11:9876"); consumer.setInstanceName("Consumber"); /** * 订阅指定topic下tags分别等于TagA或TagC或TagD */ consumer.subscribe("TopicTest1", "TagA || TagC || TagD"); /** * 订阅指定topic下全部消息<br> * 注意:一个consumer对象能够订阅多个topic */ consumer.subscribe("TopicTest2", "*"); consumer.registerMessageListener( new MessageListenerConcurrently() { public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) { System.out.println(Thread.currentThread().getName() + " Receive New Messages: " + msgs.size()); MessageExt msg = msgs.get(0); if (msg.getTopic().equals("TopicTest1")) { // 执行TopicTest1的消费逻辑 if (msg.getTags() != null && msg.getTags().equals("TagA")) { // 执行TagA的消费 System.out.println(new String(msg.getBody())); } else if (msg.getTags() != null && msg.getTags().equals("TagC")) { // 执行TagC的消费 System.out.println(new String(msg.getBody())); } else if (msg.getTags() != null && msg.getTags().equals("TagD")) { // 执行TagD的消费 System.out.println(new String(msg.getBody())); } } else if (msg.getTopic().equals("TopicTest2")) { System.out.println(new String(msg.getBody())); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); /** * Consumer对象在使用以前必需要调用start初始化,初始化一次便可<br> */ consumer.start(); System.out.println("ConsumerStarted."); } }
分别运行以上两个程序,正常会输出以下信息:
生产者端:
SendResult [sendStatus=SEND_OK, msgId=C0A80008326873D16E9381E943250000,offsetMsgId=C0A8000B00002A9F0000000000031862, messageQueue=MessageQueue [topic=TopicTest1, brokerName=localhost.localdomain, queueId=2], queueOffset=9] SendResult [sendStatus=SEND_OK, msgId=C0A80008326873D16E9381E943540001,offsetMsgId=C0A8000B00002A9F0000000000031921, messageQueue=MessageQueue [topic=TopicTest2, brokerName=localhost.localdomain, queueId=0], queueOffset=10] SendResult [sendStatus=SEND_OK, msgId=C0A80008326873D16E9381E943630002,offsetMsgId=C0A8000B00002A9F00000000000319E1, messageQueue=MessageQueue [topic=TopicTest3, brokerName=localhost.localdomain, queueId=2], queueOffset=11]
消费者端:
ConsumerStarted. ConsumeMessageThread_1 Receive New Messages: 1 Hello MetaQA ConsumeMessageThread_4 Receive New Messages: 1 Hello MetaQA
com.alibaba.rocketmq.client.exception.MQClientException: No route info of this topic, TopicTest1
See http://docs.aliyun.com/cn#/pub/ons/faq/exceptions&topic_not_exist for further details.
若是出现以上错误,是由于启动Broker时后面的主机和端口未指定
screen bash mqbroker -n localhost:9876
JDK:JDK_1.8.111_x64
O S:Centos_6.5_x64
RocketMQ:3.5.8