diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index 0a402d8322abe..60762d8400a73 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -236,7 +236,11 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { Consumer consumer = current.getKey(); List entriesWithSameKey = current.getValue(); int entriesWithSameKeyCount = entriesWithSameKey.size(); - final int availablePermits = consumer == null ? 0 : Math.max(consumer.getAvailablePermits(), 0); + int availablePermits = consumer == null ? 0 : Math.max(consumer.getAvailablePermits(), 0); + if (consumer != null && consumer.getMaxUnackedMessages() > 0) { + availablePermits = Math.min(availablePermits, + consumer.getMaxUnackedMessages() - consumer.getUnackedMessages()); + } int maxMessagesForC = Math.min(entriesWithSameKeyCount, availablePermits); int messagesForC = getRestrictedMaxEntriesForConsumer(consumer, entriesWithSameKey, maxMessagesForC, readType, consumerStickyKeyHashesMap.get(consumer)); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 9af19cc672d6b..2f572a841b09a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -149,6 +149,16 @@ public Object[][] ackReceiptEnabled() { return new Object[][] { { true }, { false } }; } + @DataProvider(name = "ackReceiptEnabledAndSubscriptionTypes") + public Object[][] ackReceiptEnabledAndSubscriptionTypes() { + return new Object[][] { + {true, SubscriptionType.Shared}, + {true, SubscriptionType.Key_Shared}, + {false, SubscriptionType.Shared}, + {false, SubscriptionType.Key_Shared}, + }; + } + @AfterMethod(alwaysRun = true) @Override protected void cleanup() throws Exception { @@ -1591,8 +1601,9 @@ public void testConsumerBlockingWithUnAckedMessagesMultipleIteration(boolean ack } } - @Test(dataProvider = "ackReceiptEnabled") - public void testMaxUnAckMessagesLowerThanPermits(boolean ackReceiptEnabled) throws PulsarClientException { + @Test(dataProvider = "ackReceiptEnabledAndSubscriptionTypes") + public void testMaxUnAckMessagesLowerThanPermits(boolean ackReceiptEnabled, SubscriptionType subType) + throws PulsarClientException { final int maxUnacks = 10; pulsar.getConfiguration().setMaxUnackedMessagesPerConsumer(maxUnacks); final String topic = "persistent://my-property/my-ns/testMaxUnAckMessagesLowerThanPermits"; @@ -1600,7 +1611,7 @@ public void testMaxUnAckMessagesLowerThanPermits(boolean ackReceiptEnabled) thro @Cleanup Consumer consumer = pulsarClient.newConsumer(Schema.STRING) .topic(topic).subscriptionName("sub") - .subscriptionType(SubscriptionType.Shared) + .subscriptionType(subType) .isAckReceiptEnabled(ackReceiptEnabled) .acknowledgmentGroupTime(0, TimeUnit.SECONDS) .subscribe();