From 3ca7df6fc7e5dc7701cff379ffaf2ecc81dadea4 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 18 Mar 2025 21:16:38 +0800 Subject: [PATCH 1/2] [fix][broker] Restore the behavior to dispatch batch messages according to consumer permits --- .../persistent/PersistentDispatcherMultipleConsumers.java | 2 +- .../org/apache/pulsar/broker/service/BatchMessageTest.java | 4 ++++ 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 6f3fe19f0a104..7ef39724cf27d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -850,7 +850,7 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis int maxAdditionalUnackedMessages = Math.max(c.getMaxUnackedMessages() - c.getUnackedMessages(), 0); maxMessagesInThisBatch = Math.min(maxMessagesInThisBatch, maxAdditionalUnackedMessages); } - int maxEntriesInThisBatch = Math.min(availablePermits, + int maxEntriesInThisBatch = Math.min(availablePermits / avgBatchSizePerMsg, // use the average batch size per message to calculate the number of entries to // dispatch. round up to the next integer without using floating point arithmetic. (maxMessagesInThisBatch + avgBatchSizePerMsg - 1) / avgBatchSizePerMsg); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java index e5f9e43b8bb4a..b821c0cc66351 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java @@ -1012,6 +1012,10 @@ public void testBatchMessageDispatchingAccordingToPermits() throws Exception { } FutureUtil.waitForAll(sendFutureList).get(); + Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { + assertTrue(consumer1.numMessagesInQueue() > 0); + assertTrue(consumer2.numMessagesInQueue() > 0); + }); assertEquals(consumer1.numMessagesInQueue(), batchMessages, batchMessages); assertEquals(consumer2.numMessagesInQueue(), batchMessages, batchMessages); From 0a960a2d733e8bfb717770a608c3d7628b7a7b59 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 19 Mar 2025 10:24:04 +0800 Subject: [PATCH 2/2] Fix the maxMessagesInThisBatch --- .../persistent/PersistentDispatcherMultipleConsumers.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 7ef39724cf27d..2af04044aae45 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -843,14 +843,14 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis c, c.getAvailablePermits()); } - int maxMessagesInThisBatch = - Math.max(remainingMessages, serviceConfig.getDispatcherMaxRoundRobinBatchSize()); + int maxMessagesInThisBatch = Math.min(remainingMessages, availablePermits); if (c.getMaxUnackedMessages() > 0) { // Calculate the maximum number of additional unacked messages allowed int maxAdditionalUnackedMessages = Math.max(c.getMaxUnackedMessages() - c.getUnackedMessages(), 0); maxMessagesInThisBatch = Math.min(maxMessagesInThisBatch, maxAdditionalUnackedMessages); } - int maxEntriesInThisBatch = Math.min(availablePermits / avgBatchSizePerMsg, + // TODO: add tests to verify dispatcherMaxRoundRobinBatchSize is respected + int maxEntriesInThisBatch = Math.min(serviceConfig.getDispatcherMaxRoundRobinBatchSize(), // use the average batch size per message to calculate the number of entries to // dispatch. round up to the next integer without using floating point arithmetic. (maxMessagesInThisBatch + avgBatchSizePerMsg - 1) / avgBatchSizePerMsg);