mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-01 15:47:01 +08:00
[ISSUE #4127] [BrokerOuterAPI] Anonymous new Runnable() can be replaced with lambda
This commit is contained in:
committed by
GitHub
parent
4f2b708179
commit
519c0b8a46
@@ -928,7 +928,7 @@ public class BrokerController {
|
||||
|
||||
if (!PermName.isWriteable(this.getBrokerConfig().getBrokerPermission())
|
||||
|| !PermName.isReadable(this.getBrokerConfig().getBrokerPermission())) {
|
||||
ConcurrentHashMap<String, TopicConfig> topicConfigTable = new ConcurrentHashMap<String, TopicConfig>();
|
||||
ConcurrentHashMap<String, TopicConfig> topicConfigTable = new ConcurrentHashMap<>();
|
||||
for (TopicConfig topicConfig : topicConfigWrapper.getTopicConfigTable().values()) {
|
||||
TopicConfig tmp =
|
||||
new TopicConfig(topicConfig.getTopicName(), topicConfig.getReadQueueNums(), topicConfig.getWriteQueueNums(),
|
||||
|
||||
@@ -63,7 +63,7 @@ public class BrokerOuterAPI {
|
||||
private final TopAddressing topAddressing = new TopAddressing(MixAll.getWSAddr());
|
||||
private String nameSrvAddr = null;
|
||||
private BrokerFixedThreadPoolExecutor brokerOuterExecutor = new BrokerFixedThreadPoolExecutor(4, 10, 1, TimeUnit.MINUTES,
|
||||
new ArrayBlockingQueue<Runnable>(32), new ThreadFactoryImpl("brokerOutApi_thread_", true));
|
||||
new ArrayBlockingQueue<>(32), new ThreadFactoryImpl("brokerOutApi_thread_", true));
|
||||
|
||||
public BrokerOuterAPI(final NettyClientConfig nettyClientConfig) {
|
||||
this(nettyClientConfig, null);
|
||||
@@ -142,21 +142,18 @@ public class BrokerOuterAPI {
|
||||
requestHeader.setBodyCrc32(bodyCrc32);
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(nameServerAddressList.size());
|
||||
for (final String namesrvAddr : nameServerAddressList) {
|
||||
brokerOuterExecutor.execute(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
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();
|
||||
brokerOuterExecutor.execute(() -> {
|
||||
try {
|
||||
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();
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -269,46 +266,43 @@ public class BrokerOuterAPI {
|
||||
if (nameServerAddressList != null && nameServerAddressList.size() > 0) {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(nameServerAddressList.size());
|
||||
for (final String namesrvAddr : nameServerAddressList) {
|
||||
brokerOuterExecutor.execute(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
QueryDataVersionRequestHeader requestHeader = new QueryDataVersionRequestHeader();
|
||||
requestHeader.setBrokerAddr(brokerAddr);
|
||||
requestHeader.setBrokerId(brokerId);
|
||||
requestHeader.setBrokerName(brokerName);
|
||||
requestHeader.setClusterName(clusterName);
|
||||
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.QUERY_DATA_VERSION, requestHeader);
|
||||
request.setBody(topicConfigWrapper.getDataVersion().encode());
|
||||
RemotingCommand response = remotingClient.invokeSync(namesrvAddr, request, timeoutMills);
|
||||
DataVersion nameServerDataVersion = null;
|
||||
Boolean changed = false;
|
||||
switch (response.getCode()) {
|
||||
case ResponseCode.SUCCESS: {
|
||||
QueryDataVersionResponseHeader queryDataVersionResponseHeader =
|
||||
(QueryDataVersionResponseHeader) response.decodeCommandCustomHeader(QueryDataVersionResponseHeader.class);
|
||||
changed = queryDataVersionResponseHeader.getChanged();
|
||||
byte[] body = response.getBody();
|
||||
if (body != null) {
|
||||
nameServerDataVersion = DataVersion.decode(body, DataVersion.class);
|
||||
if (!topicConfigWrapper.getDataVersion().equals(nameServerDataVersion)) {
|
||||
changed = true;
|
||||
}
|
||||
}
|
||||
if (changed == null || changed) {
|
||||
changedList.add(Boolean.TRUE);
|
||||
brokerOuterExecutor.execute(() -> {
|
||||
try {
|
||||
QueryDataVersionRequestHeader requestHeader = new QueryDataVersionRequestHeader();
|
||||
requestHeader.setBrokerAddr(brokerAddr);
|
||||
requestHeader.setBrokerId(brokerId);
|
||||
requestHeader.setBrokerName(brokerName);
|
||||
requestHeader.setClusterName(clusterName);
|
||||
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.QUERY_DATA_VERSION, requestHeader);
|
||||
request.setBody(topicConfigWrapper.getDataVersion().encode());
|
||||
RemotingCommand response = remotingClient.invokeSync(namesrvAddr, request, timeoutMills);
|
||||
DataVersion nameServerDataVersion = null;
|
||||
Boolean changed = false;
|
||||
switch (response.getCode()) {
|
||||
case ResponseCode.SUCCESS: {
|
||||
QueryDataVersionResponseHeader queryDataVersionResponseHeader =
|
||||
(QueryDataVersionResponseHeader) response.decodeCommandCustomHeader(QueryDataVersionResponseHeader.class);
|
||||
changed = queryDataVersionResponseHeader.getChanged();
|
||||
byte[] body = response.getBody();
|
||||
if (body != null) {
|
||||
nameServerDataVersion = DataVersion.decode(body, DataVersion.class);
|
||||
if (!topicConfigWrapper.getDataVersion().equals(nameServerDataVersion)) {
|
||||
changed = true;
|
||||
}
|
||||
}
|
||||
default:
|
||||
break;
|
||||
if (changed == null || changed) {
|
||||
changedList.add(Boolean.TRUE);
|
||||
}
|
||||
}
|
||||
log.warn("Query data version from name server {} OK,changed {}, broker {},name server {}", namesrvAddr, changed, topicConfigWrapper.getDataVersion(), nameServerDataVersion == null ? "" : nameServerDataVersion);
|
||||
} catch (Exception e) {
|
||||
changedList.add(Boolean.TRUE);
|
||||
log.error("Query data version from name server {} Exception, {}", namesrvAddr, e);
|
||||
} finally {
|
||||
countDownLatch.countDown();
|
||||
default:
|
||||
break;
|
||||
}
|
||||
log.warn("Query data version from name server {} OK,changed {}, broker {},name server {}", namesrvAddr, changed, topicConfigWrapper.getDataVersion(), nameServerDataVersion == null ? "" : nameServerDataVersion);
|
||||
} catch (Exception e) {
|
||||
changedList.add(Boolean.TRUE);
|
||||
log.error("Query data version from name server {} Exception, {}", namesrvAddr, e);
|
||||
} finally {
|
||||
countDownLatch.countDown();
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user