java rocketmq--消息的產(chǎn)生(普通消息)
前言
與消息發(fā)送緊密相關的幾行代碼:
1. DefaultMQProducer producer = new DefaultMQProducer("ProducerGroupName");
2. producer.start();
3. Message msg = new Message(...)
4. SendResult sendResult = producer.send(msg);
5. producer.shutdown();
那這幾行代碼執(zhí)行時,背后都做了什么?
一. 首先是DefaultMQProducer.start
@Override
public void start() throws MQClientException {
this.defaultMQProducerImpl.start();
}
調(diào)用了默認生成消息的實現(xiàn)類 -- DefaultMQProducerImpl
調(diào)用defaultMQProducerImpl.start()方法,DefaultMQProducerImpl.start()會初始化得到MQClientInstance實例對象,MQClientInstance實例對象調(diào)用它自己的start方法會 ,啟動一些服務,如拉去消息服務PullMessageService.Start()、啟動負載平衡服務RebalanceService.Start(),比如網(wǎng)絡通信服務MQClientAPIImpl.Start()
另外,還會執(zhí)行與生產(chǎn)消息相關的信息,如注冊produceGroup、new一個TopicPublishInfo對象并以默認TopicKey為鍵值,構(gòu)成鍵值對存入DefaultMQProducerImpl的topicPublishInfoTable中。
efaultMQProducerImpl.start()后,獲取的MQClientInstance實例對象會調(diào)用sendHeartbeatToAllBroker()方法,不斷向broker發(fā)送心跳包,yin'b可以使用下面一幅圖大致描述DefaultMQProducerImpl.start()過程:

