| |
|
开发:
C++知识库
Java知识库
JavaScript
Python
PHP知识库
人工智能
区块链
大数据
移动开发
嵌入式
开发工具
数据结构与算法
开发测试
游戏开发
网络协议
系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程 数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁 |
-> 大数据 -> 【软件安装】RocketMQ在Lunix系统中的安装 -> 正文阅读 |
|
[大数据]【软件安装】RocketMQ在Lunix系统中的安装 |
0、安装前环境准备本篇是基于Linux操作系统中的安装,故先准备一台干净的Linux操作系统,并在系统上先安装好JDK和Maven,本文中所有的操作基于CentOS8进行安装演示; 1、官网下载RocketMQ安装包cd /usr/local/ mkdir source cd source/ wget https://dlcdn.apache.org/rocketmq/4.9.2/rocketmq-all-4.9.2-bin-release.zip 2、解压RocketMQ安装包,阅读其中的README.md文件unzip rocketmq-all-4.9.2-bin-release.zip cd rocketmq-4.9.2/ vim README.md 3、启动RocketMQ服务cd bin/ ./mqnamesrv 4、新起一个终端窗口,修改runbroker.sh文件cd?/usr/local/source/rocketmq-4.9.2/bin vim runbroker.sh 默认设置的虚拟机内存参数过大,正常人的虚拟机根本跑不起来,改小点可以启动。 ./mqbroker -n localhost:9876 ?5、新起一个终端窗口,测试RocketMQ功能修改tools.sh文件 cd?/usr/local/source/rocketmq-4.9.2/bin vim tools.sh? 执行?./tools.sh org.apache.rocketmq.example.quickstart.Producer 启动MQ生产者 新起一个终端窗口启动MQ消费者 cd?/usr/local/source/rocketmq-4.9.2/bin ./tools.sh org.apache.rocketmq.example.quickstart.Consumer 6、RocketMQ中各角色的解读NameServer:主要的功能是用来收集其它角色的信息,相当于一个中介,维护了一个服务的列表,知道有哪些服务还存活。底层由netty实现,提供了路由管理、服务注册、服务发现的功能,是一个无状态节点; Broker:面向Producer和Consumer接受和发送消息,向NameServer提交自己的信息。是消息中间件的消息存储、转发服务器。每个Broker节点在启动时,都会遍历NameServer列表,于每个NameServer建立长连接,注册自己的信息,之后定时上报; Producer:消息生产者。通过集群中一个节点建立长连接,获得Topic的路由信息,包括Topic下面有哪些Queue,这些Queue分布在哪些Broker上等; Consumer:消息消费者。通过NameServer集群获得Topic的路由信息,连接到对应的Broker上消费消息。注意,由于Master和Slave都可以读取消息,因此Consumer会与Master和Slave都建立连接。 7、RocketMQ的HelloWorld生产者Producer代码: package com.feenix.rocketmq.producer; 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.exception.RemotingException; import java.nio.charset.StandardCharsets; /** * @Author: Feenix * @CreateTime: 2022-02-24 16:20 * @Version: 1.0 * @Description: */ public class Producer { private static final String PRODUCER_GROUP = "feenix_group"; private static final String NAMESRV_ADDR = "192.168.159.149:9876"; private static final String TOPIC = "feenix_topic"; public static void main(String[] args) throws MQClientException, MQBrokerException, RemotingException, InterruptedException { DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP); // 设置NameServer的地址 producer.setNamesrvAddr(NAMESRV_ADDR); producer.start(); // Topic >> 消息将要发送到的地址 // body >> 真正流转的消息内容 String body = "Hello World By RocketMQ!"; Message message = new Message(TOPIC, body.getBytes()); SendResult result = producer.send(message); System.out.println("SendResult:" + result); // producer.shutdown(); } } 消费者Consumer代码: package com.feenix.rocketmq.consumer; 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 java.util.List; /** * @Author: Feenix * @CreateTime: 2022-02-24 16:44 * @Version: 1.0 * @Description: */ public class Consumer { private static final String CONSUMER_GROUP = "feenix_consumer_group"; private static final String NAMESRV_ADDR = "192.168.159.149:9876"; private static final String TOPIC = "feenix_topic"; private static final String FILTER = "*"; public static void main(String[] args) throws MQClientException { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP); consumer.setNamesrvAddr(NAMESRV_ADDR); // 每个consumer只能订阅一个topic consumer.subscribe(TOPIC, FILTER); // 注册消息监听器 consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> messageExtList, ConsumeConcurrentlyContext consumeConcurrentlyContext) { for (MessageExt messageExt : messageExtList) { System.out.println(new String(messageExt.getBody())); } // 默认情况下这条消息只会被一个consumer消费 // 当被消费之后,告诉服务器消费成功,broker会将成功消费的消息剔除掉不再消费 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); System.out.println("consumer start....."); } } |
|
|
上一篇文章 下一篇文章 查看所有文章 |
|
开发:
C++知识库
Java知识库
JavaScript
Python
PHP知识库
人工智能
区块链
大数据
移动开发
嵌入式
开发工具
数据结构与算法
开发测试
游戏开发
网络协议
系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程 数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁 |
360图书馆 购物 三丰科技 阅读网 日历 万年历 2025年1日历 | -2025/1/16 23:59:59- |
|
网站联系: qq:121756557 email:121756557@qq.com IT数码 |