RocketMQ 可视化环境搭建和基础代码使用

RocketMQ 是一款分布式消息中间件,最初是由阿里巴巴消息中间件团队研发并大规模应用于生产系统,满足线上海量消息堆积的需求, 在 2016 年底捐赠给 Apache 开源基金会成为孵化项目,经过不到一年时间正式成为了 Apache 顶级项目。<br />早期阿里曾经基于 ActiveMQ 研发消息系统, 随着业务消息的规模增大,瓶颈逐渐显现,后来也考虑过Kafka,但因为在低延迟和高可靠性方面没有选择,最后才自主研发了 RocketMQ, 各方面的性能都比目前已有的消息队列要好,RocketMQ 和 Kafka 在概念和原理上都非常相似,所以也经常被拿来对比;RocketMQ 默认采用长轮询的拉模式, 单机支持千万级别的消息堆积,可以非常好的应用在海量消息系统中。<br />本文分为三部分,如下图所示:<br />image.png <a name="1xvvk"></a>

1 安装 RocketMQ—Windows 版本

<a name="RdEk4"></a>

(1)下载 Windows 安装包

Windows 版本下载地址:http://rocketmq.apache.org/release_notes/<br />image.png<br />下载并解压 rocketmq 安装包。 <a name="LHt4n"></a>

(2)配置系统环境变量

配置系统变量 ROCKETMQ_HOME=“D:\soft\rocketmq-all-4.5.1-bin-release”,如下图所示:<br />注意:每个人 rocketmq 存放目录不一样,我的在 D:\soft 下,用户根据自己的环境配置相应的系统变量。<br />image.png

因为接下来启动 mqnamesrv.cmd 中使用到了环境变量 %ROCKETMQ_HOME%,所以这里需要配置此系统变量。

<a name="2S7UG"></a>

(3)启动 namesrv

进入 rocketmq 的 bin 目录,执行 start mqnamesrv.cmd ,执行成功如下图所示:<br />image.png<br />注意:启动之后,不能关闭此窗口。 <a name="28zct"></a>

(4)启动 broker

还是在 bin 目录下执行 start mqbroker.cmd -n 127.0.0.1:9876 autoCreateTopicEnable=true ,执行成功如下图所示:<br />image.png<br />同样不要关闭以上运行窗口。<br />完成以下步骤,说明你的 RocketMQ 已经按照成功了。 <a name="0ItMT"></a>

2 安装可视化插件

<a name="5iVLy"></a>

(1)下载插件

打开连接 https://github.com/apache/rocketmq-externals.git 下载可视化插件 rocketmq-externals,如下图所示:<br />image.png<br />点击 Download ZIP 进行下载。

我为大家准备了国内百度云的下载链接,方便大家使用。 百度链接:https://pan.baidu.com/s/1sMO6W-562IFJF1uUBQFXYg   提取码:fuzy 

<a name="mNB2M"></a>

(2)配置插件

下载完成之后,进入 rocketmq-externals\rocketmq-console\src\main\resources\application.properties 进行配置,如下图所示:<br />image.png<br />其中主要的字段说明如下:

  • server.port=8066:此可视化插件的运行端口。
  • rocketmq.config.namesrvAddr=127.0.0.1:9876:rocketmq 的链接信息。 <a name="lQyXw"></a>

(3)编译插件

进入 rocketmq-externals\rocketmq-console 文件夹,执行 mvn clean package -Dmaven.test.skip=true<br /> 编译项目。<br />编译成功如下图所示:<br />image.png<br />编译阶段有可能出现以下两个问题,没有找到 mvn 命令,或编译超级慢的问题,以下提供解决方案。 <a name="umxKe"></a>

问题一:mvn 非可以运行的命令

解决方案:这是因为没有安装 Maven 或者没有配置 Maven 的环境变量导致的,下载 Maven 安装包,增加环境变量 MAVEN_HOME=maven安装目录  ,给 path 中添加 %MAVEN_HOME%\bin ,重新启动命令行工具(CMD)重新执行命令。 <a name="dLAy7"></a>

