RocketMQ生產(chǎn)者調(diào)用start發(fā)送消息原理示例
RocketMQ發(fā)送消息
我們?cè)谑褂肦ocketMQ發(fā)送消息時(shí),一般都會(huì)使用DefaultMQProducer,類型的代碼如下:
DefaultMQProducer producer = new DefaultMQProducer("producer_group");
producer.setNamesrvAddr("42.192.50.8:9876");
try {
producer.start();
producer.send(new Message("topic", "ping".getBytes(StandardCharsets.UTF_8)));
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.shutdown();
}
上述代碼中,在消息發(fā)送之前調(diào)用了start()方法,如果不調(diào)用start()方法,直接發(fā)送消息,那么會(huì)出現(xiàn)以下報(bào)錯(cuò):

報(bào)錯(cuò)消息里面很明顯地告知我們,目前這個(gè)DefaultMQProducer狀態(tài)沒(méi)有準(zhǔn)備好,還不能發(fā)送消息。為了一探究竟,我們得去看看start()里面究竟做了什么操作呢?
start()里面究竟做了什么操作
我們根據(jù)源碼一路走下來(lái),可以追蹤到DefaultMQProducerImpl.start(final boolean startFactory)這個(gè)方法:
public void start(final boolean startFactory) throws MQClientException {
switch (this.serviceState) {
case CREATE_JUST:
this.serviceState = ServiceState.START_FAILED;
this.checkConfig();
if (!this.defaultMQProducer.getProducerGroup().equals(MixAll.CLIENT_INNER_PRODUCER_GROUP)) {
this.defaultMQProducer.changeInstanceNameToPID();
}
// 創(chuàng)建MQClientInstance
this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQProducer, rpcHook);
// 注冊(cè)Producer到MQClientInstance中
boolean registerOK = mQClientFactory.registerProducer(this.defaultMQProducer.getProducerGroup(), this);
if (!registerOK) {
this.serviceState = ServiceState.CREATE_JUST;
throw new MQClientException("The producer group[" + this.defaultMQProducer.getProducerGroup()
+ "] has been created before, specify another name please." + FAQUrl.suggestTodo(FAQUrl.GROUP_NAME_DUPLICATE_URL),
null);
}
this.topicPublishInfoTable.put(this.defaultMQProducer.getCreateTopicKey(), new TopicPublishInfo());
// 啟動(dòng)MQClientInstance實(shí)例
if (startFactory) {
mQClientFactory.start();
}
log.info("the producer [{}] start OK. sendMessageWithVIPChannel={}", this.defaultMQProducer.getProducerGroup(),
this.defaultMQProducer.isSendMessageWithVIPChannel());
this.serviceState = ServiceState.RUNNING;
break;
case RUNNING:
case START_FAILED:
case SHUTDOWN_ALREADY:
throw new MQClientException("The producer service state not OK, maybe started once, "
+ this.serviceState
+ FAQUrl.suggestTodo(FAQUrl.CLIENT_SERVICE_NOT_OK),
null);
default:
break;
}
上述代碼主要做了以下幾點(diǎn):
1.創(chuàng)建MQClientInstance實(shí)例;
2.注冊(cè)Producer到MQClientInstance實(shí)例中;
3.啟動(dòng)MQClientInstance實(shí)例;
MQClientInstance實(shí)例并不是每次都會(huì)創(chuàng)建的,它創(chuàng)建出來(lái)也會(huì)緩存的MQClientManager中,不過(guò)根據(jù)源碼來(lái)看的話,每次創(chuàng)建Producer都會(huì)對(duì)應(yīng)創(chuàng)建一個(gè)新的MQClientInstance實(shí)例,所以一般情況下不建議一個(gè)應(yīng)用服務(wù)中重復(fù)創(chuàng)建Producer;
最終start()方法的關(guān)鍵實(shí)現(xiàn)邏輯還是需要進(jìn)入MQClientInstance.start()中:
public void start() throws MQClientException {
synchronized (this) {
switch (this.serviceState) {
case CREATE_JUST:
this.serviceState = ServiceState.START_FAILED;
// 如果namesrv地址為null,那么就需要自己找namesrv地址
if (null == this.clientConfig.getNamesrvAddr()) {
this.mQClientAPIImpl.fetchNameServerAddr();
}
// 開(kāi)啟一個(gè)請(qǐng)求響應(yīng)渠道,沒(méi)猜錯(cuò)的話,應(yīng)該是netty實(shí)現(xiàn)的
this.mQClientAPIImpl.start();
// 開(kāi)啟定時(shí)任務(wù)
this.startScheduledTask();
// 開(kāi)啟拉消息服務(wù)
this.pullMessageService.start();
// 開(kāi)啟負(fù)載均衡服務(wù)
this.rebalanceService.start();
// 再開(kāi)啟一個(gè)默認(rèn)生產(chǎn)者,這個(gè)生產(chǎn)者不需要啟動(dòng)MQClientInstance實(shí)例
this.defaultMQProducer.getDefaultMQProducerImpl().start(false);
log.info("the client factory [{}] start OK", this.clientId);
this.serviceState = ServiceState.RUNNING;
break;
case START_FAILED:
throw new MQClientException("The Factory object[" + this.getClientId() + "] has been created before, and failed.", null);
default:
break;
}
}
}
看樣子,這才是start()方法真正要做的事情:
1.找namesrv地址,應(yīng)該是后面需要使用namesrv地址查詢對(duì)應(yīng)的broker;
2.開(kāi)啟Netty客戶端的初始化,包括與namesrv建立信道;另外開(kāi)啟兩個(gè)定時(shí)任務(wù),一個(gè)清除列表中過(guò)期的請(qǐng)求,第二個(gè)就是篩選可用的namesrv服務(wù);
3.開(kāi)啟一些定時(shí)任務(wù);包括如果沒(méi)有設(shè)置namesrv地址的話,會(huì)從指定站點(diǎn)拉namesrv地址;清除下線broker并發(fā)送心跳給所有的broker等工作;
4.因?yàn)楫?dāng)前是生產(chǎn)者,所以pullMessageService很快就結(jié)束;
5.生產(chǎn)者不需要做負(fù)載均衡,所以rebalanceService很快也結(jié)束;
6.給默認(rèn)創(chuàng)建的生產(chǎn)者執(zhí)行一下start()方法,其實(shí)啥也沒(méi)做;
上述大多數(shù)任務(wù)都是給消費(fèi)者使用的,作為生產(chǎn)者,唯一起作用的就是前三步,查找namesrv地址、第二步與namesrv建立通信以及第三步對(duì)broker的一些定時(shí)清理工作;不過(guò)沒(méi)有發(fā)生消息之前,是不會(huì)從遠(yuǎn)程獲取任何數(shù)據(jù)的。所以綜上所述,start()方法里面只做了以下兩件事情:
1.與namesrv建立通信渠道,它甚至都沒(méi)有從namesrv獲取任何數(shù)據(jù);
2.啟動(dòng)一些定時(shí)任務(wù),包括清理下線的broker;
小結(jié)
雖然在生產(chǎn)者中,start()方法里面真正做的事情比較少,但是卻是非常有必要的。發(fā)送消息之前,我們沒(méi)有使用start()方法,導(dǎo)致消息發(fā)送失敗,是因?yàn)樯a(chǎn)者與namesrv之間的通信渠道沒(méi)有建立。
以上就是RocketMQ生產(chǎn)者調(diào)用start發(fā)送消息原理示例的詳細(xì)內(nèi)容,更多關(guān)于RocketMQ調(diào)用start發(fā)送消息的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
Spring中@Autowired自動(dòng)注入map詳解
這篇文章主要介紹了Spring中@Autowired自動(dòng)注入map詳解, spring是支持基于接口實(shí)現(xiàn)類的直接注入的,支持注入map,list等集合中,不用做其他的配置,直接注入,需要的朋友可以參考下2023-10-10
詳解Spring batch 入門學(xué)習(xí)教程(附源碼)
本篇文章主要介紹了Spring batch 入門學(xué)習(xí)教程(附源碼),小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2017-11-11
關(guān)于Java中finalize析構(gòu)方法的作用詳解
構(gòu)造方法用于創(chuàng)建和初始化類對(duì)象,也就是說(shuō),構(gòu)造方法負(fù)責(zé)”生出“一個(gè)類對(duì)象,并可以在對(duì)象出生時(shí)進(jìn)行必要的操作,在這篇文章中會(huì)給大家簡(jiǎn)單介紹一下析構(gòu)方法,需要的朋友可以參考下2023-05-05
Java?熱更新?Groovy?實(shí)踐及踩坑指南(推薦)
Apache的Groovy是Java平臺(tái)上設(shè)計(jì)的面向?qū)ο缶幊陶Z(yǔ)言,這門動(dòng)態(tài)語(yǔ)言擁有類似Python、Ruby和Smalltalk中的一些特性,可以作為Java平臺(tái)的腳本語(yǔ)言使用,這篇文章主要介紹了Java?熱更新?Groovy?實(shí)踐及踩坑指南,需要的朋友可以參考下2022-09-09
Java開(kāi)發(fā)工具IntelliJ IDEA安裝圖解
這篇文章主要介紹了Java開(kāi)發(fā)工具IntelliJ IDEA安裝圖解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-11-11
Intellj?idea新建的java源文件夾不是藍(lán)色的圖文解決辦法
idea打開(kāi)java項(xiàng)目后新建的模塊中,java文件夾需要變成藍(lán)色,這篇文章主要給大家介紹了關(guān)于Intellj?idea新建的java源文件夾不是藍(lán)色的相關(guān)資料,文中通過(guò)圖文介紹的非常詳細(xì),需要的朋友可以參考下2024-02-02
深入探究 spring-boot-starter-parent的作用
這篇文章主要介紹了spring-boot-starter-parent的作用詳解,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,感興趣的小伙伴可以跟著小編一起來(lái)學(xué)習(xí)一下2023-05-05

