diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java index d72807a920f2b..e1878ee1491b2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/BatchMessageIndexAckTest.java @@ -18,6 +18,16 @@ */ package org.apache.pulsar.client.impl; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doReturn; +import static org.testng.Assert.assertEquals; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.NavigableMap; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.client.api.DigestType; @@ -32,24 +42,15 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; +import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.common.util.FutureUtil; +import org.awaitility.Awaitility; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.DataProvider; import org.testng.annotations.Test; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.NavigableMap; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.mockito.ArgumentMatchers.anyLong; -import static org.mockito.Mockito.doReturn; - @Slf4j @Test(groups = "broker-impl") public class BatchMessageIndexAckTest extends ProducerConsumerBase { @@ -345,4 +346,45 @@ public void testDoNotRecycleAckSetMultipleTimes() throws Exception { producer.close(); consumer.close(); } + + @Test + public void testAcknowledgeCumulative() throws Exception { + final String topic = "persistent://my-property/my-ns/testAcknowledgeCumulative"; + int messageNumber= 10; + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("sub") + .enableBatchIndexAcknowledgment(true) + .subscribe(); + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .batchingMaxMessages(messageNumber) + .blockIfQueueFull(true) + .create(); + + for (int i = 0; i < 10; i++) { + producer.sendAsync("message-" + i); + } + producer.flush(); + + int count = 0; + while (true) { + Message message = consumer.receive(5, TimeUnit.SECONDS); + if (message == null) { + break; + } + consumer.acknowledgeCumulative(message.getMessageId()); + count++; + } + + assertEquals(count, messageNumber); + Awaitility.await().untilAsserted(() -> { + TopicStats stats = admin.topics().getStats(topic, true, true); + assertEquals(stats.getSubscriptions().get("sub").getMsgBacklog(), 0); + assertEquals(stats.getSubscriptions().get("sub").getBacklogSize(), 0); + }); + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java index b44670c14631f..cc15bbd55b8d2 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PersistentAcknowledgmentsGroupingTracker.java @@ -622,7 +622,18 @@ protected LastCumulativeAck initialValue() { private boolean flushRequired = false; public synchronized void update(final MessageIdImpl messageId, final BitSetRecyclable bitSetRecyclable) { - if (messageId.compareTo(this.messageId) > 0) { + MessageIdImpl newMessageId = messageId; + MessageIdImpl lastMessageId = this.messageId; + if (newMessageId instanceof BatchMessageIdImpl && !(lastMessageId instanceof BatchMessageIdImpl)) { + lastMessageId = + new BatchMessageIdImpl(lastMessageId.ledgerId, lastMessageId.entryId, lastMessageId.partitionIndex, + Integer.MAX_VALUE); + } else if (!(newMessageId instanceof BatchMessageIdImpl) && (lastMessageId instanceof BatchMessageIdImpl)) { + newMessageId = + new BatchMessageIdImpl(newMessageId.ledgerId, newMessageId.entryId, newMessageId.partitionIndex, + Integer.MAX_VALUE); + } + if (newMessageId.compareTo(lastMessageId) > 0) { if (this.bitSetRecyclable != null && this.bitSetRecyclable != bitSetRecyclable) { this.bitSetRecyclable.recycle(); }