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 3a48f96b987..5c5600d1552 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.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 filterDataByTopic = new ConcurrentHashMap<>(256); + private final transient ConcurrentMap + subscriptionFilterData = new ConcurrentHashMap<>(256); + private transient BrokerController brokerController; private transient BloomFilter bloomFilter; @@ -113,24 +118,22 @@ public void register(final String consumerGroup, final Collection groupFilterData = getByGroup(consumerGroup); - - Iterator iterator = groupFilterData.iterator(); - while (iterator.hasNext()) { - ConsumerFilterData filterData = iterator.next(); + Set 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 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 boolean register(final String topic, final String consumerGroup, final St 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 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 getByGroup(final String consumerGroup) { @@ -196,14 +213,12 @@ public Collection getByGroup(final String consumerGroup) { } public final Collection 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 void decode(final String jsonString) { if (!bloomChanged) { this.filterDataByTopic = load.filterDataByTopic; } + + // rebuild subscriptionFilterData from filterDataByTopic + for (Entry entry : this.filterDataByTopic.entrySet()) { + for (Entry 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 String encode(final boolean prettyFormat) { } public void clean() { - Iterator> topicIterator = this.filterDataByTopic.entrySet().iterator(); - while (topicIterator.hasNext()) { - Map.Entry filterDataMapByTopic = topicIterator.next(); + Iterator> consumerIterator = this.subscriptionFilterData.entrySet().iterator(); + while (consumerIterator.hasNext()) { + Map.Entry subscriptionFilterHandlerEntry = consumerIterator.next(); Iterator> filterDataIterator - = filterDataMapByTopic.getValue().getGroupFilterData().entrySet().iterator(); + = subscriptionFilterHandlerEntry.getValue().getTopicSqlFilterData().entrySet().iterator(); while (filterDataIterator.hasNext()) { Map.Entry 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 FilterDataMapByTopic(String topic) { 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 getGroupFilterData() { + return this.groupFilterData; + } - if (data == null || data.isDead()) { - return; - } + public void setGroupFilterData(final ConcurrentHashMap 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 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 boolean register(String consumerGroup, String expression, String type, Bl 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 @@ protected void reAlive(ConsumerFilterData filterData) { log.info("Re alive consumer filter: {}, oldDeadTime: {}", filterData, oldDeadTime); } - public final ConsumerFilterData get(String consumerGroup) { - return this.groupFilterData.get(consumerGroup); - } - - public final ConcurrentMap getGroupFilterData() { - return this.groupFilterData; - } - - public void setGroupFilterData(final ConcurrentHashMap groupFilterData) { - this.groupFilterData = groupFilterData; + public Map getTopicSqlFilterData() { + return topicSqlFilterData; } - public String getTopic() { - return topic; + public void setTopicSqlFilterData(Map 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 c01d8299dcf..d7cd26bd8a9 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 void testRegister_bySubscriptionData() { 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 void testRegister_tag() { assertThat(filterManager.register("topic0", "CID_0", "*", null, System.currentTimeMillis())).isFalse(); - Collection filterDatas = filterManager.getByGroup("CID_0"); + ConsumerFilterData filterDatas = filterManager.get("topic0", "CID_0"); - assertThat(filterDatas).isNullOrEmpty(); + assertThat(filterDatas).isNull(); } @Test @@ -269,4 +271,41 @@ public void testPersist_clean() { } } + @Test + public void testRegister_bySubscriptionData_shrinkMakesDead() { + ConsumerFilterManager filterManager = new ConsumerFilterManager(); + List 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 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(); + } + }