前言
RocketMQ是阿里巴巴旗下一款開源的MQ框架,經歷過雙十一考驗、Java編程語言實作,有非常好完整生態系統,RocketMQ作為一款純java、分布式、佇列模型的開源訊息中間件,支持事務訊息、順序訊息、批量訊息、定時訊息、訊息回溯等
本篇文章第一部分屬于一些核心概念和作業流程的講解;第二部分就是純手動搭建了一套環境;第三部分是基于環境進行測驗和集成到SpringBoot
核心概念
-
NameServer:可以理解為是一個注冊中心,主要是用來保存topic路由資訊,管理Broker,在NameServer的集群中,NameServer與NameServer之間是沒有任何通信的,
-
Broker:核心的一個角色,主要是用來保存topic的資訊,接受生產者產生的訊息,持久化訊息,在一個Broker集群中,相同的BrokerName可以稱為一個Broker組,一個Broker組中,BrokerId為0的為主節點,其它的為從節點,BrokerName和BrokerId是可以在Broker啟動時通過組態檔配置的,每個Broker組只存放一部分訊息,
-
生產者:生產訊息的一方就是生產者
-
生產者組:一個生產者組可以有很多生產者,只需要在創建生產者的時候指定生產者組,那么這個生產者就在那個生產者組
-
消費者:用來消費生產者訊息的一方
-
消費者組:跟生產者一樣,每個消費者都有所在的消費者組,一個消費者組可以有很多的消費者,不同的消費者組消費訊息是互不影響的,
-
topic(主題):可以理解為一個訊息的集合的名字,生產者在發送訊息的時候需要指定發到哪個topic下,消費者消費訊息的時候也需要知道自己消費的是哪些topic底下的訊息,
-
Tag(子主題):比topic低一級,可以用來區分同一topic下的不同業務型別的訊息,發送訊息的時候也需要指定,
這里有組的概念是因為可以用來做到不同的生產者組或者消費者組有不同的配置,這樣就可以使得生產者或者消費者更加靈活,
作業流程

通過這張圖就可以很清楚的知道,RocketMQ大致的作業流程:
-
Broker啟動的時候,會往每臺NameServer(因為NameServer之間不通信,所以每臺都得注冊)注冊自己的資訊,這些資訊包括自己的ip和埠號,自己這臺Broker有哪些topic等資訊,
-
Producer在啟動之后會跟會NameServer建立連接,定期從NameServer中獲取Broker的資訊,當發送訊息的時候,會根據訊息需要發送到哪個topic去找對應的Broker地址,如果有的話,就向這臺Broker發送請求;沒有找到的話,就看根據是否允許自動創建topic來決定是否發送訊息,
-
Broker在接收到Producer的訊息之后,會將訊息存起來,持久化,如果有從節點的話,也會主動同步給從節點,實作資料的備份
-
Consumer啟動之后也會跟NameServer建立連接,定期從NameServer中獲取Broker和對應topic的資訊,然后根據自己需要訂閱的topic資訊找到對應的Broker的地址,然后跟Broker建立連接,獲取訊息,進行消費
就跟上面的圖一樣,整體的作業流程還是比較簡單的,這里簡化了很多概念,主要是為了好理解,
環境搭建
通過上面分析,我們知道,在RocketMQ中有NameServer、Broker、生產者、消費者四種角色,而生產者和消費者實際上就是業務系統,所以這里不需要搭建,真正要搭建的就是NameServer和Broker,但是為了方便RocketMQ資料的可視化,這里多搭建一套可視化的服務,
搭建程序比較簡單,按照步驟一步一步來就可以完成,如果提示一些命令不存在,那么直接通過yum安裝這些命令就行,
一、準備
需要準備一個linux服務器,需要先安裝好JDK
關閉防火墻
systemctl stop firewalld
systemctl disable firewalld
下載并解壓RocketMQ
1、創建一個目錄,用來存放rocketmq相關的東西
mkdir /usr/rocketmq cd /usr/rocketmq
2、下載并解壓rocketmq
下載
wget https://archive.apache.org/dist/rocketmq/4.7.1/rocketmq-all-4.7.1-bin-release.zip
解壓
unzip rocketmq-all-4.7.1-bin-release.zip
如果提示unzip: Command Not Found
通過yum命令安裝,如果已經安裝了,請忽略
yum install -y unzip zip
看到這一個檔案夾就完成了

