From 519c0b8a46eac983ccb6848b0f0b5fbd2cf7db15 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9D=8E=E6=99=93=E5=8F=8C=20Li=20Xiao=20Shuang?= <644968328@qq.com> Date: Thu, 21 Apr 2022 10:48:07 +0800 Subject: [PATCH] [ISSUE #4127] [BrokerOuterAPI] Anonymous new Runnable() can be replaced with lambda --- .../rocketmq/broker/BrokerController.java | 2 +- .../rocketmq/broker/out/BrokerOuterAPI.java | 96 +++++++++---------- 2 files changed, 46 insertions(+), 52 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java index 772eec6dc4..a29552056a 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java @@ -928,7 +928,7 @@ public class BrokerController { if (!PermName.isWriteable(this.getBrokerConfig().getBrokerPermission()) || !PermName.isReadable(this.getBrokerConfig().getBrokerPermission())) { - ConcurrentHashMap topicConfigTable = new ConcurrentHashMap(); + ConcurrentHashMap topicConfigTable = new ConcurrentHashMap<>(); for (TopicConfig topicConfig : topicConfigWrapper.getTopicConfigTable().values()) { TopicConfig tmp = new TopicConfig(topicConfig.getTopicName(), topicConfig.getReadQueueNums(), topicConfig.getWriteQueueNums(), diff --git a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java index de7f3fce81..c5b53a7772 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java @@ -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(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(); } });