diff --git a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java index b487b8757f4..119c5952b0c 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java @@ -97,6 +97,9 @@ public void addPartialSubscription(String clientId, String group, String topic, if (!liteLifecycleManager.isSubscriptionActive(topic, lmqName)) { continue; } + if (!isLiteTopicSubscribed(clientGroup, lmqName) && getActiveSubscriptionNum() >= maxCount) { + throw new LiteQuotaException("lite subscription quota exceeded " + maxCount); + } thisSub.addLiteTopic(lmqName); // First remove the old subscription if (LiteMetadataUtil.isSubLiteExclusive(group, brokerController)) { @@ -147,6 +150,10 @@ public void addCompleteSubscription(String clientId, String group, String topic, removeTopicGroup(clientGroup, lmqName, false); }); lmqNameNew.forEach(lmqName -> { + long maxCount = brokerController.getBrokerConfig().getMaxLiteSubscriptionCount(); + if (!isLiteTopicSubscribed(clientGroup, lmqName) && getActiveSubscriptionNum() >= maxCount) { + throw new LiteQuotaException("lite subscription quota exceeded " + maxCount); + } thisSub.addLiteTopic(lmqName); addTopicGroup(clientGroup, lmqName); }); @@ -269,6 +276,11 @@ public void cleanSubscription(String lmqName, boolean notifyClient) { } } + protected boolean isLiteTopicSubscribed(ClientGroup clientGroup, String lmqName) { + Set topicGroupSet = liteTopic2Group.get(lmqName); + return topicGroupSet != null && topicGroupSet.contains(clientGroup); + } + protected void addTopicGroup(ClientGroup clientGroup, String lmqName) { Set topicGroupSet = liteTopic2Group .computeIfAbsent(lmqName, k -> ConcurrentHashMap.newKeySet());