This is an automated email from the ASF dual-hosted git repository.

lizhimins 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 2daf0e2ca9 [ISSUE #10713] Optimize ConsumerFilterManager register CPU 
consumption (#10714)
2daf0e2ca9 is described below

commit 2daf0e2ca91a1592d18235d43e5d709d1c35d15f
Author: ymwneu <[email protected]>
AuthorDate: Tue Aug 4 12:04:19 2026 +0800

    [ISSUE #10713] Optimize ConsumerFilterManager register CPU consumption 
(#10714)
    
    - Refactor filter data index from topic-based to consumerGroup-based 
(SubscriptionFilterHandler)
    - Add topic existence check before registering filter
    - Generate BloomFilterData only when enableCalcFilterBitMap is enabled
    - Fix thread safety: use ConcurrentHashMap for topicSqlFilterData
    - Fix TOCTOU race conditions: replace containsKey+get with single get
    - Rebuild subscriptionFilterData from filterDataByTopic in decode
    - Add test for subscription shrink marking removed topics as dead
---
 .../broker/filter/ConsumerFilterManager.java       | 253 +++++++++++++--------
 .../broker/filter/ConsumerFilterManagerTest.java   |  45 +++-
 2 files changed, 194 insertions(+), 104 deletions(-)

diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/filter/ConsumerFilterManager.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/filter/ConsumerFilterManager.java
index 3a48f96b98..5c5600d155 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/filter/ConsumerFilterManager.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/filter/ConsumerFilterManager.java
@@ -22,11 +22,13 @@ import java.util.HashSet;
 import java.util.Iterator;
 import java.util.Map;
 import java.util.Map.Entry;
+import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentMap;
 import org.apache.rocketmq.broker.BrokerController;
 import org.apache.rocketmq.broker.BrokerPathConfigHelper;
 import org.apache.rocketmq.common.ConfigManager;
+import org.apache.rocketmq.common.TopicConfig;
 import org.apache.rocketmq.common.constant.LoggerName;
 import org.apache.rocketmq.common.filter.ExpressionType;
 import org.apache.rocketmq.filter.FilterFactory;
@@ -49,6 +51,9 @@ public class ConsumerFilterManager extends ConfigManager {
     private ConcurrentMap<String/*Topic*/, FilterDataMapByTopic>
         filterDataByTopic = new ConcurrentHashMap<>(256);
 
+    private final transient ConcurrentMap<String/*ConsumerID*/, 
SubscriptionFilterHandler>
+            subscriptionFilterData = new ConcurrentHashMap<>(256);
+
     private transient BrokerController brokerController;
     private transient BloomFilter bloomFilter;
 
@@ -113,24 +118,22 @@ public class ConsumerFilterManager extends ConfigManager {
         }
 
         // make illegal topic dead.
-        Collection<ConsumerFilterData> groupFilterData = 
getByGroup(consumerGroup);
-
-        Iterator<ConsumerFilterData> iterator = groupFilterData.iterator();
-        while (iterator.hasNext()) {
-            ConsumerFilterData filterData = iterator.next();
+        Set<String> curSubList = new HashSet<>();
+        for (SubscriptionData subscriptionData : subList) {
+            curSubList.add(subscriptionData.getTopic());
+        }
 
-            boolean exist = false;
-            for (SubscriptionData subscriptionData : subList) {
-                if (subscriptionData.getTopic().equals(filterData.getTopic())) 
{
-                    exist = true;
-                    break;
+        SubscriptionFilterHandler subscriptionFilterHandler = 
this.subscriptionFilterData.get(consumerGroup);
+        if (null != subscriptionFilterHandler) {
+            for (Map.Entry<String, ConsumerFilterData> entry : 
subscriptionFilterHandler.getTopicSqlFilterData().entrySet()) {
+                if (!curSubList.contains(entry.getKey())) {
+                    ConsumerFilterData filterData = entry.getValue();
+                    if (filterData != null) {
+                        filterData.setDeadTime(System.currentTimeMillis());
+                        log.info("Consumer filter changed: {}, make illegal 
topic dead:{}", consumerGroup, filterData);
+                    }
                 }
             }
-
-            if (!exist && !filterData.isDead()) {
-                filterData.setDeadTime(System.currentTimeMillis());
-                log.info("Consumer filter changed: {}, make illegal topic 
dead:{}", consumerGroup, filterData);
-            }
         }
     }
 
@@ -144,34 +147,48 @@ public class ConsumerFilterManager extends ConfigManager {
             return false;
         }
 
-        FilterDataMapByTopic filterDataMapByTopic = 
this.filterDataByTopic.get(topic);
-
-        if (filterDataMapByTopic == null) {
-            FilterDataMapByTopic temp = new FilterDataMapByTopic(topic);
-            FilterDataMapByTopic prev = 
this.filterDataByTopic.putIfAbsent(topic, temp);
-            filterDataMapByTopic = prev != null ? prev : temp;
+        if (null != this.brokerController) {
+            TopicConfig topicConfig = 
this.brokerController.getTopicConfigManager().selectTopicConfig(topic);
+            if (null == topicConfig) {
+                return false;
+            }
         }
 
-        BloomFilterData bloomFilterData = bloomFilter.generate(consumerGroup + 
"#" + topic);
+        SubscriptionFilterHandler subscriptionFilterHandler = 
this.subscriptionFilterData.get(consumerGroup);
+        if (subscriptionFilterHandler == null) {
+            SubscriptionFilterHandler temp = new 
SubscriptionFilterHandler(consumerGroup);
+            SubscriptionFilterHandler prev = 
this.subscriptionFilterData.putIfAbsent(consumerGroup, temp);
+            subscriptionFilterHandler = prev != null ? prev : temp;
+        }
 
-        return filterDataMapByTopic.register(consumerGroup, expression, type, 
bloomFilterData, clientVersion);
+        BloomFilterData bloomFilterData = null;
+        if (this.brokerController == null
+                || 
this.brokerController.getBrokerConfig().isEnableCalcFilterBitMap()) {
+            bloomFilterData = bloomFilter.generate(consumerGroup + "#" + 
topic);
+        }
+        ConsumerFilterData consumerFilterData = 
subscriptionFilterHandler.register(consumerGroup, expression, type, 
bloomFilterData, clientVersion, topic);
+        if (null == consumerFilterData) {
+            return false;
+        }
+        this.filterDataByTopic.putIfAbsent(topic, new 
FilterDataMapByTopic(topic));
+        this.filterDataByTopic.get(topic).put(consumerFilterData);
+        return true;
     }
 
     public void unRegister(final String consumerGroup) {
-        for (Entry<String, FilterDataMapByTopic> entry : 
filterDataByTopic.entrySet()) {
-            entry.getValue().unRegister(consumerGroup);
+        SubscriptionFilterHandler handler = 
this.subscriptionFilterData.get(consumerGroup);
+        if (handler != null) {
+            handler.unRegister();
         }
     }
 
     public ConsumerFilterData get(final String topic, final String 
consumerGroup) {
-        if (!this.filterDataByTopic.containsKey(topic)) {
-            return null;
-        }
-        if (this.filterDataByTopic.get(topic).getGroupFilterData().isEmpty()) {
+        SubscriptionFilterHandler handler = 
this.subscriptionFilterData.get(consumerGroup);
+        if (handler == null) {
             return null;
         }
 
-        return 
this.filterDataByTopic.get(topic).getGroupFilterData().get(consumerGroup);
+        return handler.getTopicSqlFilterData().get(topic);
     }
 
     public Collection<ConsumerFilterData> getByGroup(final String 
consumerGroup) {
@@ -196,14 +213,12 @@ public class ConsumerFilterManager extends ConfigManager {
     }
 
     public final Collection<ConsumerFilterData> get(final String topic) {
-        if (!this.filterDataByTopic.containsKey(topic)) {
-            return null;
-        }
-        if (this.filterDataByTopic.get(topic).getGroupFilterData().isEmpty()) {
+        FilterDataMapByTopic mapByTopic = this.filterDataByTopic.get(topic);
+        if (mapByTopic == null || mapByTopic.getGroupFilterData().isEmpty()) {
             return null;
         }
 
-        return this.filterDataByTopic.get(topic).getGroupFilterData().values();
+        return mapByTopic.getGroupFilterData().values();
     }
 
     public BloomFilter getBloomFilter() {
@@ -275,6 +290,19 @@ public class ConsumerFilterManager extends ConfigManager {
             if (!bloomChanged) {
                 this.filterDataByTopic = load.filterDataByTopic;
             }
+
+            // rebuild subscriptionFilterData from filterDataByTopic
+            for (Entry<String, FilterDataMapByTopic> entry : 
this.filterDataByTopic.entrySet()) {
+                for (Entry<String, ConsumerFilterData> groupEntry : 
entry.getValue().getGroupFilterData().entrySet()) {
+                    ConsumerFilterData data = groupEntry.getValue();
+                    if (data == null) {
+                        continue;
+                    }
+                    SubscriptionFilterHandler handler = 
this.subscriptionFilterData
+                        .computeIfAbsent(data.getConsumerGroup(), 
SubscriptionFilterHandler::new);
+                    handler.getTopicSqlFilterData().put(data.getTopic(), data);
+                }
+            }
         }
     }
 
@@ -288,26 +316,34 @@ public class ConsumerFilterManager extends ConfigManager {
     }
 
     public void clean() {
-        Iterator<Map.Entry<String, FilterDataMapByTopic>> topicIterator = 
this.filterDataByTopic.entrySet().iterator();
-        while (topicIterator.hasNext()) {
-            Map.Entry<String, FilterDataMapByTopic> filterDataMapByTopic = 
topicIterator.next();
+        Iterator<Map.Entry<String, SubscriptionFilterHandler>> 
consumerIterator = this.subscriptionFilterData.entrySet().iterator();
+        while (consumerIterator.hasNext()) {
+            Map.Entry<String, SubscriptionFilterHandler> 
subscriptionFilterHandlerEntry = consumerIterator.next();
 
             Iterator<Map.Entry<String, ConsumerFilterData>> filterDataIterator
-                = 
filterDataMapByTopic.getValue().getGroupFilterData().entrySet().iterator();
+                    = 
subscriptionFilterHandlerEntry.getValue().getTopicSqlFilterData().entrySet().iterator();
 
             while (filterDataIterator.hasNext()) {
                 Map.Entry<String, ConsumerFilterData> filterDataByGroup = 
filterDataIterator.next();
 
                 ConsumerFilterData filterData = filterDataByGroup.getValue();
                 if (filterData.howLongAfterDeath() >= (this.brokerController 
== null ? MS_24_HOUR : 
this.brokerController.getBrokerConfig().getFilterDataCleanTimeSpan())) {
-                    log.info("Remove filter consumer {}, died too long!", 
filterDataByGroup.getValue());
+                    log.info("Remove filter consumer {}, died too long!", 
filterDataByGroup.getKey());
                     filterDataIterator.remove();
+
+                    FilterDataMapByTopic mapByTopic = 
this.filterDataByTopic.get(filterData.getTopic());
+                    if (mapByTopic != null) {
+                        log.info("Remove filter data {} {} from 
filterDataByTopic", filterData.getTopic(), filterData.getConsumerGroup());
+                        
mapByTopic.getGroupFilterData().remove(filterData.getConsumerGroup());
+                        if (mapByTopic.getGroupFilterData().isEmpty()) {
+                            
this.filterDataByTopic.remove(filterData.getTopic());
+                        }
+                    }
                 }
             }
-
-            if 
(filterDataMapByTopic.getValue().getGroupFilterData().isEmpty()) {
-                log.info("Topic has no consumer, remove it! {}", 
filterDataMapByTopic.getKey());
-                topicIterator.remove();
+            if 
(subscriptionFilterHandlerEntry.getValue().getTopicSqlFilterData().isEmpty()) {
+                log.info("subscriptionFilterData Remove filter consumer {}", 
subscriptionFilterHandlerEntry.getKey());
+                consumerIterator.remove();
             }
         }
     }
@@ -334,75 +370,101 @@ public class ConsumerFilterManager extends ConfigManager 
{
             this.topic = topic;
         }
 
-        public void unRegister(String consumerGroup) {
-            if (!this.groupFilterData.containsKey(consumerGroup)) {
-                return;
+        public void put(ConsumerFilterData consumerFilterData) {
+            if (null != consumerFilterData) {
+                
this.groupFilterData.put(consumerFilterData.getConsumerGroup(), 
consumerFilterData);
             }
+        }
+
+        public final ConsumerFilterData get(String consumerGroup) {
+            return this.groupFilterData.get(consumerGroup);
+        }
 
-            ConsumerFilterData data = this.groupFilterData.get(consumerGroup);
+        public final ConcurrentMap<String, ConsumerFilterData> 
getGroupFilterData() {
+            return this.groupFilterData;
+        }
 
-            if (data == null || data.isDead()) {
-                return;
-            }
+        public void setGroupFilterData(final ConcurrentHashMap<String, 
ConsumerFilterData> groupFilterData) {
+            this.groupFilterData = groupFilterData;
+        }
+
+        public String getTopic() {
+            return topic;
+        }
 
-            long now = System.currentTimeMillis();
+        public void setTopic(final String topic) {
+            this.topic = topic;
+        }
+    }
 
-            log.info("Unregister consumer filter: {}, deadTime: {}", data, 
now);
 
-            data.setDeadTime(now);
+    public static class SubscriptionFilterHandler {
+
+        private Map<String/*Topic*/, ConsumerFilterData> topicSqlFilterData = 
new ConcurrentHashMap<>();
+
+        final private String consumerId;
+
+        public SubscriptionFilterHandler(String consumerId) {
+            this.consumerId = consumerId;
         }
 
-        public boolean register(String consumerGroup, String expression, 
String type, BloomFilterData bloomFilterData,
-            long clientVersion) {
-            ConsumerFilterData old = this.groupFilterData.get(consumerGroup);
+        public void unRegister() {
+            for (ConsumerFilterData data : topicSqlFilterData.values()) {
+                if (data != null && !data.isDead()) {
+                    long now = System.currentTimeMillis();
+                    log.info("Unregister consumer filter: {}, deadTime: {}", 
data, now);
+                    data.setDeadTime(now);
+                }
+            }
+        }
 
+        public ConsumerFilterData register(String consumerGroup, String 
expression, String type, BloomFilterData bloomFilterData,
+                                           long clientVersion, String topic) {
+            ConsumerFilterData old = this.topicSqlFilterData.get(topic);
             if (old == null) {
                 ConsumerFilterData consumerFilterData = build(topic, 
consumerGroup, expression, type, clientVersion);
                 if (consumerFilterData == null) {
-                    return false;
+                    return null;
                 }
                 consumerFilterData.setBloomFilterData(bloomFilterData);
-
-                old = this.groupFilterData.putIfAbsent(consumerGroup, 
consumerFilterData);
+                old = this.topicSqlFilterData.putIfAbsent(topic, 
consumerFilterData);
                 if (old == null) {
                     log.info("New consumer filter registered: {}", 
consumerFilterData);
-                    return true;
+                    return consumerFilterData;
                 } else {
                     if (clientVersion <= old.getClientVersion()) {
                         if (!type.equals(old.getExpressionType()) || 
!expression.equals(old.getExpression())) {
                             log.warn("Ignore consumer({} : {}) 
filter(concurrent), because of version {} <= {}, but maybe info 
changed!old={}:{}, ignored={}:{}",
-                                consumerGroup, topic,
-                                clientVersion, old.getClientVersion(),
-                                old.getExpressionType(), old.getExpression(),
-                                type, expression);
+                                    consumerGroup, topic,
+                                    clientVersion, old.getClientVersion(),
+                                    old.getExpressionType(), 
old.getExpression(),
+                                    type, expression);
                         }
                         if (clientVersion == old.getClientVersion() && 
old.isDead()) {
                             reAlive(old);
-                            return true;
+                            return old;
                         }
-
-                        return false;
+                        return null;
                     } else {
-                        this.groupFilterData.put(consumerGroup, 
consumerFilterData);
+                        this.topicSqlFilterData.put(topic, consumerFilterData);
                         log.info("New consumer filter registered(concurrent): 
{}, old: {}", consumerFilterData, old);
-                        return true;
+                        return consumerFilterData;
                     }
                 }
             } else {
                 if (clientVersion <= old.getClientVersion()) {
                     if (!type.equals(old.getExpressionType()) || 
!expression.equals(old.getExpression())) {
                         log.info("Ignore consumer({}:{}) filter, because of 
version {} <= {}, but maybe info changed!old={}:{}, ignored={}:{}",
-                            consumerGroup, topic,
-                            clientVersion, old.getClientVersion(),
-                            old.getExpressionType(), old.getExpression(),
-                            type, expression);
+                                consumerGroup, topic,
+                                clientVersion, old.getClientVersion(),
+                                old.getExpressionType(), old.getExpression(),
+                                type, expression);
                     }
                     if (clientVersion == old.getClientVersion() && 
old.isDead()) {
                         reAlive(old);
-                        return true;
+                        return old;
                     }
-
-                    return false;
+                    return null;
                 }
 
                 boolean change = !old.getExpression().equals(expression) || 
!old.getExpressionType().equals(type);
@@ -418,23 +480,19 @@ public class ConsumerFilterManager extends ConfigManager {
                     ConsumerFilterData consumerFilterData = build(topic, 
consumerGroup, expression, type, clientVersion);
                     if (consumerFilterData == null) {
                         // new expression compile error, remove old, let 
client report error.
-                        this.groupFilterData.remove(consumerGroup);
-                        return false;
+                        this.topicSqlFilterData.remove(topic);
+                        return null;
                     }
                     consumerFilterData.setBloomFilterData(bloomFilterData);
-
-                    this.groupFilterData.put(consumerGroup, 
consumerFilterData);
-
-                    log.info("Consumer filter info change, old: {}, new: {}, 
change: {}",
-                        old, consumerFilterData, change);
-
-                    return true;
+                    this.topicSqlFilterData.put(topic, consumerFilterData);
+                    log.info("Consumer filter info change, old: {}, new: {}, 
change: true", old, consumerFilterData);
+                    return consumerFilterData;
                 } else {
                     old.setClientVersion(clientVersion);
                     if (old.isDead()) {
                         reAlive(old);
                     }
-                    return true;
+                    return old;
                 }
             }
         }
@@ -445,24 +503,17 @@ public class ConsumerFilterManager extends ConfigManager {
             log.info("Re alive consumer filter: {}, oldDeadTime: {}", 
filterData, oldDeadTime);
         }
 
-        public final ConsumerFilterData get(String consumerGroup) {
-            return this.groupFilterData.get(consumerGroup);
-        }
-
-        public final ConcurrentMap<String, ConsumerFilterData> 
getGroupFilterData() {
-            return this.groupFilterData;
-        }
-
-        public void setGroupFilterData(final ConcurrentHashMap<String, 
ConsumerFilterData> groupFilterData) {
-            this.groupFilterData = groupFilterData;
+        public Map<String, ConsumerFilterData> getTopicSqlFilterData() {
+            return topicSqlFilterData;
         }
 
-        public String getTopic() {
-            return topic;
+        public void setTopicSqlFilterData(Map<String, ConsumerFilterData> 
topicSqlFilterData) {
+            this.topicSqlFilterData = topicSqlFilterData;
         }
 
-        public void setTopic(final String topic) {
-            this.topic = topic;
+        public String getConsumerId() {
+            return consumerId;
         }
     }
+
 }
diff --git 
a/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
 
b/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
index c01d8299dc..d7cd26bd8a 100644
--- 
a/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
+++ 
b/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
@@ -177,7 +177,9 @@ public class ConsumerFilterManagerTest {
             ConsumerFilterData filterData = iterator.next();
 
             assertThat(filterData).isNotNull();
-            
assertThat(filterManager.getBloomFilter().isValid(filterData.getBloomFilterData())).isTrue();
+            if (null != filterData.getBloomFilterData()) {
+                
assertThat(filterManager.getBloomFilter().isValid(filterData.getBloomFilterData())).isTrue();
+            }
         }
     }
 
@@ -187,9 +189,9 @@ public class ConsumerFilterManagerTest {
 
         assertThat(filterManager.register("topic0", "CID_0", "*", null, 
System.currentTimeMillis())).isFalse();
 
-        Collection<ConsumerFilterData> filterDatas = 
filterManager.getByGroup("CID_0");
+        ConsumerFilterData filterDatas = filterManager.get("topic0", "CID_0");
 
-        assertThat(filterDatas).isNullOrEmpty();
+        assertThat(filterDatas).isNull();
     }
 
     @Test
@@ -269,4 +271,41 @@ public class ConsumerFilterManagerTest {
         }
     }
 
+    @Test
+    public void testRegister_bySubscriptionData_shrinkMakesDead() {
+        ConsumerFilterManager filterManager = new ConsumerFilterManager();
+        List<SubscriptionData> subscriptionDatas = new ArrayList<>();
+        for (int i = 0; i < 3; i++) {
+            try {
+                subscriptionDatas.add(
+                    FilterAPI.build("topic" + i, "a is not null and a > " + i, 
ExpressionType.SQL92)
+                );
+            } catch (Exception e) {
+                assertThat(true).isFalse();
+            }
+        }
+
+        filterManager.register("CID_0", subscriptionDatas);
+
+        ConsumerFilterData topic2Data = filterManager.get("topic2", "CID_0");
+        assertThat(topic2Data).isNotNull();
+        assertThat(topic2Data.isDead()).isFalse();
+
+        // shrink: remove topic2 from subscription list
+        List<SubscriptionData> shrunkList = new 
ArrayList<>(subscriptionDatas.subList(0, 2));
+        filterManager.register("CID_0", shrunkList);
+
+        // topic2 should be marked dead
+        assertThat(topic2Data.isDead()).isTrue();
+
+        // topic0 and topic1 should still be alive
+        ConsumerFilterData topic0Data = filterManager.get("topic0", "CID_0");
+        assertThat(topic0Data).isNotNull();
+        assertThat(topic0Data.isDead()).isFalse();
+
+        ConsumerFilterData topic1Data = filterManager.get("topic1", "CID_0");
+        assertThat(topic1Data).isNotNull();
+        assertThat(topic1Data.isDead()).isFalse();
+    }
+
 }

Reply via email to