This is an automated email from the ASF dual-hosted git repository.
lwclover pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git
The following commit(s) were added to refs/heads/develop by this push:
new 42a6d576ab Remove synchronized from deleteTopic. (#9997)
42a6d576ab is described below
commit 42a6d576abda5761f56c74623a2ae49d0d8b91b5
Author: rongtong <[email protected]>
AuthorDate: Thu Sep 17 18:34:49 2026 +0800
Remove synchronized from deleteTopic. (#9997)
Co-authored-by: RongtongJin <[email protected]>
---
.../broker/offset/ConsumerOffsetManager.java | 14 +++++++-------
.../broker/processor/AdminBrokerProcessor.java | 2 +-
.../broker/processor/PopInflightMessageCounter.java | 21 ++++++++++++---------
3 files changed, 20 insertions(+), 17 deletions(-)
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java
b/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java
index 1d3bf7bed0..598cd26d44 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java
@@ -86,21 +86,21 @@ public class ConsumerOffsetManager extends ConfigManager {
}
public void cleanOffsetByTopic(String topic) {
- Iterator<Entry<String, ConcurrentMap<Integer, Long>>> it =
this.offsetTable.entrySet().iterator();
- while (it.hasNext()) {
- Entry<String, ConcurrentMap<Integer, Long>> next = it.next();
- String topicAtGroup = next.getKey();
+
+ this.offsetTable.entrySet().removeIf(entry -> {
+ String topicAtGroup = entry.getKey();
if (topicAtGroup.contains(topic)) {
String[] arrays = topicAtGroup.split(TOPIC_GROUP_SEPARATOR);
if (arrays.length == 2 && topic.equals(arrays[0])) {
- it.remove();
removeConsumerOffset(topicAtGroup);
pullOffsetTable.remove(topicAtGroup);
resetOffsetTable.remove(topicAtGroup);
- LOG.warn("Clean topic's offset, {}, {}", topicAtGroup,
next.getValue());
+ LOG.warn("Clean topic's offset, {}, {}", topicAtGroup,
entry.getValue());
+ return true;
}
}
- }
+ return false;
+ });
}
public void scanUnsubscribedTopic() {
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
index 8083d7307c..604e3e18ca 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
@@ -762,7 +762,7 @@ public class AdminBrokerProcessor implements
NettyRequestProcessor {
return response;
}
- private synchronized RemotingCommand deleteTopic(ChannelHandlerContext ctx,
+ private RemotingCommand deleteTopic(ChannelHandlerContext ctx,
RemotingCommand request) throws RemotingCommandException {
final RemotingCommand response =
RemotingCommand.createResponseCommand(null);
DeleteTopicRequestHeader requestHeader =
diff --git
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopInflightMessageCounter.java
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopInflightMessageCounter.java
index 6749af3d75..8190a01fa7 100644
---
a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopInflightMessageCounter.java
+++
b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopInflightMessageCounter.java
@@ -24,7 +24,6 @@ import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
import org.apache.rocketmq.store.pop.PopCheckPoint;
import java.util.Map;
-import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
@@ -91,31 +90,35 @@ public class PopInflightMessageCounter {
}
public void clearInFlightMessageNumByGroupName(String group) {
- Set<String> topicGroupKey = this.topicInFlightMessageNum.keySet();
- for (String key : topicGroupKey) {
+ // Use removeIf for thread-safe removal from ConcurrentHashMap
+ this.topicInFlightMessageNum.entrySet().removeIf(entry -> {
+ String key = entry.getKey();
if (key.contains(group)) {
Pair<String, String> topicAndGroup = splitKey(key);
if (topicAndGroup != null &&
topicAndGroup.getObject2().equals(group)) {
- this.topicInFlightMessageNum.remove(key);
log.info("PopInflightMessageCounter#clearInFlightMessageNumByGroupName: clean
by group, topic={}, group={}",
topicAndGroup.getObject1(),
topicAndGroup.getObject2());
+ return true;
}
}
- }
+ return false;
+ });
}
public void clearInFlightMessageNumByTopicName(String topic) {
- Set<String> topicGroupKey = this.topicInFlightMessageNum.keySet();
- for (String key : topicGroupKey) {
+ // Use removeIf for thread-safe removal from ConcurrentHashMap
+ this.topicInFlightMessageNum.entrySet().removeIf(entry -> {
+ String key = entry.getKey();
if (key.contains(topic)) {
Pair<String, String> topicAndGroup = splitKey(key);
if (topicAndGroup != null &&
topicAndGroup.getObject1().equals(topic)) {
- this.topicInFlightMessageNum.remove(key);
log.info("PopInflightMessageCounter#clearInFlightMessageNumByTopicName: clean
by topic, topic={}, group={}",
topicAndGroup.getObject1(),
topicAndGroup.getObject2());
+ return true;
}
}
- }
+ return false;
+ });
}
public void clearInFlightMessageNum(String topic, String group, int
queueId) {