RocketMQ 學習

来源:https://www.cnblogs.com/donleo123/archive/2023/03/30/17273204.html
-Advertisement-
Play Games

前言 RocketMQ是阿裡巴巴旗下一款開源的MQ框架,經歷過雙十一考驗、Java編程語言實現,有非常好完整生態系統。RocketMQ作為一款純java、分散式、隊列模型的開源消息中間件,支持事務消息、順序消息、批量消息、定時消息、消息回溯等 本篇文章第一部分屬於一些核心概念和工作流程的講解;第二部 ...


前言

  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下)。

  1. cd rocketmq-all-4.7.1-bin-release/conf/2m-2s-sync/
  2. vim broker-a.properties
  3. 在最後加一行diskMaxUsedSpaceRatio=99(所有節點的配置文件都加一下),表示剩餘磁碟比例不足99才報錯
  4. :wq 保存退出

  5. 重啟mq

 重新發送消息Ok了

 

 

作者:donleo123 出處:https://www.cnblogs.com/donleo123/ 本文如對您有幫助,還請多推薦下此文,如有錯誤歡迎指正,相互學習,共同進步。
您的分享是我們最大的動力!

-Advertisement-
Play Games
更多相關文章
  • React Router 備忘清單 IT寶庫整理的React Router開發速查清單適合初學者的綜合 React Router 6.x 備忘清單入門,為開發人員分享快速參考備忘單。 開發速查表大綱 入門 安裝使用 添加路由器 根路由 處理未找到錯誤 contacts 用戶界面 嵌套路由 客戶端路由 ...
  • Redis 備忘清單 IT寶庫整理的Redis開發速查備忘清單 - 本備忘單旨在快速理解 redis 所涉及的主要概念,提供了最常用的SQL語句,供您參考。入門,為開發人員分享快速參考備忘單。 開發速查表大綱 入門 介紹 小試 數據類型 Redis服務相關的命令設置 COMMAND 一些引用(可能有 ...
  • 事件系統 文章為本人理解,如有理解不到位之處,煩請各位指正。 @ Qt的事件迴圈,應該是所有Qter都避不開的一個點,所以,這篇博客,咱們來瞭解源碼中一些關於Qt中事件迴圈的部分。 先拋出幾個疑問,根據源代碼,下麵一一進行解析。 事件迴圈是什麼? 事件是怎麼產生的? 事件是如何處理的? 什麼是事件循 ...
  • 自2022年11月30日 OpenAI 發佈 ChatGPT 以來,雖然時有唱衰的聲音出現,但在OpenAI不斷推陳出新,陸續發佈了OpenAPI、GPT-4、ChatGPT Plugins之後,似乎讓大家看到了一個聊天機器人往操作系統入口進軍的升緯之路。 ChatGPT能被認為是操作系統級別的入口 ...
  • C++中的explicit關鍵字只能用於修飾只有一個參數的類構造函數,它的作用是表明該構造函數是顯示的,而非隱式的,跟它相對應的另一個關鍵字是implicit,意思是隱藏的,類構造函數預設情況下即聲明為implicit(隱式)。 那麼顯示聲明的構造函數和隱式聲明的有什麼區別呢? 來看下麵的例子: c ...
  • Java重寫toString的意義 一.toString()方法 toString()方法在Object類里定義的,其返回值類型為String類型,返回類名和它的引用地址. 在進行String類與其他類型的連接操作時,自動調用toString()方法,demo如下: Date time = new ...
  • 一、問題引入 單鏈表的實現【01】:Student-Management-System 只體現了項目功能實現,未對代碼部分做出說明。 故新增隨筆進行補充說明代碼部分。 重構代碼,迭代版本:Student Mangement System(Version 2.0) 二、解決過程 基於單鏈表實現就離不開 ...
  • 在java,c#類的成員修飾符包括,公有、私有、程式集可用的、受保護的。 對於python來說,只有兩個成員修飾符:公有成員,私有成員 成員修飾符是來修飾誰呢?當然是修飾成員了。那麼python類的成員包括什麼呢? python成員: 欄位,方法,屬性 每個類成員的修飾符有兩種: 公有成員:內部外部 ...
一周排行
    -Advertisement-
    Play Games
  • 概述:在C#中,++i和i++都是自增運算符,其中++i先增加值再返回,而i++先返回值再增加。應用場景根據需求選擇,首碼適合先增後用,尾碼適合先用後增。詳細示例提供清晰的代碼演示這兩者的操作時機和實際應用。 在C#中,++i 和 i++ 都是自增運算符,但它們在操作上有細微的差異,主要體現在操作的 ...
  • 上次發佈了:Taurus.MVC 性能壓力測試(ap 壓測 和 linux 下wrk 壓測):.NET Core 版本,今天計劃準備壓測一下 .NET 版本,來測試並記錄一下 Taurus.MVC 框架在 .NET 版本的性能,以便後續持續優化改進。 為了方便對比,本文章的電腦環境和測試思路,儘量和... ...
  • .NET WebAPI作為一種構建RESTful服務的強大工具,為開發者提供了便捷的方式來定義、處理HTTP請求並返迴響應。在設計API介面時,正確地接收和解析客戶端發送的數據至關重要。.NET WebAPI提供了一系列特性,如[FromRoute]、[FromQuery]和[FromBody],用 ...
  • 原因:我之所以想做這個項目,是因為在之前查找關於C#/WPF相關資料時,我發現講解圖像濾鏡的資源非常稀缺。此外,我註意到許多現有的開源庫主要基於CPU進行圖像渲染。這種方式在處理大量圖像時,會導致CPU的渲染負擔過重。因此,我將在下文中介紹如何通過GPU渲染來有效實現圖像的各種濾鏡效果。 生成的效果 ...
  • 引言 上一章我們介紹了在xUnit單元測試中用xUnit.DependencyInject來使用依賴註入,上一章我們的Sample.Repository倉儲層有一個批量註入的介面沒有做單元測試,今天用這個示例來演示一下如何用Bogus創建模擬數據 ,和 EFCore 的種子數據生成 Bogus 的優 ...
  • 一、前言 在自己的項目中,涉及到實時心率曲線的繪製,項目上的曲線繪製,一般很難找到能直接用的第三方庫,而且有些還是定製化的功能,所以還是自己繪製比較方便。很多人一聽到自己畫就害怕,感覺很難,今天就分享一個完整的實時心率數據繪製心率曲線圖的例子;之前的博客也分享給DrawingVisual繪製曲線的方 ...
  • 如果你在自定義的 Main 方法中直接使用 App 類並啟動應用程式,但發現 App.xaml 中定義的資源沒有被正確載入,那麼問題可能在於如何正確配置 App.xaml 與你的 App 類的交互。 確保 App.xaml 文件中的 x:Class 屬性正確指向你的 App 類。這樣,當你創建 Ap ...
  • 一:背景 1. 講故事 上個月有個朋友在微信上找到我,說他們的軟體在客戶那邊隔幾天就要崩潰一次,一直都沒有找到原因,讓我幫忙看下怎麼回事,確實工控類的軟體環境複雜難搞,朋友手上有一個崩潰的dump,剛好丟給我來分析一下。 二:WinDbg分析 1. 程式為什麼會崩潰 windbg 有一個厲害之處在於 ...
  • 前言 .NET生態中有許多依賴註入容器。在大多數情況下,微軟提供的內置容器在易用性和性能方面都非常優秀。外加ASP.NET Core預設使用內置容器,使用很方便。 但是筆者在使用中一直有一個頭疼的問題:服務工廠無法提供請求的服務類型相關的信息。這在一般情況下並沒有影響,但是內置容器支持註冊開放泛型服 ...
  • 一、前言 在項目開發過程中,DataGrid是經常使用到的一個數據展示控制項,而通常表格的最後一列是作為操作列存在,比如會有編輯、刪除等功能按鈕。但WPF的原始DataGrid中,預設只支持固定左側列,這跟大家習慣性操作列放最後不符,今天就來介紹一種簡單的方式實現固定右側列。(這裡的實現方式參考的大佬 ...