From 7d04ab0742474957952aadacce922905bb9e4b2c Mon Sep 17 00:00:00 2001 From: lipenghui Date: Thu, 19 Dec 2019 19:03:45 +0800 Subject: [PATCH 1/2] Fix wrong redelivery count while redeliver when consumer disconnected. --- .../pulsar/broker/service/Consumer.java | 7 ++- .../service/InMemoryRedeliveryTracker.java | 12 ++++ .../broker/service/RedeliveryTracker.java | 4 ++ .../service/RedeliveryTrackerDisabled.java | 10 +++ ...PersistentDispatcherMultipleConsumers.java | 4 +- ...sistentDispatcherSingleActiveConsumer.java | 2 +- .../api/ExposeMessageRedeliveryCountTest.java | 61 +++++++++++++++++++ 7 files changed, 96 insertions(+), 4 deletions(-) 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..7aecc5052f22b 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 add(Position position) { + PositionImpl positionImpl = (PositionImpl) position; + trackerCache.put(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..90cb674a3e701 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 add(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..3e7b4550fc51a 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 add(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..d36fc1f1b3158 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.add(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.add(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..5a19872404072 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::add); redeliverUnacknowledgedMessages(consumer); } 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..5c66d337b1488 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 { @@ -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); + } + + } } From e05600fcb3b91010c46c3e32f3170d7460bb2c1d Mon Sep 17 00:00:00 2001 From: lipenghui Date: Mon, 23 Dec 2019 18:00:30 +0800 Subject: [PATCH 2/2] Fix unit tests --- .../pulsar/broker/service/InMemoryRedeliveryTracker.java | 4 ++-- .../org/apache/pulsar/broker/service/RedeliveryTracker.java | 2 +- .../pulsar/broker/service/RedeliveryTrackerDisabled.java | 2 +- .../persistent/PersistentDispatcherMultipleConsumers.java | 4 ++-- .../PersistentDispatcherSingleActiveConsumer.java | 2 +- .../org/apache/pulsar/client/api/DeadLetterTopicTest.java | 6 +++--- .../pulsar/client/api/ExposeMessageRedeliveryCountTest.java | 4 ++-- 7 files changed, 12 insertions(+), 12 deletions(-) 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 7aecc5052f22b..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 @@ -70,8 +70,8 @@ public boolean contains(Position position) { } @Override - public void add(Position position) { + public void addIfAbsent(Position position) { PositionImpl positionImpl = (PositionImpl) position; - trackerCache.put(positionImpl.getLedgerId(), positionImpl.getEntryId(), 0, 0L); + 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 90cb674a3e701..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 @@ -36,5 +36,5 @@ public interface RedeliveryTracker { boolean contains(Position position); - void add(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 3e7b4550fc51a..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 @@ -59,7 +59,7 @@ public boolean contains(Position position) { } @Override - public void add(Position position) { + 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 d36fc1f1b3158..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 @@ -225,7 +225,7 @@ public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceE } consumer.getPendingAcks().forEach((ledgerId, entryId, batchSize, none) -> { messagesToRedeliver.add(ledgerId, entryId); - redeliveryTracker.add(PositionImpl.get(ledgerId, entryId)); + redeliveryTracker.addIfAbsent(PositionImpl.get(ledgerId, entryId)); }); totalAvailablePermits -= consumer.getAvailablePermits(); readMoreEntries(); @@ -639,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.add(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 5a19872404072..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::add); + 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 5c66d337b1488..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 @@ -52,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(); @@ -89,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();