问题二:编译超慢的问题

解决方案:这是因为使用 Maven 数据源为国外源的问题导致的,只需要配置阿里的 Maven 源即可。<br />打开 Maven 目录下的 conf/setting.xml 给 mirrors 节点下添加如下内容:

1<mirror> 2 <id>alimaven</id> 3 <name>aliyun maven</name> 4 <url>http://maven.aliyun.com/nexus/content/groups/public/</url> 5 <mirrorOf>central</mirrorOf> 6</mirror>

<a name="rypWH"></a>

(4)运行插件

编译成功之后,进入 target 文件夹,执行 java -jar rocketmq-console-ng-1.0.1.jar 启动程序。<br />启动成功之后,在浏览器输入地址 http://127.0.0.1:8066 进行访问,效果如下图:<br />image.png <a name="8IIeb"></a>

3 基础使用

<a name="EEp8i"></a>

(1)添加引用 jar 包

pom.xml 添加以下代码:

1<!-- https://mvnrepository.com/artifact/com.alibaba.rocketmq/rocketmq-client --> 2<dependency> 3 <groupId>com.alibaba.rocketmq</groupId> 4 <artifactId>rocketmq-client</artifactId> 5 <version>3.6.2.Final</version> 6</dependency>

<a name="gb2C2"></a>

(2)添加生产者和消费者代码

1public class RocketMQDemo { 2 static final String MQ_NAMESRVADDR = "localhost:9876"; 3 public static void main(String[] args) { 4 // 分组名 5 String groupName = "myGroup-1"; 6 // 主题名 7 String topicName = "myTopic-1"; 8 // 标签名 9 String tagName = "myTag-1"; 10 new Thread(() -> { 11 try { 12 producer(groupName, topicName, tagName); 13 } catch (InterruptedException e) { 14 e.printStackTrace(); 15 } catch (RemotingException e) { 16 e.printStackTrace(); 17 } catch (MQClientException e) { 18 e.printStackTrace(); 19 } catch (MQBrokerException e) { 20 e.printStackTrace(); 21 } 22 }).start(); 23 new Thread(() -> { 24 try { 25 consumer(groupName, topicName, tagName); 26 } catch (MQClientException e) { 27 e.printStackTrace(); 28 } 29 }).start(); 30 } 31 32 /** 33 * @Description 生产者 34 * @Author wanglei 35 * @Param [groupName 分组名, topicName 主题名, tagName 标签名] 36 **/ 37 public static void producer(String groupName, String topicName, String tagName) throws InterruptedException, RemotingException, MQClientException, MQBrokerException { 38 DefaultMQProducer producer = new DefaultMQProducer(groupName); 39 producer.setNamesrvAddr(MQ_NAMESRVADDR); 40 producer.start(); 41 String body = "Hello, 老王"; 42 Message message = new Message(topicName, tagName, body.getBytes()); 43 producer.send(message); 44 producer.shutdown(); 45 } 46 47 /** 48 * @Description 消费者 49 * @Author wanglei 50 * @Param [groupName 分组名, topicName 主题名, tagName 标签名] 51 **/ 52 public static void consumer(String groupName, String topicName, String tagName) throws MQClientException { 53 DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(groupName); 54 consumer.setNamesrvAddr(MQ_NAMESRVADDR); 55 consumer.subscribe(topicName, tagName); 56 consumer.registerMessageListener(new MessageListenerConcurrently() { 57 @Override 58 public ConsumeConcurrentlyStatus consumeMessage( 59 List<MessageExt> msgs, ConsumeConcurrentlyContext context) { 60 for (MessageExt msg : msgs) { 61 System.out.println(new String(msg.getBody())); 62 } 63 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; 64 } 65 }); 66 consumer.start(); 67 } 68}

以上程序执行结果如下:

Hello, 老王

image.png

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )