Rocketmq之NameServer地址配置及更新過程
前言
NameServer在整個Rocketmq的模塊劃分中占據(jù)重要的地位,起到類似于注冊中心的作用。
BrokerServer啟動時需要向NameServer注冊自身元數(shù)據(jù)信息以及主題Topic信息,而Producer發(fā)送消息到BrokerServer、Consumer從BrokerServer訂閱消息,則需要經(jīng)過NameServer才能確定最終要進行數(shù)據(jù)通訊的BrokerServer的地址,所以,BrokerServer、Producer、Consumer程序啟動均需要配置NameServer的地址。
本篇文章就來聊一聊,BrokerServer、Producer、Consumer程序如何設置NameServer的地址以及NameServer的地址是否可以動態(tài)更新。
配置方式
首要我們需要清楚一點,代理服務器BrokerServer的配置信息最后會封裝到BrokerConfig對象中,而生產(chǎn)者Producer、消費者Consumer的配置信息最終會封裝到ClientConfig對象中,那么在這兩個類中應該有NameServer地址相關(guān)的變量。
ClientConfig代碼片段
public class ClientConfig {
// ...省略部分代碼
private String namesrvAddr = NameServerAddressUtils.getNameServerAddresses();
// ...省略部分代碼
}
NameServerAddressUtils#getNameServerAddresses
public class NameServerAddressUtils {
// ...省略部分代碼
public static String getNameServerAddresses() {
return System.getProperty(MixAll.NAMESRV_ADDR_PROPERTY, System.getenv(MixAll.NAMESRV_ADDR_ENV));
}
// ...省略部分代碼
}
BrokerConfig代碼片段
public class BrokerConfig {
// ...省略部分代碼
@ImportantField
private String namesrvAddr = System.getProperty(MixAll.NAMESRV_ADDR_PROPERTY, System.getenv(MixAll.NAMESRV_ADDR_ENV));
// ...省略部分代碼
}
可見,不管是ClientConfig,還是BrokerConfig,NameServer地址變量的初始值,默認先取系統(tǒng)屬性rocketmq.namesrv.addr的值,如果該系統(tǒng)屬性未設置,則再取環(huán)境變量NAMESRV_ADDR的值,如果兩者均未設置,那么配置類中NameServer地址變量初始值為null。
注:
通過上述分析,我們可以知曉兩種配置NameServer地址的方法,不管你的程序是BrokerServer、Producer或者是Consumer
1.設置系統(tǒng)屬性:rocketmq.namesrv.addr
2.設置環(huán)境變量:NAMESRV_ADDR
如果配置類中NameServer地址變量初始值為null,那BrokerServer、Producer、Consumer啟動的時候是不是就會因為沒有這個值而啟動不了或者報錯呢?
其實并不會,還有額外的補償手段,程序會通過http請求去訪問一個特定的url獲取NameServer地址,我們繼續(xù)分析。
Producer或者Consumer
Producer或者Consumer底層都會持有一個MQClientInstance類對象,而在MQClientInstance類中,我們可以看到通過url請求NameServer地址的代碼。
原生API發(fā)送消息或者消費消息的代碼大致如下所列,我們以消息發(fā)送者的start方法為切入口進行分析。
// 原生API發(fā)送消息
DefaultMQProducer producer = new DefaultMQProducer("test-group");
producer.start();
SendResult sendResult = producer.send(new Message("test-topic", "test-message".getBytes(StandardCharsets.UTF_8)));
System.out.println(sendResult);
// 原生API消費消息
DefaultMQPushConsumer defaultMQPushConsumer = new DefaultMQPushConsumer("test-group");
defaultMQPushConsumer.registerMessageListener((MessageListenerConcurrently) (msgList, context) -> {
try {
msgList.forEach(System.out::println);
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
defaultMQPushConsumer.subscribe("test-topic", "*");
defaultMQPushConsumer.start();
DefaultMQProducer#start
@Override
public void start() throws MQClientException {
this.setProducerGroup(withNamespace(this.producerGroup));
this.defaultMQProducerImpl.start();
if (null != traceDispatcher) {
try {
traceDispatcher.start(this.getNamesrvAddr(), this.getAccessChannel());
} catch (MQClientException e) {
log.warn("trace dispatcher start failed ", e);
}
}
}
我們可以看到,DefaultMQProducer的start方法主要邏輯委托給了DefaultMQProducerImpl的start方法
DefaultMQProducerImpl#start
public void start() throws MQClientException {
this.start(true);
}
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();
}
this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQProducer, rpcHook);
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());
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;
}
this.mQClientFactory.sendHeartbeatToAllBrokerWithLock();
this.startScheduledTask();
}
在DefaultMQProducerImpl的start(final boolean startFactory)方法中,創(chuàng)建了MQClientInstance對象實例,并調(diào)用了其start方法
MQClientInstance#start
public void start() throws MQClientException {
synchronized (this) {
switch (this.serviceState) {
case CREATE_JUST:
// 省略部分代碼
// 沒有手動指定NameSrv的值,從遠端服務器獲取并更新本地緩存
if (null == this.clientConfig.getNamesrvAddr()) {
this.mQClientAPIImpl.fetchNameServerAddr();
}
// 省略部分代碼
// Start various schedule tasks
this.startScheduledTask();
// 省略部分代碼
break;
case START_FAILED:
throw new MQClientException("The Factory object[" + this.getClientId() + "] has been created before, and failed.", null);
default:
break;
}
}
}
MQClientInstance的start方法中,與我們該篇分析的內(nèi)容相關(guān)的大致就是上述代碼所列的兩處。
第一,判斷ClientConfig對象的namesrvAddr值是否為null,若是,調(diào)用MQClientAPIImpl的fetchNameServerAddr方法從遠端服務器拉取NameServer服務器的地址,并更新用到的地方
MQClientAPIImpl#fetchNameServerAddr
public String fetchNameServerAddr() {
try {
String addrs = this.topAddressing.fetchNSAddr();
if (addrs != null) {
if (!addrs.equals(this.nameSrvAddr)) {
log.info("name server address changed, old=" + this.nameSrvAddr + ", new=" + addrs);
this.updateNameServerAddressList(addrs);
this.nameSrvAddr = addrs;
return nameSrvAddr;
}
}
} catch (Exception e) {
log.error("fetchNameServerAddr Exception", e);
}
return nameSrvAddr;
}
TopAddressing#fetchNSAddr
public final String fetchNSAddr() {
return fetchNSAddr(true, 3000);
}
public final String fetchNSAddr(boolean verbose, long timeoutMills) {
String url = this.wsAddr;
try {
if (!UtilAll.isBlank(this.unitName)) {
url = url + "-" + this.unitName + "?nofix=1";
}
HttpTinyClient.HttpResult result = HttpTinyClient.httpGet(url, null, null, "UTF-8", timeoutMills);
if (200 == result.code) {
String responseStr = result.content;
if (responseStr != null) {
return clearNewLine(responseStr);
} else {
log.error("fetch nameserver address is null");
}
} else {
log.error("fetch nameserver address failed. statusCode=" + result.code);
}
} catch (IOException e) {
if (verbose) {
log.error("fetch name server address exception", e);
}
}
if (verbose) {
String errorMsg =
"connect to " + url + " failed, maybe the domain name " + MixAll.getWSAddr() + " not bind in /etc/hosts";
errorMsg += FAQUrl.suggestTodo(FAQUrl.NAME_SERVER_ADDR_NOT_EXIST_URL);
log.warn(errorMsg);
}
return null;
}
向遠端請求的url地址就是wsAddr變量值,那么這個變量是在什么地方賦值的呢?發(fā)現(xiàn)該變量是在TopAddressing的構(gòu)造方法中賦值,而TopAddressing對象又是在MQClientAPIImpl中創(chuàng)建,我們找到具體的創(chuàng)建邏輯
MQClientAPIImpl#Constructor
public MQClientAPIImpl(final NettyClientConfig nettyClientConfig,
final ClientRemotingProcessor clientRemotingProcessor,
RPCHook rpcHook, final ClientConfig clientConfig) {
this.clientConfig = clientConfig;
topAddressing = new TopAddressing(MixAll.getWSAddr(), clientConfig.getUnitName());
// 省略部分邏輯
}
傳入TopAddressing構(gòu)造方法的值取自MixAll的getWSAddr方法的返回值
MixAll#getWSAddr
public static String getWSAddr() {
String wsDomainName = System.getProperty("rocketmq.namesrv.domain", DEFAULT_NAMESRV_ADDR_LOOKUP);
String wsDomainSubgroup = System.getProperty("rocketmq.namesrv.domain.subgroup", "nsaddr");
String wsAddr = "http://" + wsDomainName + ":8080/rocketmq/" + wsDomainSubgroup;
if (wsDomainName.indexOf(":") > 0) {
wsAddr = "http://" + wsDomainName + "/rocketmq/" + wsDomainSubgroup;
}
return wsAddr;
}
通過分析上述方法,可以知道,默認會從http://jmenv.tbsite.net:8080/rocketmq/nsaddr這個鏈接處拉取NameServer的地址。當然,如果你想更改這個鏈接,可以修改系統(tǒng)屬性"rocketmq.namesrv.domain"和"rocketmq.namesrv.domain.subgroup"的值以達到目的。
第二,判斷ClientConfig對象的namesrvAddr值是否為null,若是,開啟定時任務,周期性地更新NameServer服務器的地址
MQClientInstance#startScheduledTask
private void startScheduledTask() {
// 沒有手動指定name-server地址的情況下,兩分鐘更新一次name-server地址
// 這個就是name-server可以動態(tài)變化的唯一途徑
if (null == this.clientConfig.getNamesrvAddr()) {
// 兩分鐘拉取更新一次name-server地址
this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
@Override
public void run() {
try {
MQClientInstance.this.mQClientAPIImpl.fetchNameServerAddr();
} catch (Exception e) {
log.error("ScheduledTask fetchNameServerAddr exception", e);
}
}
}, 1000 * 10, 1000 * 60 * 2, TimeUnit.MILLISECONDS);
}
}
可以看出,默認每隔2分鐘會調(diào)用MQClientAPIImpl的fetchNameServerAddr方法更新NameServer地址,具體邏輯上述第一條已經(jīng)解釋過不再贅述
BrokerServer
BrokerServer的啟動類是BrokerStartup,在BrokerStartup的main方法中,會創(chuàng)建BrokerController的實例對象并調(diào)用其start方法,而在創(chuàng)建BrokerController的實例對象時會一并調(diào)用其initialize方法進行初始化,就在這個初始化方法中,有跟我們該篇分析相關(guān)的內(nèi)容。
BrokerStartup#main
public static void main(String[] args) {
start(createBrokerController(args));
}
BrokerStartup#createBrokerController
public static BrokerController createBrokerController(String[] args) {
// 省略部分代碼
try {
// 省略部分代碼
boolean initResult = controller.initialize();
// 省略部分代碼
return controller;
} catch (Throwable e) {
e.printStackTrace();
System.exit(-1);
}
return null;
}
BrokerController#initialize
public boolean initialize() throws CloneNotSupportedException {
// 將本地文件中存儲的數(shù)據(jù)加載至內(nèi)存
boolean result = this.topicConfigManager.load();
result = result && this.consumerOffsetManager.load();
result = result && this.subscriptionGroupManager.load();
result = result && this.consumerFilterManager.load();
// 省略部分代碼
result = result && this.messageStore.load();
if (result) {
// 省略部分代碼
if (this.brokerConfig.getNamesrvAddr() != null) {
this.brokerOuterAPI.updateNameServerAddressList(this.brokerConfig.getNamesrvAddr());
log.info("Set user specified name server address: {}", this.brokerConfig.getNamesrvAddr());
} else if (this.brokerConfig.isFetchNamesrvAddrByAddressServer()) {
// 沒有明確指定name-server的地址,且配置了允許從地址服務器獲取name-server地址
// 每隔2分鐘從name-server地址服務器拉取最新的配置
this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
@Override
public void run() {
try {
BrokerController.this.brokerOuterAPI.fetchNameServerAddr();
} catch (Throwable e) {
log.error("ScheduledTask fetchNameServerAddr exception", e);
}
}
}, 1000 * 10, 1000 * 60 * 2, TimeUnit.MILLISECONDS);
}
// 省略部分代碼
}
return result;
}
可以看到,BrokerConfig中如果namesrvAddr變量值為null,默認每隔2分鐘會調(diào)用BrokerOuterAPI的fetchNameServerAddr方法更新NameServer地址
BrokerOuterAPI#fetchNameServerAddr
public String fetchNameServerAddr() {
try {
String addrs = this.topAddressing.fetchNSAddr();
if (addrs != null) {
if (!addrs.equals(this.nameSrvAddr)) {
log.info("name server address changed, old: {} new: {}", this.nameSrvAddr, addrs);
this.updateNameServerAddressList(addrs);
this.nameSrvAddr = addrs;
return nameSrvAddr;
}
}
} catch (Exception e) {
log.error("fetchNameServerAddr Exception", e);
}
return nameSrvAddr;
}
該方法與MQClientAPIImpl的fetchNameServerAddr方法邏輯幾乎一模一樣,也不再進行贅述。
注:
第三種配置NameServer地址的方法,那就是系統(tǒng)啟動的時候,不配置系統(tǒng)屬性rocketmq.namesrv.addr和環(huán)境變量NAMESRV_ADDR,這樣系統(tǒng)會自動訪問一個url從遠端服務器拉取NameServer的地址,同時會開啟相應的定時任務定時刷新NameServer地址。
總結(jié)
BrokerServer、Producer、Consumer程序啟動獲取NameServer地址的三種方式
- 1.配置系統(tǒng)屬性:rocketmq.namesrv.addr
- 2.配置環(huán)境變量:NAMESRV_ADDR
- 3.不配置系統(tǒng)屬性rocketmq.namesrv.addr和環(huán)境變量NAMESRV_ADDR,讓程序自動訪問一個url從遠端服務器拉取NameServer的地址,同時會開啟相應的定時任務定時刷新NameServer地址。該url地址默認為http://jmenv.tbsite.net:8080/rocketmq/nsaddr,可以通過修改系統(tǒng)屬性rocketmq.namesrv.domain和rocketmq.namesrv.domain.subgroup以達到修改url地址的目的。
只有上述第3種方式才可以實現(xiàn)程序運行過程中動態(tài)更新NameServer地址值。
以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。
- RocketMQ中的NameServer詳細解析
- SpringBoot定時監(jiān)聽RocketMQ的NameServer問題及解決方案
- Linux解決RocketMQ中NameServer啟動問題的方法詳解
- RocketMQ?NameServer架構(gòu)設計啟動流程
- RocketMQ NameServer保障數(shù)據(jù)一致性實現(xiàn)方法講解
- RocketMQ生產(chǎn)者一個應用不能發(fā)送多個NameServer消息解決
- RocketMQ NameServer 核心源碼解析
- RocketMQ之NameServer架構(gòu)設計及啟動關(guān)閉流程源碼分析
- java開發(fā)RocketMQ之NameServer路由管理源碼分析
相關(guān)文章
Spring Cloud升級最新Finchley版本的所有坑
springboot之security?FilterSecurityInterceptor的使用要點記錄

