From b7f0162ffade63de1beec8e9271ae013401a666b Mon Sep 17 00:00:00 2001 From: HuiTong Date: Fri, 19 Aug 2022 12:26:47 +0800 Subject: [PATCH] fix thread-safety problem of admin tools (#4843) --- .../org/apache/rocketmq/common/admin/ConsumeStats.java | 10 ++++++---- .../apache/rocketmq/common/admin/TopicStatsTable.java | 10 ++++++---- .../rocketmq/common/protocol/body/TopicList.java | 5 +++-- .../rocketmq/tools/admin/DefaultMQAdminExtImpl.java | 4 ++-- 4 files changed, 17 insertions(+), 12 deletions(-) diff --git a/common/src/main/java/org/apache/rocketmq/common/admin/ConsumeStats.java b/common/src/main/java/org/apache/rocketmq/common/admin/ConsumeStats.java index 6b1c49290f..ae7e18dd2f 100644 --- a/common/src/main/java/org/apache/rocketmq/common/admin/ConsumeStats.java +++ b/common/src/main/java/org/apache/rocketmq/common/admin/ConsumeStats.java @@ -16,14 +16,16 @@ */ package org.apache.rocketmq.common.admin; -import java.util.HashMap; import java.util.Iterator; +import java.util.Map; import java.util.Map.Entry; +import java.util.concurrent.ConcurrentHashMap; + import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.remoting.protocol.RemotingSerializable; public class ConsumeStats extends RemotingSerializable { - private HashMap offsetTable = new HashMap(); + private Map offsetTable = new ConcurrentHashMap(); private double consumeTps = 0; public long computeTotalDiff() { @@ -39,11 +41,11 @@ public class ConsumeStats extends RemotingSerializable { return diffTotal; } - public HashMap getOffsetTable() { + public Map getOffsetTable() { return offsetTable; } - public void setOffsetTable(HashMap offsetTable) { + public void setOffsetTable(Map offsetTable) { this.offsetTable = offsetTable; } diff --git a/common/src/main/java/org/apache/rocketmq/common/admin/TopicStatsTable.java b/common/src/main/java/org/apache/rocketmq/common/admin/TopicStatsTable.java index 729075c061..42a8872dcb 100644 --- a/common/src/main/java/org/apache/rocketmq/common/admin/TopicStatsTable.java +++ b/common/src/main/java/org/apache/rocketmq/common/admin/TopicStatsTable.java @@ -16,18 +16,20 @@ */ package org.apache.rocketmq.common.admin; -import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.remoting.protocol.RemotingSerializable; public class TopicStatsTable extends RemotingSerializable { - private HashMap offsetTable = new HashMap(); + private Map offsetTable = new ConcurrentHashMap(); - public HashMap getOffsetTable() { + public Map getOffsetTable() { return offsetTable; } - public void setOffsetTable(HashMap offsetTable) { + public void setOffsetTable(Map offsetTable) { this.offsetTable = offsetTable; } } diff --git a/common/src/main/java/org/apache/rocketmq/common/protocol/body/TopicList.java b/common/src/main/java/org/apache/rocketmq/common/protocol/body/TopicList.java index baf8312479..9b9144e2c7 100644 --- a/common/src/main/java/org/apache/rocketmq/common/protocol/body/TopicList.java +++ b/common/src/main/java/org/apache/rocketmq/common/protocol/body/TopicList.java @@ -16,12 +16,13 @@ */ package org.apache.rocketmq.common.protocol.body; -import java.util.HashSet; import java.util.Set; +import java.util.concurrent.CopyOnWriteArraySet; + import org.apache.rocketmq.remoting.protocol.RemotingSerializable; public class TopicList extends RemotingSerializable { - private Set topicList = new HashSet(); + private Set topicList = new CopyOnWriteArraySet<>(); private String brokerAddr; public Set getTopicList() { diff --git a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java index 3cea455a2d..bb08e01194 100644 --- a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java +++ b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java @@ -797,7 +797,7 @@ public class DefaultMQAdminExtImpl implements MQAdminExt, MQAdminExtInner { } if (!hasConsumed) { - HashMap topicStatus = this.mqClientInstance.getMQClientAPIImpl().getTopicStatsInfo(brokerAddr, topic, timeoutMillis).getOffsetTable(); + Map topicStatus = this.mqClientInstance.getMQClientAPIImpl().getTopicStatsInfo(brokerAddr, topic, timeoutMillis).getOffsetTable(); for (int i = 0; i < queueData.getReadQueueNums(); i++) { MessageQueue queue = new MessageQueue(topic, queueData.getBrokerName(), i); OffsetWrapper offsetWrapper = new OffsetWrapper(); @@ -1107,7 +1107,7 @@ public class DefaultMQAdminExtImpl implements MQAdminExt, MQAdminExtInner { return adminToolExecute(new AdminToolHandler() { @Override public AdminToolResult doExecute() throws Exception { - final List spanSet = new ArrayList(); + final List spanSet = new CopyOnWriteArrayList<>(); TopicRouteData topicRouteData = examineTopicRouteInfo(topic); if (topicRouteData == null || topicRouteData.getBrokerDatas() == null || topicRouteData.getBrokerDatas().size() == 0) {