然后進入rocketmq-all-4.7.1-bin-release檔案夾
cd rocketmq-all-4.7.1-bin-release
RocketMQ的東西都在這了

二、搭建NameServer
在啟動NameServer之前,強烈建議修改一下啟動時的jvm引數,因為默認的引數都比較大,為了避免記憶體不夠,建議修改小,當然,如果你的記憶體足夠大,可以忽略,
vi bin/runserver.sh
修改畫圈的這一行

可以設定小一點
-server -Xms512m -Xmx512m -Xmn256m -XX:MetaspaceSize=32m -XX:MaxMetaspaceSize=50m
啟動NameServer
修改完之后,執行如下命令就可以啟動NameServer了
nohup sh bin/mqnamesrv &
查看NameServer日志
tail -f ~/logs/rocketmqlogs/namesrv.log
如果看到如下的日志,就說明啟動成功了

關閉NameServer
sh bin/mqshutdown namesrv
三、搭建Broker
這里啟動單機版的Broker
修改jvm引數
跟啟動NameServer一樣,也建議去修改jvm引數
vi bin/runbroker.sh
將畫圈的地方設定小點,當然也別太小啊

可以這樣設定
-server -Xms1g -Xmx1g -Xmn512m
修改Broker組態檔broker.conf
這里需要改一下Broker組態檔,需要指定NameServer的地址,因為需要Broker需要往NameServer注冊
vi conf/broker.conf
Broker組態檔

這里就能看出Broker的配置了,什么Broker集群的名稱啊,Broker的名稱啊,Broker的id啊,都跟前面說的對上了,
在檔案末尾追加地址
namesrvAddr = localhost:9876
因為NameServer跟Broker在同一臺機器,所以是localhost,NameServer埠默認的是9876,
不過這里我還建議再修改一處資訊,因為Broker向NameServer進行注冊的時候,帶過去的ip如果不指定就會自動獲取,但是自動獲取的有個坑,就是有可能你的電腦無法訪問到這個自動獲取的ip,所以我建議手動指定你的電腦可以訪問到的服務器ip,
我的虛擬機的ip是192.168.3.158,所以就指定為192.168.3.158,如下
brokerIP1 = 192.168.3.158 brokerIP2 = 192.168.3.158
開啟自動創建Topic
autoCreateTopicEnable = true
如果以上都配置的話,最終的組態檔應該如下,紅圈的為新加的

啟動Broker
nohup sh bin/mqbroker -c conf/broker.conf &
-c 引數就是指定組態檔
查看日志
tail -f ~/logs/rocketmqlogs/broker.log
當看到如下日志就說明啟動成功了

關閉Broker
sh bin/mqshutdown broker
查看Broker 與NameServer是否運行
jps

說明Broker與NameServer是運行狀態
四、搭建可視化控制臺
其實前面NameServer和Broker搭建完成之后,就可以用來收發訊息了,但是為了更加直觀,可以搭一套可視化的服務,
可視化服務其實就是一個jar包,啟動就行了,
jar包可以從這獲取
鏈接:https://pan.baidu.com/s/16s1qwY2qzE2bxR81t5Wm6w
提取碼:s0sd
將jar包上傳到服務器,放到/usr/rocketmq的目錄底下,當然放哪都無所謂,這里只是為了方便,因為rocketmq的東西都在這里
然后進入/usr/rocketmq下,執行如下命名
nohup java -jar -server -Xms256m -Xmx256m -Drocketmq.config.namesrvAddr=localhost:9876 -Dserver.port=8088 rocketmq-console-ng-1.0.1.jar &
rocketmq.config.namesrvAddr就是用來指定NameServer的地址的
查看日志
tail -f ~/logs/consolelogs/rocketmq-console.log
當看到如下日志,就說明啟動成功了

