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) {

Reply via email to