上圖中的三個部分中涉及的內(nèi)容:
1.1 初始化MQClientInstance
一個客戶端只能產(chǎn)生一個MQClientInstance實例對象,產(chǎn)生方式使用了工廠模式與單例模式。MQClientInstance.start()方法啟動一些服務,源碼如下:
public void start() throws MQClientException {
synchronized (this) {
switch (this.serviceState) {
case CREATE_JUST:
this.serviceState = ServiceState.START_FAILED;
// If not specified,looking address from name server
if (null == this.clientConfig.getNamesrvAddr()) {
this.mQClientAPIImpl.fetchNameServerAddr();
}
// Start request-response channel
this.mQClientAPIImpl.start();
// Start various schedule tasks
this.startScheduledTask();
// Start pull service
this.pullMessageService.start();
// Start rebalance service
this.rebalanceService.start();
// Start push service
this.defaultMQProducer.getDefaultMQProducerImpl().start(false);
log.info("the client factory [{}] start OK", this.clientId);
this.serviceState = ServiceState.RUNNING;
break;
case RUNNING:
break;
case SHUTDOWN_ALREADY:
break;
case START_FAILED:
throw new MQClientException("The Factory object[" + this.getClientId() + "] has been created before, and failed.", null);
default:
break;
}
}
}
1.2 注冊producer
該過程會將這個當前producer對象注冊到MQClientInstance實例對象的的producerTable中。一個jvm(一個客戶端)中一個producerGroup只能有一個實例,MQClientInstance操作producerTable大概有如下幾個方法:
- -- selectProducer
- -- updateTopicRouteInfoFromNameServer
- -- prepareHeartbeatData
- -- isNeedUpdateTopicRouteInfo
- -- shutdown
注:
根據(jù)不同的clientId,MQClientManager將給出不同的MQClientInstance;
根據(jù)不同的group,MQClientInstance將給出不同的MQProducer和MQConsumer
1.3 向路由信息表中添加路由
topicPublishInfoTable定義:
public class DefaultMQProducerImpl implements MQProducerInner {
private final Logger log = ClientLogger.getLog();
private final Random random = new Random();
private final DefaultMQProducer defaultMQProducer;
private final ConcurrentMap<String/* topic */, TopicPublishInfo> topicPublishInfoTable = new ConcurrentHashMap<String, TopicPublishInfo>();
它是一個以topic為key的Map型數(shù)據(jù)結(jié)構(gòu),DefaultMQProducerImpl.start()時會默認創(chuàng)建一個key=MixAll.DEFAULT_TOPIC的TopicPublishInfo存放到topicPublishInfoTable中。
1.4 發(fā)送心跳包
MQClientInstance向broker發(fā)送心跳包時,調(diào)用sendHeartbeatToAllBroker( ),以及從MQClientInstance實例對象的brokerAddrTable中拿到所有broker地址,向這些broker發(fā)送心跳包。
sendHeartbeatToAllBroker會涉及到prepareHeartbeatData()方法,該方法會生成heartbeatData數(shù)據(jù),發(fā)送心跳包時,heartbeatData作為心跳包的body。與producer相關的部分代碼如下:
// Producer
for (Map.Entry<String/* group */, MQProducerInner> entry : this.producerTable.entrySet()) {
MQProducerInner impl = entry.getValue();
if (impl != null) {
ProducerData producerData = new ProducerData();
producerData.setGroupName(entry.getKey());
heartbeatData.getProducerDataSet().add(producerData);
}
二、. SendResult sendResult = producer.send(msg)
首先會調(diào)用DefaultMQProducer.send(msg) ,繼而調(diào)用sendDefaultImpl:
public SendResult send(Message msg,
long timeout) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
return this.sendDefaultImpl(msg, CommunicationMode.SYNC, null, timeout);
}
sendDefaultImpl做了啥?
2.1. 獲取topicPublishInfo
根據(jù)msg的topic從topicPublishInfoTable獲取對應的topicPublishInfo,如果沒有則更新路由信息,從nameserver端拉取最新路由信息。從nameserver端拉取最新路由信息大致為:
首先getTopicRouteInfoFromNameServer,然后topicRouteData2TopicPublishInfo。

2.2 選擇消息發(fā)送的隊列
普通消息:默認方式下,selectOneMessageQueue從topicPublishInfo中的messageQueueList中選擇一個隊列(MessageQueue)進行發(fā)送消息,默認采用長輪詢的方式選擇隊列 。
它的機制如下:正常情況下,順序選擇queue進行發(fā)送;如果某一個節(jié)點發(fā)生了超時,則下次選擇queue時,跳過相同的broker。不同的隊列選擇策略形成了生產(chǎn)消息的幾種模式,如順序消息,事務消息。
順序消息:將一組需要有序消費的消息發(fā)往同一個broker的同一個隊列上即可實現(xiàn)順序消息,假設相同訂單號的支付,退款需要放到同一個隊列,那么就可以在send的時候,自己實現(xiàn)MessageQueueSelector,根據(jù)參數(shù)arg字段來選擇queue。
private SendResult sendSelectImpl(
Message msg,
MessageQueueSelector selector,
Object arg,
final CommunicationMode communicationMode,
final SendCallback sendCallback, final long timeout
) throws MQClientException, RemotingException, MQBrokerException, InterruptedException { 。。。}
事務消息:只有在消息發(fā)送成功,并且本地操作執(zhí)行成功時,才發(fā)送提交事務消息,做事務提交,消息發(fā)送失敗,直接發(fā)送回滾消息,進行回滾,具體如何實現(xiàn)后面會單獨成文分析。
2.3 封裝消息體通信包,發(fā)送數(shù)據(jù)包
首先,根據(jù)獲取的MessageQueue中的getBrokerName,調(diào)用findBrokerAddressInPublish得到該消息存放對應的broker地址,如果沒有找到則跟新路由信息,重新獲取地址 :
brokerAddrTable.get(brokerName).get(MixAll.MASTER_ID)
可知獲取的broker均為master(id=0)
然后, 將與該消息相關信息打包成RemotingCommand數(shù)據(jù)包,其RequestCode.SEND_MESSAGE
根據(jù)獲取的broke地址,將數(shù)據(jù)包到對應的broker,默認是發(fā)送超時時間為3s。
封裝消息請求包的包頭:
SendMessageRequestHeader requestHeader = new SendMessageRequestHeader(); requestHeader.setProducerGroup(this.defaultMQProducer.getProducerGroup()); requestHeader.setTopic(msg.getTopic()); requestHeader.setDefaultTopic(this.defaultMQProducer.getCreateTopicKey()); requestHeader.setDefaultTopicQueueNums(this.defaultMQProducer.getDefaultTopicQueueNums()); requestHeader.setQueueId(mq.getQueueId()); requestHeader.setSysFlag(sysFlag); requestHeader.setBornTimestamp(System.currentTimeMillis()); requestHeader.setFlag(msg.getFlag()); requestHeader.setProperties(MessageDecoder.messageProperties2String(msg.getProperties())); requestHeader.setReconsumeTimes(0); requestHeader.setUnitMode(this.isUnitMode()); requestHeader.setBatch(msg instanceof MessageBatch);
發(fā)送消息包(普通消息默認為同步方式):
SendResult sendResult = null;
switch (communicationMode) {
case SYNC:
sendResult = this.mQClientFactory.getMQClientAPIImpl().sendMessage(
brokerAddr,
mq.getBrokerName(),
msg,
requestHeader,
timeout,
communicationMode,
context,
this);
break;
處理來自broker端的響應數(shù)據(jù)包:
private SendResult sendMessageSync(
final String addr,
final String brokerName,
final Message msg,
final long timeoutMillis,
final RemotingCommand request
) throws RemotingException, MQBrokerException, InterruptedException {
RemotingCommand response = this.remotingClient.invokeSync(addr, request, timeoutMillis);
assert response != null;
return this.processSendResponse(brokerName, msg, response);
}
broker端處理request數(shù)據(jù)包后會將消息存儲到commitLog,具體過程后續(xù)分析。
以上就是本文的全部內(nèi)容,希望對大家的學習有所幫助,也希望大家多多支持腳本之家。
相關文章
使用Spring事件監(jiān)聽機制實現(xiàn)跨模塊調(diào)用的步驟詳解
Spring 事件監(jiān)聽機制是 Spring 框架中用于在應用程序的不同組件之間進行通信的一種機制,Spring 事件監(jiān)聽機制基于觀察者設計模式,使得應用程序的各個部分可以解耦,提高模塊化和可維護性,本文給大家介紹了使用Spring事件監(jiān)聽機制實現(xiàn)跨模塊調(diào)用,需要的朋友可以參考下2024-06-06
springboot jpaRepository為何一定要對Entity序列化
這篇文章主要介紹了springboot jpaRepository為何一定要對Entity序列化,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-12-12
SpringBoot通過參數(shù)注解自動獲取當前用戶信息的方法
這篇文章主要介紹了SpringBoot通過參數(shù)注解自動獲取當前用戶信息的方法,文中使用HandlerMethodArgumentResolver 類來實現(xiàn)這個功能,并通過代碼示例講解的非常詳細,需要的朋友可以參考下2024-03-03
Java中List集合去除重復數(shù)據(jù)的方法匯總
這篇文章主要給大家介紹了關于Java中List集合去除重復數(shù)據(jù)的方法,文中通過圖文介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2021-02-02
POI導出之Excel實現(xiàn)單元格的背景色填充問題
這篇文章主要介紹了POI導出之Excel實現(xiàn)單元格的背景色填充問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2023-03-03

