From 95202f1e6ffcd384b6aaefa5b4271b122d2eec21 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=AD=E5=B0=8F=E6=BC=AA?= <644120242@qq.com> Date: Sun, 13 Mar 2022 17:13:42 +0800 Subject: [PATCH 1/9] docs: Add the link to Apache RocketMQ MQTT in the readme.md (#3971) --- README.md | 1 + 1 file changed, 1 insertion(+) diff --git a/README.md b/README.md index 773c4e85ae..c780bdaf55 100644 --- a/README.md +++ b/README.md @@ -47,6 +47,7 @@ It offers a variety of features: * [RocketMQ Docker](https://github.com/apache/rocketmq-docker) * [RocketMQ Dashboard](https://github.com/apache/rocketmq-dashboard) * [RocketMQ Connect](https://github.com/apache/rocketmq-connect) +* [RocketMQ MQTT](https://github.com/apache/rocketmq-mqtt) * [RocketMQ Incubating Community Projects](https://github.com/apache/rocketmq-externals) ---------- From 895bd95e82e597ebb6c970e86522e238618a8610 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=9C=A8=E7=BA=A2=E5=B0=98=E4=B8=AD=E6=88=90=E4=BB=99?= <1178404986@qq.com> Date: Mon, 14 Mar 2022 01:21:16 -0500 Subject: [PATCH 2/9] fix a flaky test (#3959) --- .../java/org/apache/rocketmq/common/filter/FilterAPITest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/src/test/java/org/apache/rocketmq/common/filter/FilterAPITest.java b/common/src/test/java/org/apache/rocketmq/common/filter/FilterAPITest.java index 5190f88f47..bf14207c30 100644 --- a/common/src/test/java/org/apache/rocketmq/common/filter/FilterAPITest.java +++ b/common/src/test/java/org/apache/rocketmq/common/filter/FilterAPITest.java @@ -57,7 +57,7 @@ public class FilterAPITest { assertThat(ExpressionType.isTagType(subscriptionData.getExpressionType())).isTrue(); assertThat(subscriptionData.getTagsSet()).isNotNull(); - assertThat(subscriptionData.getTagsSet()).containsExactly("A", "B"); + assertThat(subscriptionData.getTagsSet()).containsExactlyInAnyOrder("A", "B"); } catch (Exception e) { e.printStackTrace(); assertThat(Boolean.FALSE).isTrue(); From 8cd998829ab2c6f94ca6ad6d285f56f455e97790 Mon Sep 17 00:00:00 2001 From: rongtong Date: Mon, 14 Mar 2022 19:29:15 +0800 Subject: [PATCH 3/9] Fix using wrong offset when deliver in ScheduleService (#3967) --- .../store/schedule/ScheduleMessageService.java | 18 ++++++++---------- 1 file changed, 8 insertions(+), 10 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/schedule/ScheduleMessageService.java b/store/src/main/java/org/apache/rocketmq/store/schedule/ScheduleMessageService.java index d5b4e8d390..f14b198abb 100644 --- a/store/src/main/java/org/apache/rocketmq/store/schedule/ScheduleMessageService.java +++ b/store/src/main/java/org/apache/rocketmq/store/schedule/ScheduleMessageService.java @@ -447,9 +447,9 @@ public class ScheduleMessageService extends ConfigManager { boolean deliverSuc; if (ScheduleMessageService.this.enableAsyncDeliver) { - deliverSuc = this.asyncDeliver(msgInner, msgExt.getMsgId(), offset, offsetPy, sizePy); + deliverSuc = this.asyncDeliver(msgInner, msgExt.getMsgId(), nextOffset, offsetPy, sizePy); } else { - deliverSuc = this.syncDeliver(msgInner, msgExt.getMsgId(), offset, offsetPy, sizePy); + deliverSuc = this.syncDeliver(msgInner, msgExt.getMsgId(), nextOffset, offsetPy, sizePy); } if (!deliverSuc) { @@ -787,24 +787,22 @@ public class ScheduleMessageService extends ConfigManager { public enum ProcessStatus { /** * In process, the processing result has not yet been returned. - * */ + */ RUNNING, /** * Put message success. - * */ + */ SUCCESS, /** - * Put message exception. - * When autoResend is true, the message will be resend. - * */ + * Put message exception. When autoResend is true, the message will be resend. + */ EXCEPTION, /** - * Skip put message. - * When the message cannot be looked, the message will be skipped. - * */ + * Skip put message. When the message cannot be looked, the message will be skipped. + */ SKIP, } } From 3378da69edd0d3e5c7ab4409318ae46d03493dcc Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Mon, 14 Mar 2022 20:15:43 +0800 Subject: [PATCH 4/9] [ISSUE #3882]Nameserver change modify `topicQueueTable` in `RouteInfoManager` (#3881) * 1. nameserver change. modify `topicQueueTable` in `RouteInfoManager` add brokerName to QueueData mapping to speed up brokerName related logic * fix checkstyle * fix unit test * refactor route info manager unit test * 1. when `pickupTopicRouteData` only add filter server address if filter server table not empty. 2. merge with patch #3893 * 1. `RouteInfoManger` only add filter server list when filterServerTable contains brokerAddr * 1. add master change info log * 1. add master change info log * 1. add delete topic test 2. add slave change to master test --- .../processor/DefaultRequestProcessor.java | 25 +- .../namesrv/routeinfo/RouteInfoManager.java | 354 ++++++++---------- .../RouteInfoManagerBrokerPermTest.java | 111 ++++++ .../RouteInfoManagerBrokerRegisterTest.java | 124 ++++++ .../RouteInfoManagerStaticRegisterTest.java | 153 ++++++++ .../routeinfo/RouteInfoManagerTest.java | 162 -------- .../routeinfo/RouteInfoManagerTestBase.java | 188 ++++++++++ 7 files changed, 746 insertions(+), 371 deletions(-) create mode 100644 namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerPermTest.java create mode 100644 namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerRegisterTest.java create mode 100644 namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerStaticRegisterTest.java delete mode 100644 namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTest.java create mode 100644 namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTestBase.java diff --git a/namesrv/src/main/java/org/apache/rocketmq/namesrv/processor/DefaultRequestProcessor.java b/namesrv/src/main/java/org/apache/rocketmq/namesrv/processor/DefaultRequestProcessor.java index bde03489d3..8068c72d8e 100644 --- a/namesrv/src/main/java/org/apache/rocketmq/namesrv/processor/DefaultRequestProcessor.java +++ b/namesrv/src/main/java/org/apache/rocketmq/namesrv/processor/DefaultRequestProcessor.java @@ -27,6 +27,8 @@ import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.UtilAll; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.help.FAQUrl; +import org.apache.rocketmq.common.protocol.body.ClusterInfo; +import org.apache.rocketmq.common.protocol.body.TopicList; import org.apache.rocketmq.common.protocol.header.namesrv.AddWritePermOfBrokerRequestHeader; import org.apache.rocketmq.common.protocol.header.namesrv.AddWritePermOfBrokerResponseHeader; import org.apache.rocketmq.logging.InternalLogger; @@ -270,7 +272,7 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen Boolean changed = this.namesrvController.getRouteInfoManager().isBrokerTopicConfigChanged(requestHeader.getBrokerAddr(), dataVersion); if (!changed) { - this.namesrvController.getRouteInfoManager().updateBrokerInfoUpdateTimestamp(requestHeader.getBrokerAddr()); + this.namesrvController.getRouteInfoManager().updateBrokerInfoUpdateTimestamp(requestHeader.getBrokerAddr(), System.currentTimeMillis()); } DataVersion nameSeverDataVersion = this.namesrvController.getRouteInfoManager().queryBrokerTopicConfig(requestHeader.getBrokerAddr()); @@ -376,7 +378,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen private RemotingCommand getBrokerClusterInfo(ChannelHandlerContext ctx, RemotingCommand request) { final RemotingCommand response = RemotingCommand.createResponseCommand(null); - byte[] content = this.namesrvController.getRouteInfoManager().getAllClusterInfo(); + ClusterInfo clusterInfo = this.namesrvController.getRouteInfoManager().getAllClusterInfo(); + byte[] content = clusterInfo.encode(); response.setBody(content); response.setCode(ResponseCode.SUCCESS); @@ -427,7 +430,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen private RemotingCommand getAllTopicListFromNameserver(ChannelHandlerContext ctx, RemotingCommand request) { final RemotingCommand response = RemotingCommand.createResponseCommand(null); - byte[] body = this.namesrvController.getRouteInfoManager().getAllTopicList(); + TopicList allTopicList = this.namesrvController.getRouteInfoManager().getAllTopicList(); + byte[] body = allTopicList.encode(); response.setBody(body); response.setCode(ResponseCode.SUCCESS); @@ -474,7 +478,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen final GetTopicsByClusterRequestHeader requestHeader = (GetTopicsByClusterRequestHeader) request.decodeCommandCustomHeader(GetTopicsByClusterRequestHeader.class); - byte[] body = this.namesrvController.getRouteInfoManager().getTopicsByCluster(requestHeader.getCluster()); + TopicList topicsByCluster = this.namesrvController.getRouteInfoManager().getTopicsByCluster(requestHeader.getCluster()); + byte[] body = topicsByCluster.encode(); response.setBody(body); response.setCode(ResponseCode.SUCCESS); @@ -486,7 +491,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen RemotingCommand request) throws RemotingCommandException { final RemotingCommand response = RemotingCommand.createResponseCommand(null); - byte[] body = this.namesrvController.getRouteInfoManager().getSystemTopicList(); + TopicList systemTopicList = this.namesrvController.getRouteInfoManager().getSystemTopicList(); + byte[] body = systemTopicList.encode(); response.setBody(body); response.setCode(ResponseCode.SUCCESS); @@ -498,7 +504,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen RemotingCommand request) throws RemotingCommandException { final RemotingCommand response = RemotingCommand.createResponseCommand(null); - byte[] body = this.namesrvController.getRouteInfoManager().getUnitTopics(); + TopicList unitTopics = this.namesrvController.getRouteInfoManager().getUnitTopics(); + byte[] body = unitTopics.encode(); response.setBody(body); response.setCode(ResponseCode.SUCCESS); @@ -510,7 +517,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen RemotingCommand request) throws RemotingCommandException { final RemotingCommand response = RemotingCommand.createResponseCommand(null); - byte[] body = this.namesrvController.getRouteInfoManager().getHasUnitSubTopicList(); + TopicList hasUnitSubTopicList = this.namesrvController.getRouteInfoManager().getHasUnitSubTopicList(); + byte[] body = hasUnitSubTopicList.encode(); response.setBody(body); response.setCode(ResponseCode.SUCCESS); @@ -522,7 +530,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen throws RemotingCommandException { final RemotingCommand response = RemotingCommand.createResponseCommand(null); - byte[] body = this.namesrvController.getRouteInfoManager().getHasUnitSubUnUnitTopicList(); + TopicList hasUnitSubUnUnitTopicList = this.namesrvController.getRouteInfoManager().getHasUnitSubUnUnitTopicList(); + byte[] body = hasUnitSubUnUnitTopicList.encode(); response.setBody(body); response.setCode(ResponseCode.SUCCESS); diff --git a/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java b/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java index 982d543946..2069f9674c 100644 --- a/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java +++ b/namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManager.java @@ -17,6 +17,8 @@ package org.apache.rocketmq.namesrv.routeinfo; import io.netty.channel.Channel; + +import java.util.ArrayList; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; @@ -28,6 +30,8 @@ import java.util.Set; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.function.Predicate; + import org.apache.rocketmq.common.DataVersion; import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.TopicConfig; @@ -50,25 +54,25 @@ public class RouteInfoManager { private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.NAMESRV_LOGGER_NAME); private final static long BROKER_CHANNEL_EXPIRED_TIME = 1000 * 60 * 2; private final ReadWriteLock lock = new ReentrantReadWriteLock(); - private final HashMap> topicQueueTable; + private final HashMap> topicQueueTable; private final HashMap brokerAddrTable; private final HashMap> clusterAddrTable; private final HashMap brokerLiveTable; private final HashMap/* Filter Server */> filterServerTable; public RouteInfoManager() { - this.topicQueueTable = new HashMap>(1024); + this.topicQueueTable = new HashMap>(1024); this.brokerAddrTable = new HashMap(128); this.clusterAddrTable = new HashMap>(32); this.brokerLiveTable = new HashMap(256); this.filterServerTable = new HashMap>(256); } - public byte[] getAllClusterInfo() { + public ClusterInfo getAllClusterInfo() { ClusterInfo clusterInfoSerializeWrapper = new ClusterInfo(); clusterInfoSerializeWrapper.setBrokerAddrTable(this.brokerAddrTable); clusterInfoSerializeWrapper.setClusterAddrTable(this.clusterAddrTable); - return clusterInfoSerializeWrapper.encode(); + return clusterInfoSerializeWrapper; } public void deleteTopic(final String topic) { @@ -84,7 +88,7 @@ public class RouteInfoManager { } } - public byte[] getAllTopicList() { + public TopicList getAllTopicList() { TopicList topicList = new TopicList(); try { try { @@ -97,18 +101,18 @@ public class RouteInfoManager { log.error("getAllTopicList Exception", e); } - return topicList.encode(); + return topicList; } public RegisterBrokerResult registerBroker( - final String clusterName, - final String brokerAddr, - final String brokerName, - final long brokerId, - final String haServerAddr, - final TopicConfigSerializeWrapper topicConfigWrapper, - final List filterServerList, - final Channel channel) { + final String clusterName, + final String brokerAddr, + final String brokerName, + final long brokerId, + final String haServerAddr, + final TopicConfigSerializeWrapper topicConfigWrapper, + final List filterServerList, + final Channel channel) { RegisterBrokerResult result = new RegisterBrokerResult(); try { try { @@ -136,19 +140,25 @@ public class RouteInfoManager { while (it.hasNext()) { Entry item = it.next(); if (null != brokerAddr && brokerAddr.equals(item.getValue()) && brokerId != item.getKey()) { + log.debug("remove entry {} from brokerData", item); it.remove(); } } String oldAddr = brokerData.getBrokerAddrs().put(brokerId, brokerAddr); + if (MixAll.MASTER_ID == brokerId) { + log.info("cluster [{}] brokerName [{}] master address change from {} to {}", + brokerData.getCluster(), brokerData.getBrokerName(), oldAddr, brokerAddr); + } + registerFirst = registerFirst || (null == oldAddr); if (null != topicConfigWrapper - && MixAll.MASTER_ID == brokerId) { + && MixAll.MASTER_ID == brokerId) { if (this.isBrokerTopicConfigChanged(brokerAddr, topicConfigWrapper.getDataVersion()) - || registerFirst) { + || registerFirst) { ConcurrentMap tcTable = - topicConfigWrapper.getTopicConfigTable(); + topicConfigWrapper.getTopicConfigTable(); if (tcTable != null) { for (Map.Entry entry : tcTable.entrySet()) { this.createAndUpdateQueueData(brokerName, entry.getValue()); @@ -158,11 +168,11 @@ public class RouteInfoManager { } BrokerLiveInfo prevBrokerLiveInfo = this.brokerLiveTable.put(brokerAddr, - new BrokerLiveInfo( - System.currentTimeMillis(), - topicConfigWrapper.getDataVersion(), - channel, - haServerAddr)); + new BrokerLiveInfo( + System.currentTimeMillis(), + topicConfigWrapper.getDataVersion(), + channel, + haServerAddr)); if (null == prevBrokerLiveInfo) { log.info("new broker registered, {} HAServer: {}", brokerAddr, haServerAddr); } @@ -208,10 +218,10 @@ public class RouteInfoManager { return null; } - public void updateBrokerInfoUpdateTimestamp(final String brokerAddr) { + public void updateBrokerInfoUpdateTimestamp(final String brokerAddr, long timeStamp) { BrokerLiveInfo prev = this.brokerLiveTable.get(brokerAddr); if (prev != null) { - prev.setLastUpdateTimestamp(System.currentTimeMillis()); + prev.setLastUpdateTimestamp(timeStamp); } } @@ -223,31 +233,17 @@ public class RouteInfoManager { queueData.setPerm(topicConfig.getPerm()); queueData.setTopicSysFlag(topicConfig.getTopicSysFlag()); - List queueDataList = this.topicQueueTable.get(topicConfig.getTopicName()); - if (null == queueDataList) { - queueDataList = new LinkedList(); - queueDataList.add(queueData); - this.topicQueueTable.put(topicConfig.getTopicName(), queueDataList); + Map queueDataMap = this.topicQueueTable.get(topicConfig.getTopicName()); + if (null == queueDataMap) { + queueDataMap = new HashMap<>(); + queueDataMap.put(queueData.getBrokerName(), queueData); + this.topicQueueTable.put(topicConfig.getTopicName(), queueDataMap); log.info("new topic registered, {} {}", topicConfig.getTopicName(), queueData); } else { - boolean addNewOne = true; - - Iterator it = queueDataList.iterator(); - while (it.hasNext()) { - QueueData qd = it.next(); - if (qd.getBrokerName().equals(brokerName)) { - if (qd.equals(queueData)) { - addNewOne = false; - } else { - log.info("topic changed, {} OLD: {} NEW: {}", topicConfig.getTopicName(), qd, - queueData); - it.remove(); - } - } - } - - if (addNewOne) { - queueDataList.add(queueData); + QueueData old = queueDataMap.put(queueData.getBrokerName(), queueData); + if (old != null && !old.equals(queueData)) { + log.info("topic changed, {} OLD: {} NEW: {}", topicConfig.getTopicName(), old, + queueData); } } } @@ -278,11 +274,14 @@ public class RouteInfoManager { private int operateWritePermOfBroker(final String brokerName, final int requestCode) { int topicCnt = 0; - for (Entry> entry : this.topicQueueTable.entrySet()) { - List qdList = entry.getValue(); - for (QueueData qd : qdList) { - if (qd.getBrokerName().equals(brokerName)) { + for (Map.Entry> entry : topicQueueTable.entrySet()) { + String topic = entry.getKey(); + Map queueDataMap = entry.getValue(); + + if (queueDataMap != null) { + QueueData qd = queueDataMap.get(brokerName); + if (qd != null) { int perm = qd.getPerm(); switch (requestCode) { case RequestCode.WIPE_WRITE_PERM_OF_BROKER: @@ -293,6 +292,7 @@ public class RouteInfoManager { break; } qd.setPerm(perm); + topicCnt++; } } @@ -302,17 +302,17 @@ public class RouteInfoManager { } public void unregisterBroker( - final String clusterName, - final String brokerAddr, - final String brokerName, - final long brokerId) { + final String clusterName, + final String brokerAddr, + final String brokerName, + final long brokerId) { try { try { this.lock.writeLock().lockInterruptibly(); BrokerLiveInfo brokerLiveInfo = this.brokerLiveTable.remove(brokerAddr); log.info("unregisterBroker, remove from brokerLiveTable {}, {}", - brokerLiveInfo != null ? "OK" : "Failed", - brokerAddr + brokerLiveInfo != null ? "OK" : "Failed", + brokerAddr ); this.filterServerTable.remove(brokerAddr); @@ -322,14 +322,14 @@ public class RouteInfoManager { if (null != brokerData) { String addr = brokerData.getBrokerAddrs().remove(brokerId); log.info("unregisterBroker, remove addr from brokerAddrTable {}, {}", - addr != null ? "OK" : "Failed", - brokerAddr + addr != null ? "OK" : "Failed", + brokerAddr ); if (brokerData.getBrokerAddrs().isEmpty()) { this.brokerAddrTable.remove(brokerName); log.info("unregisterBroker, remove name from brokerAddrTable OK, {}", - brokerName + brokerName ); removeBrokerName = true; @@ -341,13 +341,13 @@ public class RouteInfoManager { if (nameSet != null) { boolean removed = nameSet.remove(brokerName); log.info("unregisterBroker, remove name from clusterAddrTable {}, {}", - removed ? "OK" : "Failed", - brokerName); + removed ? "OK" : "Failed", + brokerName); if (nameSet.isEmpty()) { this.clusterAddrTable.remove(clusterName); log.info("unregisterBroker, remove cluster from clusterAddrTable {}", - clusterName + clusterName ); } } @@ -362,26 +362,21 @@ public class RouteInfoManager { } private void removeTopicByBrokerName(final String brokerName) { - Iterator>> itMap = this.topicQueueTable.entrySet().iterator(); - while (itMap.hasNext()) { - Entry> entry = itMap.next(); + Set noBrokerRegisterTopic = new HashSet<>(); - String topic = entry.getKey(); - List queueDataList = entry.getValue(); - Iterator it = queueDataList.iterator(); - while (it.hasNext()) { - QueueData qd = it.next(); - if (qd.getBrokerName().equals(brokerName)) { - log.info("removeTopicByBrokerName, remove one broker's topic {} {}", topic, qd); - it.remove(); - } + this.topicQueueTable.forEach((topic, queueDataMap) -> { + QueueData old = queueDataMap.remove(brokerName); + if (old != null) { + log.info("removeTopicByBrokerName, remove one broker's topic {} {}", topic, old); } - if (queueDataList.isEmpty()) { + if (queueDataMap.size() == 0) { + noBrokerRegisterTopic.add(topic); log.info("removeTopicByBrokerName, remove the topic all queue {}", topic); - itMap.remove(); } - } + }); + + noBrokerRegisterTopic.forEach(topicQueueTable::remove); } public TopicRouteData pickupTopicRouteData(final String topic) { @@ -398,27 +393,31 @@ public class RouteInfoManager { try { try { this.lock.readLock().lockInterruptibly(); - List queueDataList = this.topicQueueTable.get(topic); - if (queueDataList != null) { - topicRouteData.setQueueDatas(queueDataList); + Map queueDataMap = this.topicQueueTable.get(topic); + if (queueDataMap != null) { + topicRouteData.setQueueDatas(new ArrayList<>(queueDataMap.values())); foundQueueData = true; - Iterator it = queueDataList.iterator(); - while (it.hasNext()) { - QueueData qd = it.next(); - brokerNameSet.add(qd.getBrokerName()); - } + brokerNameSet.addAll(queueDataMap.keySet()); for (String brokerName : brokerNameSet) { BrokerData brokerData = this.brokerAddrTable.get(brokerName); if (null != brokerData) { BrokerData brokerDataClone = new BrokerData(brokerData.getCluster(), brokerData.getBrokerName(), (HashMap) brokerData - .getBrokerAddrs().clone()); + .getBrokerAddrs().clone()); brokerDataList.add(brokerDataClone); foundBrokerData = true; - for (final String brokerAddr : brokerDataClone.getBrokerAddrs().values()) { - List filterServerList = this.filterServerTable.get(brokerAddr); - filterServerMap.put(brokerAddr, filterServerList); + + // skip if filter server table is empty + if (!filterServerTable.isEmpty()) { + for (final String brokerAddr : brokerDataClone.getBrokerAddrs().values()) { + List filterServerList = this.filterServerTable.get(brokerAddr); + + // only add filter server list when not null + if (filterServerList != null) { + filterServerMap.put(brokerAddr, filterServerList); + } + } } } } @@ -439,7 +438,8 @@ public class RouteInfoManager { return null; } - public void scanNotActiveBroker() { + public int scanNotActiveBroker() { + int removeCount = 0; Iterator> it = this.brokerLiveTable.entrySet().iterator(); while (it.hasNext()) { Entry next = it.next(); @@ -449,8 +449,12 @@ public class RouteInfoManager { it.remove(); log.warn("The broker channel expired, {} {}ms", next.getKey(), BROKER_CHANNEL_EXPIRED_TIME); this.onChannelDestroy(next.getKey(), next.getValue().getChannel()); + + removeCount++; } } + + return removeCount; } public void onChannelDestroy(String remoteAddr, Channel channel) { @@ -460,7 +464,7 @@ public class RouteInfoManager { try { this.lock.readLock().lockInterruptibly(); Iterator> itBrokerLiveTable = - this.brokerLiveTable.entrySet().iterator(); + this.brokerLiveTable.entrySet().iterator(); while (itBrokerLiveTable.hasNext()) { Entry entry = itBrokerLiveTable.next(); if (entry.getValue().getChannel() == channel) { @@ -492,7 +496,7 @@ public class RouteInfoManager { String brokerNameFound = null; boolean removeBrokerName = false; Iterator> itBrokerAddrTable = - this.brokerAddrTable.entrySet().iterator(); + this.brokerAddrTable.entrySet().iterator(); while (itBrokerAddrTable.hasNext() && (null == brokerNameFound)) { BrokerData brokerData = itBrokerAddrTable.next().getValue(); @@ -505,7 +509,7 @@ public class RouteInfoManager { brokerNameFound = brokerData.getBrokerName(); it.remove(); log.info("remove brokerAddr[{}, {}] from brokerAddrTable, because channel destroyed", - brokerId, brokerAddr); + brokerId, brokerAddr); break; } } @@ -514,7 +518,7 @@ public class RouteInfoManager { removeBrokerName = true; itBrokerAddrTable.remove(); log.info("remove brokerName[{}] from brokerAddrTable, because channel destroyed", - brokerData.getBrokerName()); + brokerData.getBrokerName()); } } @@ -527,11 +531,11 @@ public class RouteInfoManager { boolean removed = brokerNames.remove(brokerNameFound); if (removed) { log.info("remove brokerName[{}], clusterName[{}] from clusterAddrTable, because channel destroyed", - brokerNameFound, clusterName); + brokerNameFound, clusterName); if (brokerNames.isEmpty()) { log.info("remove the clusterName[{}] from clusterAddrTable, because channel destroyed and no broker in this cluster", - clusterName); + clusterName); it.remove(); } @@ -541,29 +545,22 @@ public class RouteInfoManager { } if (removeBrokerName) { - Iterator>> itTopicQueueTable = - this.topicQueueTable.entrySet().iterator(); - while (itTopicQueueTable.hasNext()) { - Entry> entry = itTopicQueueTable.next(); - String topic = entry.getKey(); - List queueDataList = entry.getValue(); + String finalBrokerNameFound = brokerNameFound; + Set needRemoveTopic = new HashSet<>(); - Iterator 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); - } - } + topicQueueTable.forEach((topic, queueDataMap) -> { + QueueData old = queueDataMap.remove(finalBrokerNameFound); + log.info("remove topic[{} {}], from topicQueueTable, because channel destroyed", + topic, old); - if (queueDataList.isEmpty()) { - itTopicQueueTable.remove(); + if (queueDataMap.size() == 0) { log.info("remove topic[{}] all queue, from topicQueueTable, because channel destroyed", - topic); + topic); + needRemoveTopic.add(topic); } - } + }); + + needRemoveTopic.forEach(topicQueueTable::remove); } } finally { this.lock.writeLock().unlock(); @@ -581,9 +578,9 @@ public class RouteInfoManager { log.info("--------------------------------------------------------"); { log.info("topicQueueTable SIZE: {}", this.topicQueueTable.size()); - Iterator>> it = this.topicQueueTable.entrySet().iterator(); + Iterator>> it = this.topicQueueTable.entrySet().iterator(); while (it.hasNext()) { - Entry> next = it.next(); + Entry> next = it.next(); log.info("topicQueueTable Topic: {} {}", next.getKey(), next.getValue()); } } @@ -622,7 +619,7 @@ public class RouteInfoManager { } } - public byte[] getSystemTopicList() { + public TopicList getSystemTopicList() { TopicList topicList = new TopicList(); try { try { @@ -651,30 +648,24 @@ public class RouteInfoManager { log.error("getAllTopicList Exception", e); } - return topicList.encode(); + return topicList; } - public byte[] getTopicsByCluster(String cluster) { + public TopicList getTopicsByCluster(String cluster) { TopicList topicList = new TopicList(); try { try { this.lock.readLock().lockInterruptibly(); + Set brokerNameSet = this.clusterAddrTable.get(cluster); for (String brokerName : brokerNameSet) { - Iterator>> topicTableIt = - this.topicQueueTable.entrySet().iterator(); - while (topicTableIt.hasNext()) { - Entry> topicEntry = topicTableIt.next(); - String topic = topicEntry.getKey(); - List queueDatas = topicEntry.getValue(); - for (QueueData queueData : queueDatas) { - if (brokerName.equals(queueData.getBrokerName())) { - topicList.getTopicList().add(topic); - break; - } + this.topicQueueTable.forEach((topic, queueDataMap) -> { + if (queueDataMap.containsKey(brokerName)) { + topicList.getTopicList().add(topic); } - } + }); } + } finally { this.lock.readLock().unlock(); } @@ -682,25 +673,39 @@ public class RouteInfoManager { log.error("getAllTopicList Exception", e); } - return topicList.encode(); + return topicList; } - public byte[] getUnitTopics() { + public TopicList getUnitTopics() { + return topicQueueTableIter(qd -> TopicSysFlag.hasUnitFlag(qd.getTopicSysFlag())); + } + + public TopicList getHasUnitSubTopicList() { + return topicQueueTableIter(qd -> TopicSysFlag.hasUnitSubFlag(qd.getTopicSysFlag())); + } + + public TopicList getHasUnitSubUnUnitTopicList() { + return topicQueueTableIter(qd -> !TopicSysFlag.hasUnitFlag(qd.getTopicSysFlag()) + && TopicSysFlag.hasUnitSubFlag(qd.getTopicSysFlag())); + } + + private TopicList topicQueueTableIter(Predicate pickCondition) { TopicList topicList = new TopicList(); try { try { this.lock.readLock().lockInterruptibly(); - Iterator>> topicTableIt = - this.topicQueueTable.entrySet().iterator(); - while (topicTableIt.hasNext()) { - Entry> topicEntry = topicTableIt.next(); - String topic = topicEntry.getKey(); - List queueDatas = topicEntry.getValue(); - if (queueDatas != null && queueDatas.size() > 0 - && TopicSysFlag.hasUnitFlag(queueDatas.get(0).getTopicSysFlag())) { - topicList.getTopicList().add(topic); + + topicQueueTable.forEach((topic, queueDataMap) -> { + for (QueueData qd : queueDataMap.values()) { + if (pickCondition.test(qd)) { + topicList.getTopicList().add(topic); + } + + // we need only one queue data here + break; } - } + }); + } finally { this.lock.readLock().unlock(); } @@ -708,60 +713,7 @@ public class RouteInfoManager { log.error("getAllTopicList Exception", e); } - return topicList.encode(); - } - - public byte[] getHasUnitSubTopicList() { - TopicList topicList = new TopicList(); - try { - try { - this.lock.readLock().lockInterruptibly(); - Iterator>> topicTableIt = - this.topicQueueTable.entrySet().iterator(); - while (topicTableIt.hasNext()) { - Entry> topicEntry = topicTableIt.next(); - String topic = topicEntry.getKey(); - List queueDatas = topicEntry.getValue(); - if (queueDatas != null && queueDatas.size() > 0 - && TopicSysFlag.hasUnitSubFlag(queueDatas.get(0).getTopicSysFlag())) { - topicList.getTopicList().add(topic); - } - } - } finally { - this.lock.readLock().unlock(); - } - } catch (Exception e) { - log.error("getAllTopicList Exception", e); - } - - return topicList.encode(); - } - - public byte[] getHasUnitSubUnUnitTopicList() { - TopicList topicList = new TopicList(); - try { - try { - this.lock.readLock().lockInterruptibly(); - Iterator>> topicTableIt = - this.topicQueueTable.entrySet().iterator(); - while (topicTableIt.hasNext()) { - Entry> topicEntry = topicTableIt.next(); - String topic = topicEntry.getKey(); - List queueDatas = topicEntry.getValue(); - if (queueDatas != null && queueDatas.size() > 0 - && !TopicSysFlag.hasUnitFlag(queueDatas.get(0).getTopicSysFlag()) - && TopicSysFlag.hasUnitSubFlag(queueDatas.get(0).getTopicSysFlag())) { - topicList.getTopicList().add(topic); - } - } - } finally { - this.lock.readLock().unlock(); - } - } catch (Exception e) { - log.error("getAllTopicList Exception", e); - } - - return topicList.encode(); + return topicList; } } @@ -772,7 +724,7 @@ class BrokerLiveInfo { private String haServerAddr; public BrokerLiveInfo(long lastUpdateTimestamp, DataVersion dataVersion, Channel channel, - String haServerAddr) { + String haServerAddr) { this.lastUpdateTimestamp = lastUpdateTimestamp; this.dataVersion = dataVersion; this.channel = channel; @@ -814,6 +766,6 @@ class BrokerLiveInfo { @Override public String toString() { return "BrokerLiveInfo [lastUpdateTimestamp=" + lastUpdateTimestamp + ", dataVersion=" + dataVersion - + ", channel=" + channel + ", haServerAddr=" + haServerAddr + "]"; + + ", channel=" + channel + ", haServerAddr=" + haServerAddr + "]"; } } diff --git a/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerPermTest.java b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerPermTest.java new file mode 100644 index 0000000000..91532d8736 --- /dev/null +++ b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerPermTest.java @@ -0,0 +1,111 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.rocketmq.namesrv.routeinfo; + +import org.apache.rocketmq.common.constant.PermName; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.QueueData; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import java.lang.reflect.Field; +import java.util.HashMap; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +public class RouteInfoManagerBrokerPermTest extends RouteInfoManagerTestBase { + private static RouteInfoManager routeInfoManager; + public static String clusterName = "cluster"; + public static String brokerPrefix = "broker"; + public static String topicPrefix = "topic"; + + public static RouteInfoManagerTestBase.Cluster cluster; + + @Before + public void setup() { + routeInfoManager = new RouteInfoManager(); + cluster = registerCluster(routeInfoManager, + clusterName, + brokerPrefix, + 3, + 3, + topicPrefix, + 10); + } + + @After + public void terminate() { + routeInfoManager.printAllPeriodically(); + + for (BrokerData bd : cluster.brokerDataMap.values()) { + unregisterBrokerAll(routeInfoManager, bd); + } + } + + @Test + public void testAddWritePermOfBrokerByLock() throws Exception { + String brokerName = getBrokerName(brokerPrefix,0); + String topicName = getTopicName(topicPrefix,0); + + + QueueData qd = new QueueData(); + qd.setPerm(PermName.PERM_READ); + qd.setBrokerName(brokerName); + + HashMap> topicQueueTable = new HashMap<>(); + + Map queueDataMap = new HashMap<>(); + queueDataMap.put(brokerName, qd); + topicQueueTable.put(topicName, queueDataMap); + + Field filed = RouteInfoManager.class.getDeclaredField("topicQueueTable"); + filed.setAccessible(true); + filed.set(routeInfoManager, topicQueueTable); + + int addTopicCnt = routeInfoManager.addWritePermOfBrokerByLock(brokerName); + assertThat(addTopicCnt).isEqualTo(1); + assertThat(qd.getPerm()).isEqualTo(PermName.PERM_READ | PermName.PERM_WRITE); + + } + + @Test + public void testWipeWritePermOfBrokerByLock() throws Exception { + String brokerName = getBrokerName(brokerPrefix,0); + String topicName = getTopicName(topicPrefix,0); + + QueueData qd = new QueueData(); + qd.setPerm(PermName.PERM_READ); + qd.setBrokerName(brokerName); + + HashMap> topicQueueTable = new HashMap<>(); + + Map queueDataMap = new HashMap<>(); + queueDataMap.put(brokerName, qd); + topicQueueTable.put(topicName, queueDataMap); + + Field filed = RouteInfoManager.class.getDeclaredField("topicQueueTable"); + filed.setAccessible(true); + filed.set(routeInfoManager, topicQueueTable); + + int addTopicCnt = routeInfoManager.wipeWritePermOfBrokerByLock(brokerName); + assertThat(addTopicCnt).isEqualTo(1); + assertThat(qd.getPerm()).isEqualTo(PermName.PERM_READ); + + } +} diff --git a/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerRegisterTest.java b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerRegisterTest.java new file mode 100644 index 0000000000..19ab058554 --- /dev/null +++ b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerBrokerRegisterTest.java @@ -0,0 +1,124 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.rocketmq.namesrv.routeinfo; + +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.TopicRouteData; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.HashMap; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; + +public class RouteInfoManagerBrokerRegisterTest extends RouteInfoManagerTestBase { + private static RouteInfoManager routeInfoManager; + public static String clusterName = "cluster"; + public static String brokerPrefix = "broker"; + public static String topicPrefix = "topic"; + public static int brokerPerName = 3; + public static int brokerNameNumber = 3; + + public static RouteInfoManagerTestBase.Cluster cluster; + + @Before + public void setup() { + routeInfoManager = new RouteInfoManager(); + cluster = registerCluster(routeInfoManager, + clusterName, + brokerPrefix, + brokerNameNumber, + brokerPerName, + topicPrefix, + 10); + } + + @After + public void terminate() { + routeInfoManager.printAllPeriodically(); + + for (BrokerData bd : cluster.brokerDataMap.values()) { + unregisterBrokerAll(routeInfoManager, bd); + } + } + + @Test + public void testScanNotActiveBroker() { + for (int j = 0; j < brokerNameNumber; j++) { + String brokerName = getBrokerName(brokerPrefix, j); + + for (int i = 0; i < brokerPerName; i++) { + String brokerAddr = getBrokerAddr(clusterName, brokerName, i); + + // set not active + routeInfoManager.updateBrokerInfoUpdateTimestamp(brokerAddr, 0); + + assertEquals(1, routeInfoManager.scanNotActiveBroker()); + } + } + + } + + @Test + public void testMasterChangeFromSlave() { + String topicName = getTopicName(topicPrefix, 0); + String brokerName = getBrokerName(brokerPrefix, 0); + + String originMasterAddr = getBrokerAddr(clusterName, brokerName, MixAll.MASTER_ID); + TopicRouteData topicRouteData = routeInfoManager.pickupTopicRouteData(topicName); + BrokerData brokerDataOrigin = findBrokerDataByBrokerName(topicRouteData.getBrokerDatas(), brokerName); + + // check origin master address + Assert.assertEquals(brokerDataOrigin.getBrokerAddrs().get(MixAll.MASTER_ID), originMasterAddr); + + // master changed + String newMasterAddr = getBrokerAddr(clusterName, brokerName, 1); + registerBrokerWithTopicConfig(routeInfoManager, + clusterName, + newMasterAddr, + brokerName, + MixAll.MASTER_ID, + newMasterAddr, + cluster.topicConfig, + new ArrayList<>()); + + topicRouteData = routeInfoManager.pickupTopicRouteData(topicName); + brokerDataOrigin = findBrokerDataByBrokerName(topicRouteData.getBrokerDatas(), brokerName); + + // check new master address + assertEquals(brokerDataOrigin.getBrokerAddrs().get(MixAll.MASTER_ID), newMasterAddr); + } + + @Test + public void testUnregisterBroker() { + String topicName = getTopicName(topicPrefix, 0); + String brokerName = getBrokerName(brokerPrefix, 0); + long unregisterBrokerId = 2; + + unregisterBroker(routeInfoManager, cluster.brokerDataMap.get(brokerName), unregisterBrokerId); + + TopicRouteData topicRouteData = routeInfoManager.pickupTopicRouteData(topicName); + HashMap brokerAddrs = findBrokerDataByBrokerName(topicRouteData.getBrokerDatas(), brokerName).getBrokerAddrs(); + + assertFalse(brokerAddrs.containsKey(unregisterBrokerId)); + } +} diff --git a/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerStaticRegisterTest.java b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerStaticRegisterTest.java new file mode 100644 index 0000000000..427e74f0bf --- /dev/null +++ b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerStaticRegisterTest.java @@ -0,0 +1,153 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.rocketmq.namesrv.routeinfo; + +import org.apache.rocketmq.common.TopicConfig; +import org.apache.rocketmq.common.protocol.body.ClusterInfo; +import org.apache.rocketmq.common.protocol.body.TopicList; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.QueueData; +import org.apache.rocketmq.common.protocol.route.TopicRouteData; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; + +public class RouteInfoManagerStaticRegisterTest extends RouteInfoManagerTestBase { + private static RouteInfoManager routeInfoManager; + public static String clusterName = "cluster"; + public static String brokerPrefix = "broker"; + public static String topicPrefix = "topic"; + + public static RouteInfoManagerTestBase.Cluster cluster; + + @Before + public void setup() { + routeInfoManager = new RouteInfoManager(); + cluster = registerCluster(routeInfoManager, + clusterName, + brokerPrefix, + 3, + 3, + topicPrefix, + 10); + } + + @After + public void terminate() { + routeInfoManager.printAllPeriodically(); + + for (BrokerData bd : cluster.brokerDataMap.values()) { + unregisterBrokerAll(routeInfoManager, bd); + } + } + + @Test + public void testGetAllClusterInfo() { + ClusterInfo clusterInfo = routeInfoManager.getAllClusterInfo(); + HashMap> clusterAddrTable = clusterInfo.getClusterAddrTable(); + + assertEquals(1, clusterAddrTable.size()); + assertEquals(cluster.getAllBrokerName(), clusterAddrTable.get(clusterName)); + } + + @Test + public void testGetAllTopicList() { + TopicList topicInfo = routeInfoManager.getAllTopicList(); + + assertEquals(cluster.getAllTopicName(), topicInfo.getTopicList()); + } + + @Test + public void testGetTopicsByCluster() { + TopicList topicList = routeInfoManager.getTopicsByCluster(clusterName); + assertEquals(cluster.getAllTopicName(), topicList.getTopicList()); + } + + @Test + public void testPickupTopicRouteData() { + String topic = getTopicName(topicPrefix, 0); + + TopicRouteData topicRouteData = routeInfoManager.pickupTopicRouteData(topic); + + TopicConfig topicConfig = cluster.topicConfig.get(topic); + + // check broker data + Collections.sort(topicRouteData.getBrokerDatas()); + List ans = new ArrayList<>(cluster.brokerDataMap.values()); + Collections.sort(ans); + + assertEquals(topicRouteData.getBrokerDatas(), ans); + + // check queue data + HashSet allBrokerNameInQueueData = new HashSet<>(); + + for (QueueData queueData : topicRouteData.getQueueDatas()) { + allBrokerNameInQueueData.add(queueData.getBrokerName()); + + assertEquals(queueData.getWriteQueueNums(), topicConfig.getWriteQueueNums()); + assertEquals(queueData.getReadQueueNums(), topicConfig.getReadQueueNums()); + assertEquals(queueData.getPerm(), topicConfig.getPerm()); + assertEquals(queueData.getTopicSysFlag(), topicConfig.getTopicSysFlag()); + } + + assertEquals(allBrokerNameInQueueData, new HashSet<>(cluster.getAllBrokerName())); + } + + @Test + public void testDeleteTopic() { + String topic = getTopicName(topicPrefix, 0); + routeInfoManager.deleteTopic(topic); + + assertNull(routeInfoManager.pickupTopicRouteData(topic)); + } + + @Test + public void testGetSystemTopicList() { + TopicList topicList = routeInfoManager.getSystemTopicList(); + assertThat(topicList).isNotNull(); + } + + @Test + public void testGetUnitTopics() { + TopicList topicList = routeInfoManager.getUnitTopics(); + assertThat(topicList).isNotNull(); + } + + @Test + public void testGetHasUnitSubTopicList() { + TopicList topicList = routeInfoManager.getHasUnitSubTopicList(); + assertThat(topicList).isNotNull(); + } + + @Test + public void testGetHasUnitSubUnUnitTopicList() { + TopicList topicList = routeInfoManager.getHasUnitSubUnUnitTopicList(); + assertThat(topicList).isNotNull(); + } + +} \ No newline at end of file diff --git a/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTest.java b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTest.java deleted file mode 100644 index e0d9e1871c..0000000000 --- a/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTest.java +++ /dev/null @@ -1,162 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.rocketmq.namesrv.routeinfo; - -import io.netty.channel.Channel; -import org.apache.rocketmq.common.TopicConfig; -import org.apache.rocketmq.common.constant.PermName; -import org.apache.rocketmq.common.namesrv.RegisterBrokerResult; -import org.apache.rocketmq.common.protocol.body.TopicConfigSerializeWrapper; -import org.apache.rocketmq.common.protocol.route.QueueData; -import org.apache.rocketmq.common.protocol.route.TopicRouteData; -import org.junit.After; -import org.junit.Assert; -import org.junit.Before; -import org.junit.Test; - -import java.lang.reflect.Field; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.concurrent.ConcurrentHashMap; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.mock; - -public class RouteInfoManagerTest { - - private static RouteInfoManager routeInfoManager; - - @Before - public void setup() { - routeInfoManager = new RouteInfoManager(); - testRegisterBroker(); - } - - @After - public void terminate() { - routeInfoManager.printAllPeriodically(); - routeInfoManager.unregisterBroker("default-cluster", "127.0.0.1:10911", "default-broker", 1234); - } - - @Test - public void testGetAllClusterInfo() { - byte[] clusterInfo = routeInfoManager.getAllClusterInfo(); - assertThat(clusterInfo).isNotNull(); - } - - @Test - public void testGetAllTopicList() { - byte[] topicInfo = routeInfoManager.getAllTopicList(); - Assert.assertTrue(topicInfo != null); - assertThat(topicInfo).isNotNull(); - } - - @Test - public void testRegisterBroker() { - TopicConfigSerializeWrapper topicConfigSerializeWrapper = new TopicConfigSerializeWrapper(); - ConcurrentHashMap topicConfigConcurrentHashMap = new ConcurrentHashMap<>(); - TopicConfig topicConfig = new TopicConfig(); - topicConfig.setWriteQueueNums(8); - topicConfig.setTopicName("unit-test"); - topicConfig.setPerm(6); - topicConfig.setReadQueueNums(8); - topicConfig.setOrder(false); - topicConfigConcurrentHashMap.put("unit-test", topicConfig); - topicConfigSerializeWrapper.setTopicConfigTable(topicConfigConcurrentHashMap); - Channel channel = mock(Channel.class); - RegisterBrokerResult registerBrokerResult = routeInfoManager.registerBroker("default-cluster", "127.0.0.1:10911", "default-broker", 1234, "127.0.0.1:1001", - topicConfigSerializeWrapper, new ArrayList(), channel); - assertThat(registerBrokerResult).isNotNull(); - } - - @Test - public void testWipeWritePermOfBrokerByLock() throws Exception { - List qdList = new ArrayList<>(); - QueueData qd = new QueueData(); - qd.setPerm(PermName.PERM_READ | PermName.PERM_WRITE); - qd.setBrokerName("broker-a"); - qdList.add(qd); - HashMap> topicQueueTable = new HashMap<>(); - topicQueueTable.put("topic-a", qdList); - - Field filed = RouteInfoManager.class.getDeclaredField("topicQueueTable"); - filed.setAccessible(true); - filed.set(routeInfoManager, topicQueueTable); - - int addTopicCnt = routeInfoManager.wipeWritePermOfBrokerByLock("broker-a"); - assertThat(addTopicCnt).isEqualTo(1); - assertThat(qd.getPerm()).isEqualTo(PermName.PERM_READ); - - } - - @Test - public void testPickupTopicRouteData() { - TopicRouteData result = routeInfoManager.pickupTopicRouteData("unit_test"); - assertThat(result).isNull(); - } - - @Test - public void testGetSystemTopicList() { - byte[] topicList = routeInfoManager.getSystemTopicList(); - assertThat(topicList).isNotNull(); - } - - @Test - public void testGetTopicsByCluster() { - byte[] topicList = routeInfoManager.getTopicsByCluster("default-cluster"); - assertThat(topicList).isNotNull(); - } - - @Test - public void testGetUnitTopics() { - byte[] topicList = routeInfoManager.getUnitTopics(); - assertThat(topicList).isNotNull(); - } - - @Test - public void testGetHasUnitSubTopicList() { - byte[] topicList = routeInfoManager.getHasUnitSubTopicList(); - assertThat(topicList).isNotNull(); - } - - @Test - public void testGetHasUnitSubUnUnitTopicList() { - byte[] topicList = routeInfoManager.getHasUnitSubUnUnitTopicList(); - assertThat(topicList).isNotNull(); - } - - @Test - public void testAddWritePermOfBrokerByLock() throws Exception { - List qdList = new ArrayList<>(); - QueueData qd = new QueueData(); - qd.setPerm(PermName.PERM_READ); - qd.setBrokerName("broker-a"); - qdList.add(qd); - HashMap> topicQueueTable = new HashMap<>(); - topicQueueTable.put("topic-a", qdList); - - Field filed = RouteInfoManager.class.getDeclaredField("topicQueueTable"); - filed.setAccessible(true); - filed.set(routeInfoManager, topicQueueTable); - - int addTopicCnt = routeInfoManager.addWritePermOfBrokerByLock("broker-a"); - assertThat(addTopicCnt).isEqualTo(1); - assertThat(qd.getPerm()).isEqualTo(PermName.PERM_READ | PermName.PERM_WRITE); - - } -} \ No newline at end of file diff --git a/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTestBase.java b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTestBase.java new file mode 100644 index 0000000000..a1a56bfc2b --- /dev/null +++ b/namesrv/src/test/java/org/apache/rocketmq/namesrv/routeinfo/RouteInfoManagerTestBase.java @@ -0,0 +1,188 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.rocketmq.namesrv.routeinfo; + +import io.netty.channel.Channel; +import io.netty.channel.embedded.EmbeddedChannel; +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.TopicConfig; +import org.apache.rocketmq.common.namesrv.RegisterBrokerResult; +import org.apache.rocketmq.common.protocol.body.TopicConfigSerializeWrapper; +import org.apache.rocketmq.common.protocol.route.BrokerData; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; + +public class RouteInfoManagerTestBase { + + protected static class Cluster { + ConcurrentMap topicConfig; + Map brokerDataMap; + + public Cluster(ConcurrentMap topicConfig, Map brokerData) { + this.topicConfig = topicConfig; + this.brokerDataMap = brokerData; + } + + public Set getAllBrokerName() { + return brokerDataMap.keySet(); + } + + public Set getAllTopicName() { + return topicConfig.keySet(); + } + } + + protected Cluster registerCluster(RouteInfoManager routeInfoManager, String cluster, + String brokerNamePrefix, + int brokerNameNumber, + int brokerPerName, + String topicPrefix, + int topicNumber) { + + Map brokerDataMap = new HashMap<>(); + + // no filterServer address + List filterServerAddr = new ArrayList<>(); + + ConcurrentMap topicConfig = genTopicConfig(topicPrefix, topicNumber); + + for (int i = 0; i < brokerNameNumber; i++) { + String brokerName = getBrokerName(brokerNamePrefix, i); + + BrokerData brokerData = genBrokerData(cluster, brokerName, brokerPerName, true); + + // avoid object reference copy + ConcurrentMap topicConfigForBroker = genTopicConfig(topicPrefix, topicNumber); + + registerBrokerWithTopicConfig(routeInfoManager, brokerData, topicConfigForBroker, filterServerAddr); + + // avoid object reference copy + brokerDataMap.put(brokerData.getBrokerName(), genBrokerData(cluster, brokerName, brokerPerName, true)); + } + + return new Cluster(topicConfig, brokerDataMap); + } + + protected String getBrokerAddr(String cluster, String brokerName, long brokerNumber) { + return cluster + "-" + brokerName + ":" + brokerNumber; + } + + protected BrokerData genBrokerData(String clusterName, String brokerName, long totalBrokerNumber, boolean hasMaster) { + HashMap brokerAddrMap = new HashMap<>(); + + long startId = 0; + if (hasMaster) { + brokerAddrMap.put(MixAll.MASTER_ID, getBrokerAddr(clusterName, brokerName, MixAll.MASTER_ID)); + startId = 1; + } + + for (long i = startId; i < totalBrokerNumber; i++) { + brokerAddrMap.put(i, getBrokerAddr(clusterName, brokerName, i)); + } + + return new BrokerData(clusterName, brokerName, brokerAddrMap); + } + + protected void registerBrokerWithTopicConfig(RouteInfoManager routeInfoManager, BrokerData brokerData, + ConcurrentMap topicConfigTable, + List filterServerAddr) { + + brokerData.getBrokerAddrs().forEach((brokerId, brokerAddr) -> { + registerBrokerWithTopicConfig(routeInfoManager, brokerData.getCluster(), + brokerAddr, + brokerData.getBrokerName(), + brokerId, + brokerAddr, // set ha server address the same as brokerAddr + new ConcurrentHashMap<>(topicConfigTable), + new ArrayList<>(filterServerAddr)); + }); + } + + protected void unregisterBrokerAll(RouteInfoManager routeInfoManager, BrokerData brokerData) { + for (Map.Entry entry : brokerData.getBrokerAddrs().entrySet()) { + routeInfoManager.unregisterBroker(brokerData.getCluster(), entry.getValue(), brokerData.getBrokerName(), entry.getKey()); + } + } + + protected void unregisterBroker(RouteInfoManager routeInfoManager, BrokerData brokerData, long brokerId) { + HashMap brokerAddrs = brokerData.getBrokerAddrs(); + if (brokerAddrs.containsKey(brokerId)) { + String address = brokerAddrs.remove(brokerId); + routeInfoManager.unregisterBroker(brokerData.getCluster(), address, brokerData.getBrokerName(), brokerId); + } + } + + protected RegisterBrokerResult registerBrokerWithTopicConfig(RouteInfoManager routeInfoManager, String clusterName, + String brokerAddr, + String brokerName, + long brokerId, + String haServerAddr, + ConcurrentMap topicConfigTable, + List filterServerAddr) { + + TopicConfigSerializeWrapper topicConfigSerializeWrapper = new TopicConfigSerializeWrapper(); + topicConfigSerializeWrapper.setTopicConfigTable(topicConfigTable); + + Channel channel = new EmbeddedChannel(); + return routeInfoManager.registerBroker(clusterName, + brokerAddr, + brokerName, + brokerId, + haServerAddr, + topicConfigSerializeWrapper, + filterServerAddr, + channel); + } + + + protected String getTopicName(String topicPrefix, int topicNumber) { + return topicPrefix + "-" + topicNumber; + } + + protected ConcurrentMap genTopicConfig(String topicPrefix, int topicNumber) { + ConcurrentMap topicConfigMap = new ConcurrentHashMap<>(); + + for (int i = 0; i < topicNumber; i++) { + String topicName = getTopicName(topicPrefix, i); + + TopicConfig topicConfig = new TopicConfig(); + topicConfig.setWriteQueueNums(8); + topicConfig.setTopicName(topicName); + topicConfig.setPerm(6); + topicConfig.setReadQueueNums(8); + topicConfig.setOrder(false); + topicConfigMap.put(topicName, topicConfig); + } + + return topicConfigMap; + } + + protected String getBrokerName(String brokerNamePrefix, long brokerNameNumber) { + return brokerNamePrefix + "-" + brokerNameNumber; + } + + protected BrokerData findBrokerDataByBrokerName(List data, String brokerName) { + return data.stream().filter(bd -> bd.getBrokerName().equals(brokerName)).findFirst().orElse(null); + } + +} From adbf3db7176fea568bb96ce9d379d7e151781825 Mon Sep 17 00:00:00 2001 From: wangfan Date: Thu, 10 Mar 2022 20:29:29 +0800 Subject: [PATCH 5/9] [ISSUE #3955] delete useless check --- .../rocketmq/broker/processor/SendMessageProcessor.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/SendMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/SendMessageProcessor.java index c8ea4d3b5a..0c1dec339a 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/SendMessageProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/SendMessageProcessor.java @@ -625,8 +625,6 @@ public class SendMessageProcessor extends AbstractSendMessageProcessor implement return handlePutMessageResultFuture(putMessageResult, response, request, messageExtBatch, responseHeader, mqtraceContext, ctx, queueIdInt); } - - public boolean hasConsumeMessageHook() { return consumeMessageHookList != null && !this.consumeMessageHookList.isEmpty(); } @@ -715,9 +713,6 @@ public class SendMessageProcessor extends AbstractSendMessageProcessor implement response.setCode(-1); super.msgCheck(ctx, requestHeader, response); - if (response.getCode() != -1) { - return response; - } return response; } From 2b5132117929faadd77a364b0fe5913990d10450 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=AD=E5=B0=8F=E6=BC=AA?= <644120242@qq.com> Date: Mon, 14 Mar 2022 20:47:36 +0800 Subject: [PATCH 6/9] [ISSUE#3983] Optimize warn log output when sending exceptions are encountered. (#3984) --- .../rocketmq/client/impl/producer/DefaultMQProducerImpl.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java index f3f9caf43a..ea80478604 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java @@ -630,9 +630,6 @@ public class DefaultMQProducerImpl implements MQProducerInner { this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, false); log.warn(String.format("sendKernelImpl exception, throw exception, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e); log.warn(msg.toString()); - - log.warn("sendKernelImpl exception", e); - log.warn(msg.toString()); throw e; } } else { From 03c5a3d171eb8667a7d13409d45723005c4c22f5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=AD=E5=B0=8F=E6=BC=AA?= <644120242@qq.com> Date: Mon, 14 Mar 2022 20:48:23 +0800 Subject: [PATCH 7/9] [ISSUE #3985] Remove shuffle operation before sorting the list of 'FaultItem'. (#3986) --- .../rocketmq/client/latency/LatencyFaultToleranceImpl.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/latency/LatencyFaultToleranceImpl.java b/client/src/main/java/org/apache/rocketmq/client/latency/LatencyFaultToleranceImpl.java index 827d97265f..750759f3d8 100644 --- a/client/src/main/java/org/apache/rocketmq/client/latency/LatencyFaultToleranceImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/latency/LatencyFaultToleranceImpl.java @@ -70,12 +70,8 @@ public class LatencyFaultToleranceImpl implements LatencyFaultTolerance final FaultItem faultItem = elements.nextElement(); tmpList.add(faultItem); } - if (!tmpList.isEmpty()) { - Collections.shuffle(tmpList); - Collections.sort(tmpList); - final int half = tmpList.size() / 2; if (half <= 0) { return tmpList.get(0).getName(); @@ -84,7 +80,6 @@ public class LatencyFaultToleranceImpl implements LatencyFaultTolerance return tmpList.get(i).getName(); } } - return null; } From c1aeca291ea686a2b7f01ffec30132f5a370a2ab Mon Sep 17 00:00:00 2001 From: sunxi92 Date: Tue, 15 Mar 2022 20:11:51 +0800 Subject: [PATCH 8/9] [#3942]If both acl and message trace are enabled and the default topic RMQ_SYS_TRACE_TOPIC is used for message trace, you don't need to add the PUB permission of RMQ_SYS_TRACE_TOPIC topic to the acl config. (#3943) * If both acl and message trace are enabled and the default topic RMQ_SYS_TRACE_TOPIC is used for message trace, you don't need to add the PUB permission of RMQ_SYS_TRACE_TOPIC topic to the acl config. * Delete Chinese character in comments. --- .../rocketmq/acl/plain/PlainPermissionManager.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainPermissionManager.java b/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainPermissionManager.java index 896b6f4f68..7fb9f0e4ca 100644 --- a/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainPermissionManager.java +++ b/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainPermissionManager.java @@ -46,6 +46,7 @@ import org.apache.rocketmq.common.DataVersion; import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.PlainAccessConfig; import org.apache.rocketmq.common.constant.LoggerName; +import org.apache.rocketmq.common.topic.TopicValidator; import org.apache.rocketmq.logging.InternalLogger; import org.apache.rocketmq.logging.InternalLoggerFactory; import org.apache.rocketmq.srvutil.AclFileWatchService; @@ -664,8 +665,18 @@ public class PlainPermissionManager { if (!signature.equals(plainAccessResource.getSignature())) { throw new AclException(String.format("Check signature failed for accessKey=%s", plainAccessResource.getAccessKey())); } - // Check perm of each resource + //Skip the topic RMQ_SYS_TRACE_TOPIC permission check,if the topic RMQ_SYS_TRACE_TOPIC is used for message trace + Map resourcePermMap = plainAccessResource.getResourcePermMap(); + if (resourcePermMap != null) { + Byte permission = resourcePermMap.get(TopicValidator.RMQ_SYS_TRACE_TOPIC); + if (permission != null && permission == Permission.PUB) { + return; + } + } + + + // Check perm of each resource checkPerm(plainAccessResource, ownedAccess); } From 6ff00ed0530bbd94e6ed9cc444b0a2fdb57c7ba3 Mon Sep 17 00:00:00 2001 From: zhangjidi2016 <1017543663@qq.com> Date: Sat, 19 Mar 2022 17:56:56 +0800 Subject: [PATCH 9/9] [ISSUE #4000]Fix the warn log input in command tools (#4001) Co-authored-by: zhangjidi --- .../java/org/apache/rocketmq/client/log/ClientLogger.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/log/ClientLogger.java b/client/src/main/java/org/apache/rocketmq/client/log/ClientLogger.java index b40f6a5159..7ee2ebaa4b 100644 --- a/client/src/main/java/org/apache/rocketmq/client/log/ClientLogger.java +++ b/client/src/main/java/org/apache/rocketmq/client/log/ClientLogger.java @@ -44,6 +44,8 @@ public class ClientLogger { private static final boolean CLIENT_USE_SLF4J; + private static Appender appenderProxy = new AppenderProxy(); + //private static Appender rocketmqClientAppender = null; static { @@ -53,6 +55,7 @@ public class ClientLogger { CLIENT_LOGGER = createLogger(LoggerName.CLIENT_LOGGER_NAME); createLogger(LoggerName.COMMON_LOGGER_NAME); createLogger(RemotingHelper.ROCKETMQ_REMOTING); + Logger.getRootLogger().addAppender(appenderProxy); } else { CLIENT_LOGGER = InternalLoggerFactory.getLogger(LoggerName.CLIENT_LOGGER_NAME); } @@ -76,7 +79,6 @@ public class ClientLogger { .withRollingFileAppender(logFileName, maxFileSize, maxFileIndex) .withAsync(false, queueSize).withName(ROCKETMQ_CLIENT_APPENDER_NAME).withLayout(layout).build(); - Logger.getRootLogger().addAppender(rocketmqClientAppender); return rocketmqClientAppender; } @@ -91,7 +93,7 @@ public class ClientLogger { // createClientAppender(); //} - realLogger.addAppender(new AppenderProxy()); + realLogger.addAppender(appenderProxy); realLogger.setLevel(Level.toLevel(clientLogLevel)); realLogger.setAdditivity(additive); return logger;