前言
NameServer是整個RocketMQ的路由中心,功能類似于Zookeeper,用于服務(wù)注冊和服務(wù)發(fā)現(xiàn),是輕量級別的Zookeeper。
NameServer核心的路由表
private final HashMap<String/*topic*/, List<QueueData>> topicQueueTable;
private final HashMap<String/*brokerName*/, BroKerData> brokerAddrTable;
private final HashMap<String/*ClusterName*/, Set<String/*brokerName*/>> clusterAddrTable;
private final HashMap<String/*brokerAddr*/, BrokerLiveInfo> brokerLiveTable;
private final HashMap<String/*brokerAddr*/, List<String>/*Filter Server*/> filterServerTable;
<img src="https://user-gold-cdn.xitu.io/2020/7/8/1732c4da696d92bf?w=2196&h=2240&f=png&s=177045" alt="關(guān)鍵的表.png" title="關(guān)鍵的表.png" />
topicQueueTable:topic表示消息類型,消息發(fā)送時根據(jù)路由表進行負載均衡。
brokerAddrTable:Broker基礎(chǔ)信息:Broker名稱,所屬集群名稱,主從Broker地址。
clusterAddrTable:Broker集群信息,存儲集群中所有的Broker名稱。
brokerLiveTable:Broker狀態(tài)信息,NameServer每次收到Broker的心跳包時都會更新對應(yīng)的Broker信息。NameServer的心跳檢測主要通過掃面該表完成。
其中關(guān)鍵的四張表的結(jié)構(gòu)如上圖所示。在了解表結(jié)構(gòu)之后,可以更加容易理解源碼中對于表的操作。
路由注冊
Broker啟動時回向集群中所有的NameServer發(fā)送心跳請求,并且每隔30秒都會向所有的NameServer發(fā)送心跳包,Name Server收到Broker的心跳包時會更新brokerLiveTable中對應(yīng)Broker的BrokerLiveInfo的lastUpdateTimeStamp屬性,并且NameServer會每個十秒掃面brokerLiveTable,如果連續(xù)120s沒有收到某個Broker的心跳包,則將所有該Broker對應(yīng)的路由信息同時關(guān)閉Socket通道。
Broker發(fā)送心跳包
Broker在啟動時會創(chuàng)建一個包含單線程的線程池ScheduledThreadPoolExector對象用于執(zhí)行定時任務(wù),每30秒向所有的NameServer發(fā)送心跳包。如下面代碼所示。
/*BrokerController#start*/
// Broker每隔30秒向NameServer注冊中心發(fā)送心跳包。
this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
@Override
public void run() {
try {
BrokerController.this.registerBrokerAll(true, false, brokerConfig.isForceRegister());
} catch (Throwable e) {
log.error("registerBrokerAll Exception", e);
}
}
}, 1000 * 10, Math.max(10000, Math.min(brokerConfig.getRegisterNameServerPeriod(), 60000)), TimeUnit.MILLISECONDS);
/*BrokerouterAPI#registerBrokerAll*/
// 獲取所有的NameServer服務(wù)器的地址列表
List<String> nameServerAddressList = this.remotingClient.getNameServerAddressList();
// 遍歷所有NameServer服務(wù)器的地址列表
for (final String namesrvAddr : nameServerAddressList) {
brokerOuterExecutor.execute(new Runnable() {
@Override
public void run() {
try {
// 向NameServer服務(wù)器發(fā)送心跳包
RegisterBrokerResult result = registerBroker(namesrvAddr,oneway, timeoutMills,requestHeader,body);
if (result != null) {
registerBrokerResultList.add(result);
}
log.info("register broker[{}]to name server {} OK", brokerId, namesrvAddr);
} catch (Exception e) {
log.warn("registerBroker Exception, {}", namesrvAddr, e);
} finally {
countDownLatch.countDown();
}
}
});
}
/*BrokerouterAPI#registerBroker*/
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.REGISTER_BROKER, requestHeader);
request.setBody(body);
if (oneway) {
try {
this.remotingClient.invokeOneway(namesrvAddr, request, timeoutMills);
} catch (RemotingTooMuchRequestException e) {
// Ignore
}
return null;
}
NameServer處理來自Broker的心跳包
/*DefaultRequestProcessor#processRequest*/
switch (request.getCode()) {
case RequestCode.PUT_KV_CONFIG:
return this.putKVConfig(ctx, request);
case RequestCode.GET_KV_CONFIG:
return this.getKVConfig(ctx, request);
case RequestCode.DELETE_KV_CONFIG:
return this.deleteKVConfig(ctx, request);
case RequestCode.QUERY_DATA_VERSION:
return queryBrokerTopicConfig(ctx, request);
// 走這里
case RequestCode.REGISTER_BROKER:
Version brokerVersion = MQVersion.value2Version(request.getVersion());
if (brokerVersion.ordinal() >= MQVersion.Version.V3_0_11.ordinal()) {
return this.registerBrokerWithFilterServer(ctx, request);
} else {
return this.registerBroker(ctx, request);
}
case RequestCode.UNREGISTER_BROKER:
return this.unregisterBroker(ctx, request);
...
case RequestCode.UPDATE_NAMESRV_CONFIG:
return this.updateConfig(ctx, request);
case RequestCode.GET_NAMESRV_CONFIG:
return this.getConfig(ctx, request);
default:
break;
}
/*DefaultRequestProcessor#registerBroker*/
final RegisterBrokerRequestHeader requestHeader =
(RegisterBrokerRequestHeader) request.decodeCommandCustomHeader(RegisterBrokerRequestHeader.class);
RegisterBrokerResult result = this.namesrvController.getRouteInfoManager().registerBroker(
requestHeader.getClusterName(),
requestHeader.getBrokerAddr(),
requestHeader.getBrokerName(),
requestHeader.getBrokerId(),
requestHeader.getHaServerAddr(),
topicConfigWrapper,
null,
ctx.channel()
);
調(diào)用RouteInfoManager類的registerBroker()方法更新緩存
/*RouteInfoManager#registerBroker*/
// 防止并發(fā)修改
this.lock.writeLock().lockInterruptibly();
//--------------更新clusterAddrTable緩存-----------
Set<String> brokerNames = this.clusterAddrTable.get(clusterName);
// 如果為空,則Broker第一次注冊
if (null == brokerNames) {
brokerNames = new HashSet<String>();
this.clusterAddrTable.put(clusterName, brokerNames);
}
brokerNames.add(brokerName);
//------------更新brokerAddrTable緩存------------
boolean registerFirst = false;
BrokerData brokerData = this.brokerAddrTable.get(brokerName);
// 如果通過brokerName無法得到brokerData,說明是第一次注冊
if (null == brokerData) {
// 第一次注冊
registerFirst = true;
// 創(chuàng)建BrokerData
brokerData = new BrokerData(clusterName, brokerName, new HashMap<Long, String>());
// 將brokerData注冊到brokerAddrTable中
this.brokerAddrTable.put(brokerName, brokerData);
}
// 拿到brokerName對應(yīng)的所有地址(主從服務(wù)器地址)
Map<Long, String> brokerAddrsMap = brokerData.getBrokerAddrs();
//Switch slave to master: first remove <1, IP:PORT> in namesrv, then add <0, IP:PORT>
//The same IP:PORT must only have one record in brokerAddrTable
Iterator<Entry<Long, String>> it = brokerAddrsMap.entrySet().iterator();
// 遍歷brokerName對應(yīng)的所有地址(主從服務(wù)器地址)
while (it.hasNext()) {
Entry<Long, String> item = it.next();
// 地址相同,但key值和最新的brokerId不同(brokerId為0代表Master,大于0代表Slave),則說明發(fā)生了主從服務(wù)器broker變更,刪除舊的緩存
if (null != brokerAddr && brokerAddr.equals(item.getValue()) && brokerId != item.getKey()) {
it.remove();
}
}
// 將最新的brokerId, brokerAddr放入緩存
String oldAddr = brokerData.getBrokerAddrs().put(brokerId, brokerAddr);
// --------------更新TopicQueueTable緩存-------------
if (null != topicConfigWrapper
&& MixAll.MASTER_ID == brokerId) {// TopicQueueTable中存的是Master的信息,所以brokerId等于0并且null != topicConfigWrapper時,才更新緩存
if (this.isBrokerTopicConfigChanged(brokerAddr, topicConfigWrapper.getDataVersion())
|| registerFirst) {// 如果TopicQueueTable緩存中沒有broker的地址brokerAddr的信息或收到心跳包的時間和舊的時間不相等或broker為第一次注冊。
ConcurrentMap<String, TopicConfig> tcTable =
topicConfigWrapper.getTopicConfigTable();
if (tcTable != null) {
for (Map.Entry<String, TopicConfig> entry : tcTable.entrySet()) {
// 更新TopicQueueTable緩存
this.createAndUpdateQueueData(brokerName, entry.getValue());
}
}
}
}
// --------------更新brokerLiveTable緩存-------------
BrokerLiveInfo prevBrokerLiveInfo = this.brokerLiveTable.put(brokerAddr,
new BrokerLiveInfo(
System.currentTimeMillis(),
topicConfigWrapper.getDataVersion(),
channel,
haServerAddr));
if (null == prevBrokerLiveInfo) {
log.info("new broker registered, {} HAServer: {}", brokerAddr, haServerAddr);
}
/*RouteInfoManager#createAndUpdateQueueData*/
// 更新TopicQueueTable緩存
private void createAndUpdateQueueData(final String brokerName, final TopicConfig topicConfig) {
QueueData queueData = new QueueData();
queueData.setBrokerName(brokerName);
queueData.setWriteQueueNums(topicConfig.getWriteQueueNums());
queueData.setReadQueueNums(topicConfig.getReadQueueNums());
queueData.setPerm(topicConfig.getPerm());
queueData.setTopicSynFlag(topicConfig.getTopicSysFlag());
// ---------topicQueueTable緩存---------
List<QueueData> queueDataList = this.topicQueueTable.get(topicConfig.getTopicName());
if (null == queueDataList) {// queueDataList為空這創(chuàng)建并添加最新的queueData
queueDataList = new LinkedList<QueueData>();
queueDataList.add(queueData);
this.topicQueueTable.put(topicConfig.getTopicName(), queueDataList);
log.info("new topic registered, {} {}", topicConfig.getTopicName(), queueData);
} else {
boolean addNewOne = true;
Iterator<QueueData> it = queueDataList.iterator();
while (it.hasNext()) {
QueueData qd = it.next();
// 判斷是否已經(jīng)存在brokerame對應(yīng)的queueData信息
if (qd.getBrokerName().equals(brokerName)) {
if (qd.equals(queueData)) {// 判斷舊的queueData和新的queueData是否等價
addNewOne = false;
} else {
log.info("topic changed, {} OLD: {} NEW: {}", topicConfig.getTopicName(), qd,
queueData);
// 不等價則刪除舊的
it.remove();
}
}
}
// 如果最新queueData不等價舊的,或者沒有brokerName對應(yīng)的信息,則添加最新queueData
if (addNewOne) {
queueDataList.add(queueData);
}
}
}
路由刪除
在Broker發(fā)生宕機無法發(fā)送心跳包時,NameServer無法收到心跳包,NameServer每個10秒就會掃面brokerLiveTable,檢測每個Broker地址上次收到心跳包的時間lastUpdateTimeStamp是否超過120s,如果超過120s則認為broker失效,NameServer會移除在topicQueueTable,brokerAddrTable,clusterAddrTable,brokerLiveTable,filterServerTable對應(yīng)的失效Broker信息。
執(zhí)行路由刪除的兩個觸發(fā)點:
- NameServer每十秒的掃描
- Broker在正常狀態(tài)下執(zhí)行unregisterBroker指令
由于兩個出發(fā)點執(zhí)行都是同一段路由刪除代碼,下面通過第一個出發(fā)點進行源碼分析。
/*NamesrvController#initialize*/
// 每個10秒掃描一次Broker,移除處于不激活狀態(tài)的Broker(遍歷brokerLiveTable(HashMap<String/* brokerAddr */, BrokerLiveInfo>,BrokerLiveInfo.lastUpdateTimestamp用于記錄Broker上一次發(fā)送心跳包的時間,如果超過了120就刪除對應(yīng)的Broker)
this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
@Override
public void run() {
NamesrvController.this.routeInfoManager.scanNotActiveBroker();
}
}, 5, 10, TimeUnit.SECONDS);
/*RouteInfoManager#scanNotActiveBroker*/
Iterator<Entry<String, BrokerLiveInfo>> it = this.brokerLiveTable.entrySet().iterator();
while (it.hasNext()) {
Entry<String, BrokerLiveInfo> next = it.next();
long last = next.getValue().getLastUpdateTimestamp();
// 上次更新的時間如果超過120s則刪除
if ((last + BROKER_CHANNEL_EXPIRED_TIME) < System.currentTimeMillis()) {
// 關(guān)閉與該Broker的通道
RemotingUtil.closeChannel(next.getValue().getChannel());
// 刪除
it.remove();
log.warn("The broker channel expired, {} {}ms", next.getKey(), BROKER_CHANNEL_EXPIRED_TIME);
// 更新其他緩存中對應(yīng)的信息
this.onChannelDestroy(next.getKey(), next.getValue().getChannel());
}
}
/*RouteInfoManager#onChannelDestroy*/
if (brokerAddrFound != null && brokerAddrFound.length() > 0) {
try {
try {
this.lock.writeLock().lockInterruptibly();
// 從brokerLiveTable中刪除失效地址brokerAddrFound的信息
this.brokerLiveTable.remove(brokerAddrFound);
// 從filterServerTable中刪除失效地址brokerAddrFound的信息
this.filterServerTable.remove(brokerAddrFound);
String brokerNameFound = null;
boolean removeBrokerName = false;
Iterator<Entry<String, BrokerData>> itBrokerAddrTable =
this.brokerAddrTable.entrySet().iterator();
// // 從brokerAddrTable找到失效地址對應(yīng)得Broker名字
while (itBrokerAddrTable.hasNext() && (null == brokerNameFound)) {
BrokerData brokerData = itBrokerAddrTable.next().getValue();
Iterator<Entry<Long, String>> it = brokerData.getBrokerAddrs().entrySet().iterator();
while (it.hasNext()) {
Entry<Long, String> entry = it.next();
Long brokerId = entry.getKey();
String brokerAddr = entry.getValue();
if (brokerAddr.equals(brokerAddrFound)) {
// 失效地址的Broker名字
brokerNameFound = brokerData.getBrokerName();
// 從失效地址的Broker名字對應(yīng)的brokerData的主從地址集中刪除失效地址
it.remove();
log.info("remove brokerAddr[{}, {}] from brokerAddrTable, because channel destroyed",
brokerId, brokerAddr);
break;
}
}
// 如果失效地址的Broker名字對應(yīng)的brokerData的主從地址集中已經(jīng)沒有地址,說明Broker名字對應(yīng)的所有服務(wù)器失效,需要從brokerAddrTable中刪除失效地址的Broker名字對應(yīng)的brokerData
if (brokerData.getBrokerAddrs().isEmpty()) {
removeBrokerName = true;
itBrokerAddrTable.remove();
log.info("remove brokerName[{}] from brokerAddrTable, because channel destroyed",
brokerData.getBrokerName());
}
}
// 失效地址的Broker名字對應(yīng)的地址已經(jīng)全部失效,刪除其在clusterAddrTable中對應(yīng)的信息
if (brokerNameFound != null && removeBrokerName) {
Iterator<Entry<String, Set<String>>> it = this.clusterAddrTable.entrySet().iterator();
while (it.hasNext()) {
Entry<String, Set<String>> entry = it.next();
String clusterName = entry.getKey();
Set<String> brokerNames = entry.getValue();
boolean removed = brokerNames.remove(brokerNameFound);
if (removed) {
log.info("remove brokerName[{}], clusterName[{}] from clusterAddrTable, because channel destroyed",
brokerNameFound, clusterName);
if (brokerNames.isEmpty()) {
log.info("remove the clusterName[{}] from clusterAddrTable, because channel destroyed and no broker in this cluster",
clusterName);
it.remove();
}
break;
}
}
}
// 失效地址的Broker名字對應(yīng)的地址已經(jīng)全部失效,刪除其在topicQueueTable中對應(yīng)的信息
if (removeBrokerName) {
Iterator<Entry<String, List<QueueData>>> itTopicQueueTable =
this.topicQueueTable.entrySet().iterator();
while (itTopicQueueTable.hasNext()) {
Entry<String, List<QueueData>> entry = itTopicQueueTable.next();
String topic = entry.getKey();
List<QueueData> queueDataList = entry.getValue();
Iterator<QueueData> itQueueData = queueDataList.iterator();
while (itQueueData.hasNext()) {
QueueData queueData = itQueueData.next();
if (queueData.getBrokerName().equals(brokerNameFound)) {
itQueueData.remove();
log.info("remove topic[{} {}], from topicQueueTable, because channel destroyed",
topic, queueData);
}
}
if (queueDataList.isEmpty()) {
itTopicQueueTable.remove();
log.info("remove topic[{}] all queue, from topicQueueTable, because channel destroyed",
topic);
}
}
}
} finally {
this.lock.writeLock().unlock();
}
} catch (Exception e) {
log.error("onChannelDestroy Exception", e);
}
}
路由發(fā)現(xiàn)
通過主題先從topicQueueTable中獲取到對應(yīng)的所有Broker信息,然后再從brokerAddrTable中根據(jù)broker名稱獲取對應(yīng)的信息,最后通過從brokerAddrTable獲取的地址信息從filterServerTable獲取對應(yīng)的信息。
/*DefaultRequestProcessor#processRequest*/
switch (request.getCode()) {
case RequestCode.PUT_KV_CONFIG:
return this.putKVConfig(ctx, request);
...
case RequestCode.UNREGISTER_BROKER:
return this.unregisterBroker(ctx, request);
// 走這里
case RequestCode.GET_ROUTEINTO_BY_TOPIC:
return this.getRouteInfoByTopic(ctx, request);
case RequestCode.GET_BROKER_CLUSTER_INFO:
return this.getBrokerClusterInfo(ctx, request);
...
case RequestCode.GET_NAMESRV_CONFIG:
return this.getConfig(ctx, request);
default:
break;
}
/*DefaultRequestProcessor#getRouteInfoByTopic*/
final RemotingCommand response = RemotingCommand.createResponseCommand(null);
final GetRouteInfoRequestHeader requestHeader =
(GetRouteInfoRequestHeader) request.decodeCommandCustomHeader(GetRouteInfoRequestHeader.class);
// 通過主題先后從topicQueueTable,brokerAddrTable,filterServerTable獲取對應(yīng)的信息
TopicRouteData topicRouteData = this.namesrvController.getRouteInfoManager().pickupTopicRouteData(requestHeader.getTopic());
if (topicRouteData != null) {
if (this.namesrvController.getNamesrvConfig().isOrderMessageEnable()) {
String orderTopicConf =
this.namesrvController.getKvConfigManager().getKVConfig(NamesrvUtil.NAMESPACE_ORDER_TOPIC_CONFIG,
requestHeader.getTopic());
topicRouteData.setOrderTopicConf(orderTopicConf);
}
byte[] content = topicRouteData.encode();
response.setBody(content);
response.setCode(ResponseCode.SUCCESS);
response.setRemark(null);
return response;
}
response.setCode(ResponseCode.TOPIC_NOT_EXIST);
response.setRemark("No topic route info in name server for the topic: " + requestHeader.getTopic()
+ FAQUrl.suggestTodo(FAQUrl.APPLY_TOPIC_URL));
return response;
/*RouteInfoManager#pickupTopicRouteData*/
TopicRouteData topicRouteData = new TopicRouteData();
boolean foundQueueData = false;
boolean foundBrokerData = false;
// 用于存放主題對應(yīng)的所有broker名稱
Set<String> brokerNameSet = new HashSet<String>();
// 用于存放從brokerAddrTable獲取的brokerData
List<BrokerData> brokerDataList = new LinkedList<BrokerData>();
topicRouteData.setBrokerDatas(brokerDataList);
// 用于存放從filterServerTable獲取的信息
HashMap<String, List<String>> filterServerMap = new HashMap<String, List<String>>();
topicRouteData.setFilterServerTable(filterServerMap);
try {
try {
this.lock.readLock().lockInterruptibly();
// 獲取主題對應(yīng)的broker服務(wù)器信息,即topicQueueTable中主題對應(yīng)的所有QueueData
List<QueueData> queueDataList = this.topicQueueTable.get(topic);
if (queueDataList != null) {
topicRouteData.setQueueDatas(queueDataList);
foundQueueData = true;
Iterator<QueueData> it = queueDataList.iterator();
while (it.hasNext()) {
QueueData qd = it.next();
brokerNameSet.add(qd.getBrokerName());
}
//根據(jù)主題對應(yīng)的broker名稱,從brokerAddrTable中獲取對應(yīng)brokerData信息,brokerData包含集群名稱,broker名稱,broker主從地址信息
for (String brokerName : brokerNameSet) {
BrokerData brokerData = this.brokerAddrTable.get(brokerName);
if (null != brokerData) {
BrokerData brokerDataClone = new BrokerData(brokerData.getCluster(), brokerData.getBrokerName(), (HashMap<Long, String>) brokerData
.getBrokerAddrs().clone());
brokerDataList.add(brokerDataClone);
foundBrokerData = true;
// 根據(jù)brokerDataClone的地址信息獲取到filterServerTable中的信息
for (final String brokerAddr : brokerDataClone.getBrokerAddrs().values()) {
List<String> filterServerList = this.filterServerTable.get(brokerAddr);
filterServerMap.put(brokerAddr, filterServerList);
}
}
}
}
} finally {
this.lock.readLock().unlock();
}
} catch (Exception e) {
log.error("pickupTopicRouteData Exception", e);
}
log.debug("pickupTopicRouteData {} {}", topic, topicRouteData);
if (foundBrokerData && foundQueueData) {
return topicRouteData;
}
return null;