diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 8c0e1d5e42dc9..c447809f34c17 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -258,8 +258,11 @@ public ChannelPromise sendMessages(final List entries, EntryBatchSizes ba consumerId, entry.getLedgerId(), entry.getEntryId()); } - int redeliveryCount = redeliveryTracker - .getRedeliveryCount(PositionImpl.get(messageId.getLedgerId(), messageId.getEntryId())); + int redeliveryCount = 0; + PositionImpl position = PositionImpl.get(messageId.getLedgerId(), messageId.getEntryId()); + if (redeliveryTracker.contains(position)) { + redeliveryCount = redeliveryTracker.incrementAndGetRedeliveryCount(position); + } ctx.write(Commands.newMessage(consumerId, messageId, redeliveryCount, metadataAndPayload), ctx.voidPromise()); messageId.recycle(); messageIdBuilder.recycle(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/InMemoryRedeliveryTracker.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/InMemoryRedeliveryTracker.java index b4e3508039327..a5dab6d941309 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/InMemoryRedeliveryTracker.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/InMemoryRedeliveryTracker.java @@ -62,4 +62,16 @@ public void removeBatch(List positions) { public void clear() { trackerCache.clear(); } + + @Override + public boolean contains(Position position) { + PositionImpl positionImpl = (PositionImpl) position; + return trackerCache.containsKey(positionImpl.getLedgerId(), positionImpl.getEntryId()); + } + + @Override + public void addIfAbsent(Position position) { + PositionImpl positionImpl = (PositionImpl) position; + trackerCache.putIfAbsent(positionImpl.getLedgerId(), positionImpl.getEntryId(), 0, 0L); + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTracker.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTracker.java index 79ea510eb9703..dc5fdb32fb36a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTracker.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTracker.java @@ -33,4 +33,8 @@ public interface RedeliveryTracker { void removeBatch(List positions); void clear(); + + boolean contains(Position position); + + void addIfAbsent(Position position); } \ No newline at end of file diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTrackerDisabled.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTrackerDisabled.java index a930cd70b1ba1..ef422c927236f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTrackerDisabled.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/RedeliveryTrackerDisabled.java @@ -52,4 +52,14 @@ public void removeBatch(List positions) { public void clear() { // no-op } + + @Override + public boolean contains(Position position) { + return false; + } + + @Override + public void addIfAbsent(Position position) { + // no-op + } } 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 5d536c113e0c3..542685159233d 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 @@ -213,6 +213,7 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE } messagesToRedeliver.clear(); + redeliveryTracker.clear(); if (closeFuture != null) { log.info("[{}] All consumers removed. Subscription is disconnected", name); closeFuture.complete(null); @@ -224,6 +225,7 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE } consumer.getPendingAcks().forEach((ledgerId, entryId, batchSize, none) -> { messagesToRedeliver.add(ledgerId, entryId); + redeliveryTracker.addIfAbsent(PositionImpl.get(ledgerId, entryId)); }); totalAvailablePermits -= consumer.getAvailablePermits(); readMoreEntries(); @@ -637,7 +639,7 @@ public synchronized void redeliverUnacknowledgedMessages(Consumer consumer) { public synchronized void redeliverUnacknowledgedMessages(Consumer consumer, List positions) { positions.forEach(position -> { messagesToRedeliver.add(position.getLedgerId(), position.getEntryId()); - redeliveryTracker.incrementAndGetRedeliveryCount(position); + redeliveryTracker.addIfAbsent(position); }); if (log.isDebugEnabled()) { log.debug("[{}-{}] Redelivering unacknowledged messages for consumer {}", name, consumer, positions); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 65c4b9884d2e4..084d0b41208b8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -336,7 +336,7 @@ private synchronized void internalRedeliverUnacknowledgedMessages(Consumer consu @Override public void redeliverUnacknowledgedMessages(Consumer consumer, List positions) { // We cannot redeliver single messages to single consumers to preserve ordering. - positions.forEach(redeliveryTracker::incrementAndGetRedeliveryCount); + positions.forEach(redeliveryTracker::addIfAbsent); redeliverUnacknowledgedMessages(consumer); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java index ca2eadbdf2217..b1f1ddd6fd849 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/DeadLetterTopicTest.java @@ -58,7 +58,7 @@ public void testDeadLetterTopic() throws Exception { .topic(topic) .subscriptionName("my-subscription") .subscriptionType(SubscriptionType.Shared) - .ackTimeout(3, TimeUnit.SECONDS) + .ackTimeout(1, TimeUnit.SECONDS) .deadLetterPolicy(DeadLetterPolicy.builder().maxRedeliverCount(maxRedeliveryCount).build()) .receiverQueueSize(100) .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) @@ -134,7 +134,7 @@ public void testDeadLetterTopicWithMultiTopic() throws Exception { .topic(topic1, topic2) .subscriptionName("my-subscription") .subscriptionType(SubscriptionType.Shared) - .ackTimeout(3, TimeUnit.SECONDS) + .ackTimeout(1, TimeUnit.SECONDS) .deadLetterPolicy(DeadLetterPolicy.builder().maxRedeliverCount(maxRedeliveryCount).build()) .receiverQueueSize(100) .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) @@ -210,7 +210,7 @@ public void testDeadLetterTopicByCustomTopicName() throws Exception { .topic(topic) .subscriptionName("my-subscription") .subscriptionType(SubscriptionType.Shared) - .ackTimeout(3, TimeUnit.SECONDS) + .ackTimeout(1, TimeUnit.SECONDS) .receiverQueueSize(100) .deadLetterPolicy(DeadLetterPolicy.builder() .maxRedeliverCount(maxRedeliveryCount) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ExposeMessageRedeliveryCountTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ExposeMessageRedeliveryCountTest.java index 5156c900f815e..8ebbca0be23d6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ExposeMessageRedeliveryCountTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ExposeMessageRedeliveryCountTest.java @@ -24,6 +24,8 @@ import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.TimeUnit; public class ExposeMessageRedeliveryCountTest extends ProducerConsumerBase { @@ -50,7 +52,7 @@ public void testRedeliveryCount() throws PulsarClientException { .topic(topic) .subscriptionName("my-subscription") .subscriptionType(SubscriptionType.Shared) - .ackTimeout(3, TimeUnit.SECONDS) + .ackTimeout(1, TimeUnit.SECONDS) .receiverQueueSize(100) .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) .subscribe(); @@ -87,7 +89,7 @@ public void testRedeliveryCountWithPartitionedTopic() throws PulsarClientExcepti .topic(topic) .subscriptionName("my-subscription") .subscriptionType(SubscriptionType.Shared) - .ackTimeout(3, TimeUnit.SECONDS) + .ackTimeout(1, TimeUnit.SECONDS) .receiverQueueSize(100) .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) .subscribe(); @@ -114,4 +116,63 @@ public void testRedeliveryCountWithPartitionedTopic() throws PulsarClientExcepti admin.topics().deletePartitionedTopic(topic); } + + @Test(timeOut = 30000) + public void testRedeliveryCountWhenConsumerDisconnected() throws PulsarClientException, InterruptedException { + + String topic = "persistent://my-property/my-ns/testRedeliveryCountWhenConsumerDisconnected"; + + Consumer consumer0 = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("s1") + .subscriptionType(SubscriptionType.Shared) + .subscribe(); + + Consumer consumer1 = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("s1") + .subscriptionType(SubscriptionType.Shared) + .subscribe(); + + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .enableBatching(true) + .batchingMaxMessages(5) + .batchingMaxPublishDelay(1, TimeUnit.SECONDS) + .create(); + + final int messages = 10; + for (int i = 0; i < messages; i++) { + producer.send("my-message-" + i); + } + + List> receivedMessagesForConsumer0 = new ArrayList<>(); + List> receivedMessagesForConsumer1 = new ArrayList<>(); + + for (int i = 0; i < messages; i++) { + Message msg = consumer0.receive(1, TimeUnit.SECONDS); + if (msg != null) { + receivedMessagesForConsumer0.add(msg); + } else { + break; + } + } + + for (int i = 0; i < messages; i++) { + Message msg = consumer1.receive(1, TimeUnit.SECONDS); + if (msg != null) { + receivedMessagesForConsumer1.add(msg); + } else { + break; + } } + + Assert.assertEquals(receivedMessagesForConsumer0.size() + receivedMessagesForConsumer1.size(), messages); + + consumer0.close(); + + for (int i = 0; i < receivedMessagesForConsumer0.size(); i++) { + Assert.assertEquals(consumer1.receive().getRedeliveryCount(), 1); + } + + } }