然后在瀏覽器中輸入http://linux服務器的ip:8088/就可以看到控制臺了,如果無法訪問,可以看看防火墻有沒有關閉

通過控制臺可以查看生產者、消費者、Broker集群等資訊,非常直觀,
功能很多,這里就不一一介紹了,
停止命令
查看行程
1.jps
- -q:只輸出行程 ID
- -m:輸出傳入 main 方法的引數
- -l:輸出完全的包名,應用主類名,jar的完全路徑名
- -v:輸出jvm引數
- -V:輸出通過flag檔案傳遞到JVM中的引數
2.ps aux | grep java 來獲取java行程 id
結束行程
kill pid 或者(kill -9 pid)
- pid: jar包行程號
- kill pid: 結束行程,有局限性,例如后臺行程,守護行程等,不能結束
- kill - 9 pid : 表示強制殺死該行程;
測驗
環境搭好之后,就可以進行測驗了,
引入依賴
<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-client</artifactId> <version>4.7.1</version> </dependency>
生產者發送訊息
@Test public void sendTest() throws Exception{ //創建一個生產者,指定生產者組為ldProducer DefaultMQProducer producer = new DefaultMQProducer("ldProducer"); // 指定NameServer的地址 producer.setNamesrvAddr("192.168.3.158:9876"); // 第一次發送可能會超時,我設定的比較大 producer.setSendMsgTimeout(60000); // 啟動生產者 producer.start(); // 創建一條訊息 // topic為 ldTopic // 訊息內容為 java學習日記 // tags 為 TagA Message msg = new Message("ldTopic", "TagA", "java學習日記 ".getBytes(RemotingHelper.DEFAULT_CHARSET)); // 發送訊息并得到訊息的發送結果,然后列印 SendResult sendResult = producer.send(msg); System.out.printf("%s%n", sendResult); // 關閉生產者 producer.shutdown(); }
- 構建一個訊息生產者DefaultMQProducer實體,然后指定生產者組為ldProducer;
- 指定NameServer的地址:服務器的ip:9876,因為需要從NameServer拉取Broker的資訊
- producer.start() 啟動生產者
- 構建一個內容為三友的java日記的訊息,然后指定這個訊息往ldTopic這個topic發送
- producer.send(msg):發送訊息,列印結果
- 關閉生產者
消費者消費訊息
public class ConsumerMsg { public static void main(String[] args) throws Exception { // 通過push模式消費訊息,指定消費者組 DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ldConsumer"); consumer.setNamesrvAddr("192.168.3.158:9876"); // 訂閱這個topic下的所有的訊息 consumer.subscribe("ldTopic", "*"); // 注冊一個消費的監聽器,當有訊息的時候,會回呼這個監聽器來消費訊息 consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) { for (MessageExt msg : msgs) { System.out.printf("消費訊息:%s", new String(msg.getBody()) + "\n"); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); System.out.printf("Consumer Started.%n"); } }
- 創建一個消費者實體物件,指定消費者組為ldConsumer
- 指定NameServer的地址:服務器的ip:9876
- 訂閱 ldTopic 這個topic的所有資訊
- consumer.registerMessageListener ,這個很重要,是注冊一個監聽器,這個監聽器是當有訊息的時候就會回呼這個監聽器,處理訊息,所以需要用戶實作這個介面,然后處理訊息,
- 啟動消費者
啟動之后,消費者就會消費剛才生產者發送的訊息,于是控制臺就列印出如下資訊

再去看控制臺,已消費

SpringBoot環境下集成RocketMQ
集成
在實際專案中肯定不會像上面測驗那樣用,都是集成SpringBoot的,
1、引入依賴
<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.1.1</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <version>2.1.1.RELEASE</version> </dependency>
2、yml配置
rocketmq:
producer:
group: ldProducer
name-server: 192.168.3.158:9876
3、創建消費者
SpringBoot底下只需要實作RocketMQListener介面,然后加上@RocketMQMessageListener注解即可
@Component @RocketMQMessageListener(consumerGroup = "ldConsumer", topic = "ldDelayTaskTopic") @Slf4j public class LdRocketMQListener implements RocketMQListener<String> { @Override public void onMessage(String msg) { log.info("獲取到延遲任務訊息:{}",msg); } }
@RocketMQMessageListener需要指定消費者屬于哪個消費者組,消費哪個topic,NameServer的地址已經通過yml組態檔配置類
4、測驗
@RestController @Slf4j public class RocketMQDelayTaskController { @Resource private DefaultMQProducer producer; @GetMapping("/rocketmq/add") public void addTask(@RequestParam("task") String task) throws Exception { Message msg = new Message("ldDelayTaskTopic", "TagA", task.getBytes(RemotingHelper.DEFAULT_CHARSET)); msg.setDelayTimeLevel(2); // 發送訊息并得到訊息的發送結果,然后列印 log.info("提交延遲任務"); producer.send(msg); } }
可能遇到的問題
搭完mq單主單從集群之后,美滋滋想發一下message, 沒想到碰到一個坑爹的問題:
Caused by: org.apache.rocketmq.client.exception.MQBrokerException: CODE: 14 DESC: service not available now, maybe disk full, CL: 0.90 CQ: 0.90 INDEX: 0.90, maybe your broker machine memory too small.
org.apache.rocketmq.client.exception.MQClientException: Send [3] times, still failed, cost [549]ms, Topic: ldTopicA, BrokersSent: [broker-a, broker-a, broker-a] See http://rocketmq.apache.org/docs/faq/ for further details. at org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl.sendDefaultImpl(DefaultMQProducerImpl.java:665) at org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl.send(DefaultMQProducerImpl.java:1343) at org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl.send(DefaultMQProducerImpl.java:1289) at org.apache.rocketmq.client.producer.DefaultMQProducer.send(DefaultMQProducer.java:325) at com.example.delay.MQTest.sendTest(MQTest.java:46) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
... Caused by: org.apache.rocketmq.client.exception.MQBrokerException: CODE: 14 DESC: service not available now, maybe disk full, CL: 0.90 CQ: 0.90 INDEX: 0.90, maybe your broker machine memory too small. For more information, please visit the url, http://rocketmq.apache.org/docs/faq/ at org.apache.rocketmq.client.impl.MQClientAPIImpl.processSendResponse(MQClientAPIImpl.java:665) at org.apache.rocketmq.client.impl.MQClientAPIImpl.sendMessageSync(MQClientAPIImpl.java:505) at org.apache.rocketmq.client.impl.MQClientAPIImpl.sendMessage(MQClientAPIImpl.java:487) at org.apache.rocketmq.client.impl.MQClientAPIImpl.sendMessage(MQClientAPIImpl.java:431) at org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl.sendKernelImpl(DefaultMQProducerImpl.java:854) at org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl.sendDefaultImpl(DefaultMQProducerImpl.java:584)
看報錯應該是磁盤空間不足的問題,看到一個帖子https://bbs.csdn.net/topics/392568834,還挺符合的,雖然給出的解決方案說的沒那么詳細,但是值得一試,
查看磁盤空間

已用91%,查閱百度之后發現rocketmq原始碼的DefaultMessageStore類里,默認會把剩余磁盤的比率不足75%(rocketmq版本不同這個比率好像不一樣)當做磁盤空間不足處理,看來磁盤是有點不夠了,
先cd到rocketmq組態檔的路徑,我這里配置的是雙主雙從同步的模式,所以cd到組態檔(根據配置的不同檔案夾的路徑不一樣,但都在/conf下),
- cd rocketmq-all-4.7.1-bin-release/conf/2m-2s-sync/
- vim broker-a.properties
- 在最后加一行diskMaxUsedSpaceRatio=99(所有節點的組態檔都加一下),表示剩余磁盤比例不足99才報錯
-
:wq 保存退出
- 重啟mq

重新發送訊息Ok了
作者:donleo123 出處:https://www.cnblogs.com/donleo123/ 本文如對您有幫助,還請多推薦下此文,如有錯誤歡迎指正,相互學習,共同進步,
轉載請註明出處,本文鏈接:https://www.uj5u.com/houduan/548699.html
標籤:Java
