From 2be9fca52bdbb333e5fdc9d58269ebdce5710f2a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 01:26:07 +0800 Subject: [PATCH 01/12] - --- .../MessageRedeliveryController.java | 8 +++ ...tStickyKeyDispatcherMultipleConsumers.java | 52 ++++++++++++++++++- 2 files changed, 59 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageRedeliveryController.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageRedeliveryController.java index 5bf3f5506fa81..6380317724207 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageRedeliveryController.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageRedeliveryController.java @@ -95,6 +95,14 @@ private void removeFromHashBlocker(long ledgerId, long entryId) { } } + public Long getHash(long ledgerId, long entryId) { + LongPair value = hashesToBeBlocked.get(ledgerId, entryId); + if (value == null) { + return null; + } + return value.first; + } + public void removeAllUpTo(long markDeleteLedgerId, long markDeleteEntryId) { if (!allowOutOfOrderDelivery) { List keysToRemove = new ArrayList<>(); 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 8f05530f58bfa..2f967e461cd64 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 @@ -30,6 +30,7 @@ import java.util.Map; import java.util.NavigableSet; import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.atomic.AtomicInteger; import org.apache.bookkeeper.mledger.Entry; @@ -165,6 +166,14 @@ protected Map> initialValue() throws Exception { } }; + private static final FastThreadLocal>> localGroupedPositions = + new FastThreadLocal>>() { + @Override + protected Map> initialValue() throws Exception { + return new HashMap<>(); + } + }; + @Override protected synchronized boolean trySendMessagesToConsumers(ReadType readType, List entries) { long totalMessagesSent = 0; @@ -433,8 +442,49 @@ protected synchronized NavigableSet getMessagesToReplayNow(int max this.isDispatcherStuckOnReplays = false; return Collections.emptyNavigableSet(); } else { - return super.getMessagesToReplayNow(maxMessagesToRead); + NavigableSet res = super.getMessagesToReplayNow(maxMessagesToRead); + return filterOutMessagesWillBeDiscarded(res); + } + } + + /** + * This method is in order to avoid the scenario below: + * - Read entries from the Replay queue. + * - The Key_Shared anti-ordering mechanism filtered out all of the entries. + * - Delivery non entry to the client, but we did a BK read. + */ + private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { + // Remove invalid items. + removeConsumersFromRecentJoinedConsumers(); + if (recentlyJoinedConsumers == null || src.isEmpty()) { + return src; + } + PositionImpl firstMaxReadPos = recentlyJoinedConsumers.values().iterator().next(); + // Group by key_hash. + NavigableSet res = new TreeSet<>(); + final Map> groupedEntries = localGroupedPositions.get(); + groupedEntries.clear(); + for (PositionImpl pos : src) { + Long stickyKeyHash = redeliveryMessages.getHash(pos.getLedgerId(), pos.getEntryId()); + if (stickyKeyHash == null) { + res.add(pos); + } + Consumer c = selector.select(stickyKeyHash.intValue()); + PositionImpl currentMaxReadPosition = recentlyJoinedConsumers.get(c); + if (c == null) { + continue; + } + if (currentMaxReadPosition == null) { + res.add(pos); + } else if (pos.compareTo(firstMaxReadPos) < 0 && pos.compareTo(currentMaxReadPosition) < 0) { + res.add(pos); + } else { + // The subsequent positions will also be larger than "firstMaxReadPos", why need to check continuous. + // Because the subsequent positions may be delivered to a consumer which does not in the + // "recentlyJoinedConsumers". + } } + return res; } @Override From feacf9bec09ea27330c7eafcb1f804e6135b328a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 02:35:35 +0800 Subject: [PATCH 02/12] add test --- ...tStickyKeyDispatcherMultipleConsumers.java | 33 ++++-- .../client/api/KeySharedSubscriptionTest.java | 108 ++++++++++++++++++ 2 files changed, 132 insertions(+), 9 deletions(-) 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 2f967e461cd64..a4fa5aea4c384 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 @@ -29,6 +29,7 @@ import java.util.List; import java.util.Map; import java.util.NavigableSet; +import java.util.NoSuchElementException; import java.util.Set; import java.util.TreeSet; import java.util.concurrent.CompletableFuture; @@ -38,6 +39,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.collections4.MapUtils; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.ConsistentHashingStickyKeyConsumerSelector; @@ -257,13 +259,7 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis assert consumer != null; // checked when added to groupedEntries List entriesWithSameKey = current.getValue(); int entriesWithSameKeyCount = entriesWithSameKey.size(); - int availablePermits = Math.max(consumer.getAvailablePermits(), 0); - if (consumer.getMaxUnackedMessages() > 0) { - int remainUnAckedMessages = - // Avoid negative number - Math.max(consumer.getMaxUnackedMessages() - consumer.getUnackedMessages(), 0); - availablePermits = Math.min(availablePermits, remainUnAckedMessages); - } + int availablePermits = getAvailablePermits(consumer); int maxMessagesForC = Math.min(entriesWithSameKeyCount, availablePermits); int messagesForC = getRestrictedMaxEntriesForConsumer(consumer, entriesWithSameKey, maxMessagesForC, readType, consumerStickyKeyHashesMap.get(consumer)); @@ -447,6 +443,16 @@ protected synchronized NavigableSet getMessagesToReplayNow(int max } } + private int getAvailablePermits(Consumer c) { + int availablePermits = Math.max(c.getAvailablePermits(), 0); + if (c.getMaxUnackedMessages() > 0) { + // Avoid negative number + int remainUnAckedMessages = Math.max(c.getMaxUnackedMessages() - c.getUnackedMessages(), 0); + availablePermits = Math.min(availablePermits, remainUnAckedMessages); + } + return availablePermits; + } + /** * This method is in order to avoid the scenario below: * - Read entries from the Replay queue. @@ -456,10 +462,16 @@ protected synchronized NavigableSet getMessagesToReplayNow(int max private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { // Remove invalid items. removeConsumersFromRecentJoinedConsumers(); - if (recentlyJoinedConsumers == null || src.isEmpty()) { + if (MapUtils.isEmpty(recentlyJoinedConsumers) || src.isEmpty()) { + return src; + } + PositionImpl firstMaxReadPos = null; + try { + firstMaxReadPos = recentlyJoinedConsumers.values().iterator().next(); + } catch (NoSuchElementException noSuchElementException) { + // Avoid error due to concurrent modifying. return src; } - PositionImpl firstMaxReadPos = recentlyJoinedConsumers.values().iterator().next(); // Group by key_hash. NavigableSet res = new TreeSet<>(); final Map> groupedEntries = localGroupedPositions.get(); @@ -470,6 +482,9 @@ private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet res.add(pos); } Consumer c = selector.select(stickyKeyHash.intValue()); + if (getAvailablePermits(c) == 0) { + continue; + } PositionImpl currentMaxReadPosition = recentlyJoinedConsumers.get(c); if (c == null) { continue; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 18fb141be3178..a7c7a72505b87 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -49,11 +49,15 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import lombok.Cleanup; +import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; +import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.nonpersistent.NonPersistentStickyKeyDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentStickyKeyDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.impl.ConsumerImpl; import org.apache.pulsar.common.api.proto.KeySharedMode; import org.apache.pulsar.common.naming.TopicDomain; @@ -61,6 +65,7 @@ import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.util.Murmur3_32Hash; import org.awaitility.Awaitility; +import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -1630,4 +1635,107 @@ public void testContinueDispatchMessagesWhenMessageDelayed() throws Exception { log.info("Got {} other messages...", sum); Assert.assertEquals(sum, delayedMessages + messages); } + + private AtomicInteger injectReplayReadCounter(String topicName, String cursorName) throws Exception { + PersistentTopic persistentTopic = + (PersistentTopic) pulsar.getBrokerService().getTopic(topicName, false).join().get(); + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + ManagedCursorImpl cursor = (ManagedCursorImpl) managedLedger.openCursor(cursorName); + managedLedger.getCursors().removeCursor(cursor.getName()); + managedLedger.getActiveCursors().removeCursor(cursor.getName()); + ManagedCursorImpl spyCursor = Mockito.spy(cursor); + managedLedger.getCursors().add(spyCursor, PositionImpl.EARLIEST); + managedLedger.getActiveCursors().add(spyCursor, PositionImpl.EARLIEST); + AtomicInteger replyReadCounter = new AtomicInteger(); + Mockito.doAnswer(invocation -> { + if (!String.valueOf(invocation.getArguments()[2]).equals("Normal")) { + replyReadCounter.incrementAndGet(); + } + return invocation.callRealMethod(); + }).when(spyCursor).asyncReplayEntries(Mockito.anySet(), Mockito.any(), Mockito.any()); + Mockito.doAnswer(invocation -> { + if (!String.valueOf(invocation.getArguments()[2]).equals("Normal")) { + replyReadCounter.incrementAndGet(); + } + return invocation.callRealMethod(); + }).when(spyCursor).asyncReplayEntries(Mockito.anySet(), Mockito.any(), Mockito.any(), Mockito.anyBoolean()); + admin.topics().createSubscription(topicName, cursorName, MessageId.earliest); + return replyReadCounter; + } + + @Test + public void testNoRepeatedReadAndDiscard() throws Exception { + int delayedMessages = 100; + final String topic = BrokerTestUtil.newUniqueName("persistent://public/default/tp"); + final String subName = "my-sub"; + admin.topics().createNonPartitionedTopic(topic); + AtomicInteger replyReadCounter = injectReplayReadCounter(topic, subName); + + // Send messages. + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.INT32).topic(topic).enableBatching(false).create(); + for (int i = 0; i < delayedMessages; i++) { + MessageId messageId = producer.newMessage() + .key(String.valueOf(random.nextInt(NUMBER_OF_KEYS))) + .value(100 + i) + .send(); + log.info("Published delayed message :{}", messageId); + } + producer.close(); + + // Make ack holes. + Consumer consumer1 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + Consumer consumer2 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + List msgList1 = new ArrayList<>(); + List msgList2 = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + Message msg1 = consumer1.receive(1, TimeUnit.SECONDS); + if (msg1 != null) { + msgList1.add(msg1); + } + Message msg2 = consumer2.receive(1, TimeUnit.SECONDS); + if (msg2 != null) { + msgList2.add(msg2); + } + } + Consumer redeliverConsumer = null; + if (!msgList1.isEmpty()) { + msgList1.forEach(msg -> consumer1.acknowledgeAsync(msg)); + redeliverConsumer = consumer2; + } else { + msgList2.forEach(msg -> consumer2.acknowledgeAsync(msg)); + redeliverConsumer = consumer1; + } + + // consumer3 will be added to the "recentJoinedConsumers". + Consumer consumer3 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + redeliverConsumer.close(); + + // Verify: no repeated Read-and-discard. + Thread.sleep(5 * 1000); + int maxReplayCount = delayedMessages * 2; + log.info("Reply read count: {}", replyReadCounter.get()); + assertTrue(replyReadCounter.get() < maxReplayCount); + + // cleanup. + consumer1.close(); + consumer2.close(); + consumer3.close(); + admin.topics().delete(topic, false); + } } From 7dfda42aab0aca07878b675caf89cec13cd01be7 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 02:40:48 +0800 Subject: [PATCH 03/12] add test --- .../PersistentStickyKeyDispatcherMultipleConsumers.java | 1 - 1 file changed, 1 deletion(-) 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 a4fa5aea4c384..61b423f54d0f6 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 @@ -472,7 +472,6 @@ private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet // Avoid error due to concurrent modifying. return src; } - // Group by key_hash. NavigableSet res = new TreeSet<>(); final Map> groupedEntries = localGroupedPositions.get(); groupedEntries.clear(); From c6e903320146cbfb767242631f5085112acfe7bf Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 02:41:31 +0800 Subject: [PATCH 04/12] add test --- .../PersistentStickyKeyDispatcherMultipleConsumers.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 61b423f54d0f6..df9345d5fcd25 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 @@ -481,13 +481,13 @@ private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet res.add(pos); } Consumer c = selector.select(stickyKeyHash.intValue()); - if (getAvailablePermits(c) == 0) { + if (c == null) { continue; } - PositionImpl currentMaxReadPosition = recentlyJoinedConsumers.get(c); - if (c == null) { + if (getAvailablePermits(c) == 0) { continue; } + PositionImpl currentMaxReadPosition = recentlyJoinedConsumers.get(c); if (currentMaxReadPosition == null) { res.add(pos); } else if (pos.compareTo(firstMaxReadPos) < 0 && pos.compareTo(currentMaxReadPosition) < 0) { From b3d9d3c1fbe3cc1c8b8e49fb9ec0e62590d5fa48 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 02:43:29 +0800 Subject: [PATCH 05/12] add test --- .../PersistentStickyKeyDispatcherMultipleConsumers.java | 1 + 1 file changed, 1 insertion(+) 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 df9345d5fcd25..7ea20efeea1c3 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 @@ -482,6 +482,7 @@ private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet } Consumer c = selector.select(stickyKeyHash.intValue()); if (c == null) { + // Maybe using HashRangeExclusiveStickyKeyConsumerSelector. continue; } if (getAvailablePermits(c) == 0) { From f29aa67ad289241157f432afcaab57369aad9e31 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 14:40:50 +0800 Subject: [PATCH 06/12] fix NPE --- .../PersistentStickyKeyDispatcherMultipleConsumers.java | 4 ++++ 1 file changed, 4 insertions(+) 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 7ea20efeea1c3..4bd4858664dbd 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 @@ -410,6 +410,9 @@ && removeConsumersFromRecentJoinedConsumers()) { } private boolean removeConsumersFromRecentJoinedConsumers() { + if (MapUtils.isEmpty(recentlyJoinedConsumers)) { + return false; + } Iterator> itr = recentlyJoinedConsumers.entrySet().iterator(); boolean hasConsumerRemovedFromTheRecentJoinedConsumers = false; PositionImpl mdp = (PositionImpl) cursor.getMarkDeletedPosition(); @@ -479,6 +482,7 @@ private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet Long stickyKeyHash = redeliveryMessages.getHash(pos.getLedgerId(), pos.getEntryId()); if (stickyKeyHash == null) { res.add(pos); + continue; } Consumer c = selector.select(stickyKeyHash.intValue()); if (c == null) { From 7f70fd72aa40395a8c9109c3c8d8ca11316d4131 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 12 Mar 2024 21:24:59 +0800 Subject: [PATCH 07/12] fix bugs --- ...PersistentDispatcherMultipleConsumers.java | 28 ++++++++++++++++--- ...tStickyKeyDispatcherMultipleConsumers.java | 23 ++------------- .../client/api/KeySharedSubscriptionTest.java | 2 +- 3 files changed, 28 insertions(+), 25 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 039104fe0221a..650821de8b115 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 @@ -336,22 +336,32 @@ public synchronized void readMoreEntries() { NavigableSet messagesToReplayNow = getMessagesToReplayNow(messagesToRead); if (!messagesToReplayNow.isEmpty()) { + // In Key_Shared mode: It is possible that all resend messages will be discarded: + // - Partly because the target consumer does not have enough permits + // - Partly because of mechanism “recentJoinedConsumers”. + NavigableSet messagesToReplayNowFiltered = + filterOutMessagesWillBeDiscarded(messagesToReplayNow); + if (messagesToReplayNowFiltered.isEmpty()) { + // No messages can be delivery now. + return; + } if (log.isDebugEnabled()) { - log.debug("[{}] Schedule replay of {} messages for {} consumers", name, messagesToReplayNow.size(), - consumerList.size()); + log.debug("[{}] Schedule replay of {} messages for {} consumers", name, + messagesToReplayNowFiltered.size(), consumerList.size()); } havePendingReplayRead = true; minReplayedPosition = messagesToReplayNow.first(); Set deletedMessages = topic.isDelayedDeliveryEnabled() - ? asyncReplayEntriesInOrder(messagesToReplayNow) : asyncReplayEntries(messagesToReplayNow); + ? asyncReplayEntriesInOrder(messagesToReplayNowFiltered) + : asyncReplayEntries(messagesToReplayNowFiltered); // clear already acked positions from replay bucket deletedMessages.forEach(position -> redeliveryMessages.remove(((PositionImpl) position).getLedgerId(), ((PositionImpl) position).getEntryId())); // if all the entries are acked-entries and cleared up from redeliveryMessages, try to read // next entries as readCompletedEntries-callback was never called - if ((messagesToReplayNow.size() - deletedMessages.size()) == 0) { + if ((messagesToReplayNowFiltered.size() - deletedMessages.size()) == 0) { havePendingReplayRead = false; readMoreEntriesAsync(); } @@ -1179,6 +1189,16 @@ protected synchronized NavigableSet getMessagesToReplayNow(int max } } + /** + * This method is in order to avoid the scenario below: + * - Read entries from the Replay queue. + * - The Key_Shared anti-ordering mechanism filtered out all of the entries. + * - Delivery non entry to the client, but we did a BK read. + */ + protected NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { + return src; + } + protected synchronized boolean shouldPauseDeliveryForDelayTracker() { return delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().shouldPauseAllDeliveries(); } 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 4bd4858664dbd..5e0acdbd0eadc 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 @@ -168,14 +168,6 @@ protected Map> initialValue() throws Exception { } }; - private static final FastThreadLocal>> localGroupedPositions = - new FastThreadLocal>>() { - @Override - protected Map> initialValue() throws Exception { - return new HashMap<>(); - } - }; - @Override protected synchronized boolean trySendMessagesToConsumers(ReadType readType, List entries) { long totalMessagesSent = 0; @@ -294,7 +286,6 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis EntryBatchIndexesAcks batchIndexesAcks = EntryBatchIndexesAcks.get(messagesForC); totalEntries += filterEntriesForConsumer(entriesWithSameKey, batchSizes, sendMessageInfo, batchIndexesAcks, cursor, readType == ReadType.Replay, consumer); - consumer.sendMessages(entriesWithSameKey, batchSizes, batchIndexesAcks, sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes(), sendMessageInfo.getTotalChunkedMessages(), @@ -441,8 +432,7 @@ protected synchronized NavigableSet getMessagesToReplayNow(int max this.isDispatcherStuckOnReplays = false; return Collections.emptyNavigableSet(); } else { - NavigableSet res = super.getMessagesToReplayNow(maxMessagesToRead); - return filterOutMessagesWillBeDiscarded(res); + return super.getMessagesToReplayNow(maxMessagesToRead); } } @@ -456,13 +446,8 @@ private int getAvailablePermits(Consumer c) { return availablePermits; } - /** - * This method is in order to avoid the scenario below: - * - Read entries from the Replay queue. - * - The Key_Shared anti-ordering mechanism filtered out all of the entries. - * - Delivery non entry to the client, but we did a BK read. - */ - private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { + @Override + protected NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { // Remove invalid items. removeConsumersFromRecentJoinedConsumers(); if (MapUtils.isEmpty(recentlyJoinedConsumers) || src.isEmpty()) { @@ -476,8 +461,6 @@ private NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet return src; } NavigableSet res = new TreeSet<>(); - final Map> groupedEntries = localGroupedPositions.get(); - groupedEntries.clear(); for (PositionImpl pos : src) { Long stickyKeyHash = redeliveryMessages.getHash(pos.getLedgerId(), pos.getEntryId()); if (stickyKeyHash == null) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index a7c7a72505b87..58d0d3adda3bb 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -1721,7 +1721,7 @@ public void testNoRepeatedReadAndDiscard() throws Exception { Consumer consumer3 = pulsarClient.newConsumer(Schema.INT32) .topic(topic) .subscriptionName(subName) - .receiverQueueSize(10) + .receiverQueueSize(1000) .subscriptionType(SubscriptionType.Key_Shared) .subscribe(); redeliverConsumer.close(); From 05d558fe49f6adcbe7ec42fad061caa43a9b96ee Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 13 Mar 2024 20:12:37 +0800 Subject: [PATCH 08/12] add test: testRecentJoinedPosWillNotStuckOtherConsumer --- .../client/api/KeySharedSubscriptionTest.java | 145 ++++++++++++++++++ .../client/api/ProducerConsumerBase.java | 66 ++++++++ ...SubscriptionPauseOnAckStatPersistTest.java | 78 +--------- 3 files changed, 218 insertions(+), 71 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 58d0d3adda3bb..97aed6f12a3bc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -38,6 +38,7 @@ import java.util.Optional; import java.util.Random; import java.util.Set; +import java.util.TreeSet; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentSkipListSet; @@ -48,6 +49,7 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; import lombok.Cleanup; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; @@ -1738,4 +1740,147 @@ public void testNoRepeatedReadAndDiscard() throws Exception { consumer3.close(); admin.topics().delete(topic, false); } + + /** + * This test is in order to guarantee the feature added by https://github.com/apache/pulsar/pull/7105. + * 1. Start 3 consumers: + * - consumer1 will be closed and trigger a messages redeliver. + * - consumer2 will not ack any messages to make the new consumer joined late will be stuck due + * to the mechanism "recentlyJoinedConsumers". + * - consumer3 will always receive and ack messages. + * 2. Add consumer4 after consumer1 was close, and consumer4 will be stuck due to the mechanism + * "recentlyJoinedConsumers". + * 3. Verify: + * - (Main purpose) consumer3 can still receive messages util the cursor.readerPosition is larger than LAC. + * - no repeated Read-and-discard. + * - at last, all messages will be received. + */ + @Test(timeOut = 180 * 1000) // the test will be finished in 60s. + public void testRecentJoinedPosWillNotStuckOtherConsumer() throws Exception { + final int messagesSentCount = 100; + final Set totalReceivedMessages = new TreeSet<>(); + final String topic = BrokerTestUtil.newUniqueName("persistent://public/default/tp"); + final String subName = "my-sub"; + admin.topics().createNonPartitionedTopic(topic); + AtomicInteger replyReadCounter = injectReplayReadCounter(topic, subName); + + // Send messages. + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.INT32).topic(topic).enableBatching(false).create(); + for (int i = 0; i < messagesSentCount; i++) { + MessageId messageId = producer.newMessage() + .key(String.valueOf(random.nextInt(NUMBER_OF_KEYS))) + .value(100 + i) + .send(); + log.info("Published delayed message :{}", messageId); + } + producer.close(); + + // 1. Start 3 consumers and make ack holes. + // - one consumer will be closed and trigger a messages redeliver. + // - one consumer will not ack any messages to make the new consumer joined late will be stuck due to the + // mechanism "recentlyJoinedConsumers". + // - one consumer will always receive and ack messages. + Consumer consumer1 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + Consumer consumer2 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + Consumer consumer3 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + List msgList1 = new ArrayList<>(); + List msgList2 = new ArrayList<>(); + List msgList3 = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + Message msg1 = consumer1.receive(1, TimeUnit.SECONDS); + if (msg1 != null) { + totalReceivedMessages.add(msg1.getValue()); + msgList1.add(msg1); + } + Message msg2 = consumer2.receive(1, TimeUnit.SECONDS); + if (msg2 != null) { + totalReceivedMessages.add(msg2.getValue()); + msgList2.add(msg2); + } + Message msg3 = consumer3.receive(1, TimeUnit.SECONDS); + if (msg2 != null) { + totalReceivedMessages.add(msg3.getValue()); + msgList3.add(msg3); + } + } + Consumer consumerWillBeClose = null; + Consumer consumerAlwaysAck = null; + Consumer consumerStuck = null; + if (!msgList1.isEmpty()) { + msgList1.forEach(msg -> consumer1.acknowledgeAsync(msg)); + consumerAlwaysAck = consumer1; + consumerWillBeClose = consumer2; + consumerStuck = consumer3; + } else if (!msgList2.isEmpty()){ + msgList2.forEach(msg -> consumer2.acknowledgeAsync(msg)); + consumerAlwaysAck = consumer2; + consumerWillBeClose = consumer3; + consumerStuck = consumer1; + } else { + msgList3.forEach(msg -> consumer3.acknowledgeAsync(msg)); + consumerAlwaysAck = consumer3; + consumerWillBeClose = consumer1; + consumerStuck = consumer2; + } + + // 2. Add consumer4 after "consumerWillBeClose" was close, and consumer4 will be stuck due to the mechanism + // "recentlyJoinedConsumers". + Consumer consumer4 = pulsarClient.newConsumer(Schema.INT32) + .topic(topic) + .subscriptionName(subName) + .receiverQueueSize(1000) + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + consumerWillBeClose.close(); + + // Verify: "consumerAlwaysAck" can receive messages util the cursor.readerPosition is larger than LAC. + while (true) { + Message msg = consumerAlwaysAck.receive(2, TimeUnit.SECONDS); + if (msg == null) { + break; + } + totalReceivedMessages.add(msg.getValue()); + consumerAlwaysAck.acknowledge(msg); + } + PersistentTopic persistentTopic = + (PersistentTopic) pulsar.getBrokerService().getTopic(topic, false).join().get(); + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + ManagedCursorImpl cursor = (ManagedCursorImpl) managedLedger.openCursor(subName); + log.info("cursor_readPosition {}, LAC {}", cursor.getReadPosition(), managedLedger.getLastConfirmedEntry()); + assertTrue(((PositionImpl) cursor.getReadPosition()) + .compareTo((PositionImpl) managedLedger.getLastConfirmedEntry()) > 0); + // Verify: no repeated Read-and-discard. + Thread.sleep(5 * 1000); + int maxReplayCount = messagesSentCount * 2; + log.info("Reply read count: {}", replyReadCounter.get()); + assertTrue(replyReadCounter.get() < maxReplayCount); + // Verify: at last, all messages will be received. + ReceivedMessages receivedMessages = ackAllMessages(consumerAlwaysAck, consumerStuck, consumer4); + totalReceivedMessages.addAll(receivedMessages.messagesReceived.stream().map(p -> p.getRight()).collect( + Collectors.toList())); + assertEquals(totalReceivedMessages.size(), messagesSentCount); + + // cleanup. + consumer1.close(); + consumer2.close(); + consumer3.close(); + consumer4.close(); + admin.topics().delete(topic, false); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerBase.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerBase.java index f58c1fa26afc7..ef070250ca1aa 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerBase.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerBase.java @@ -21,9 +21,14 @@ import com.google.common.collect.Sets; import java.lang.reflect.Method; +import java.util.ArrayList; +import java.util.List; import java.util.Random; import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.function.BiFunction; +import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.TenantInfoImpl; @@ -69,4 +74,65 @@ protected String newTopicName() { return "my-property/my-ns/topic-" + Long.toHexString(random.nextLong()); } + protected ReceivedMessages receiveAndAckMessages( + BiFunction ackPredicate, + Consumer...consumers) throws Exception { + ReceivedMessages receivedMessages = new ReceivedMessages(); + while (true) { + int receivedMsgCount = 0; + for (int i = 0; i < consumers.length; i++) { + Consumer consumer = consumers[i]; + while (true) { + Message msg = consumer.receive(2, TimeUnit.SECONDS); + if (msg != null) { + receivedMsgCount++; + T v = msg.getValue(); + MessageId messageId = msg.getMessageId(); + receivedMessages.messagesReceived.add(Pair.of(msg.getMessageId(), v)); + if (ackPredicate.apply(messageId, v)) { + consumer.acknowledge(msg); + receivedMessages.messagesAcked.add(Pair.of(msg.getMessageId(), v)); + } + } else { + break; + } + } + } + // Because of the possibility of consumers getting stuck with each other, only jump out of the loop if all + // consumers could not receive messages. + if (receivedMsgCount == 0) { + break; + } + } + return receivedMessages; + } + + protected ReceivedMessages ackAllMessages(Consumer...consumers) throws Exception { + return receiveAndAckMessages((msgId, msgV) -> true, consumers); + } + + protected static class ReceivedMessages { + + List> messagesReceived = new ArrayList<>(); + + List> messagesAcked = new ArrayList<>(); + + public boolean hasReceivedMessage(T v) { + for (Pair pair : messagesReceived) { + if (pair.getRight().equals(v)) { + return true; + } + } + return false; + } + + public boolean hasAckedMessage(T v) { + for (Pair pair : messagesAcked) { + if (pair.getRight().equals(v)) { + return true; + } + } + return false; + } + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java index 22edbc36f6ce0..9a4de8ecf21cc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionPauseOnAckStatPersistTest.java @@ -21,13 +21,10 @@ import java.lang.reflect.Method; import java.util.ArrayList; import java.util.Arrays; -import java.util.List; import java.util.Map; import java.util.concurrent.TimeUnit; -import java.util.function.BiFunction; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Position; -import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.service.Dispatcher; import org.apache.pulsar.broker.service.SystemTopicBasedTopicPoliciesService; @@ -148,71 +145,10 @@ private enum SkipType{ RESET_CURSOR; } - private ReceivedMessages receiveAndAckMessages(BiFunction ackPredicate, - Consumer...consumers) throws Exception { - ReceivedMessages receivedMessages = new ReceivedMessages(); - while (true) { - int receivedMsgCount = 0; - for (int i = 0; i < consumers.length; i++) { - Consumer consumer = consumers[i]; - while (true) { - Message msg = consumer.receive(2, TimeUnit.SECONDS); - if (msg != null) { - receivedMsgCount++; - String v = msg.getValue(); - MessageId messageId = msg.getMessageId(); - receivedMessages.messagesReceived.add(Pair.of(msg.getMessageId(), v)); - if (ackPredicate.apply(messageId, v)) { - consumer.acknowledge(msg); - receivedMessages.messagesAcked.add(Pair.of(msg.getMessageId(), v)); - } - } else { - break; - } - } - } - // Because of the possibility of consumers getting stuck with each other, only jump out of the loop if all - // consumers could not receive messages. - if (receivedMsgCount == 0) { - break; - } - } - return receivedMessages; - } - - private ReceivedMessages ackAllMessages(Consumer...consumers) throws Exception { - return receiveAndAckMessages((msgId, msgV) -> true, consumers); - } - - private ReceivedMessages ackOddMessagesOnly(Consumer...consumers) throws Exception { + private ReceivedMessages ackOddMessagesOnly(Consumer...consumers) throws Exception { return receiveAndAckMessages((msgId, msgV) -> Integer.valueOf(msgV) % 2 == 1, consumers); } - private static class ReceivedMessages { - - List> messagesReceived = new ArrayList<>(); - - List> messagesAcked = new ArrayList<>(); - - public boolean hasReceivedMessage(String v) { - for (Pair pair : messagesReceived) { - if (pair.getRight().equals(v)) { - return true; - } - } - return false; - } - - public boolean hasAckedMessage(String v) { - for (Pair pair : messagesAcked) { - if (pair.getRight().equals(v)) { - return true; - } - } - return false; - } - } - @DataProvider(name = "typesOfSetDispatcherPauseOnAckStatePersistent") public Object[][] typesOfSetDispatcherPauseOnAckStatePersistent() { return new Object[][]{ @@ -367,7 +303,7 @@ public void testPauseOnAckStatPersist(SubscriptionType subscriptionType) throws // Verify: after ack messages, will unpause the dispatcher. c1.acknowledge(messageIdsSent); - ReceivedMessages receivedMessagesAfterPause = ackAllMessages(c1); + ReceivedMessages receivedMessagesAfterPause = ackAllMessages(c1); Assert.assertTrue(receivedMessagesAfterPause.hasReceivedMessage(specifiedMessage)); Assert.assertTrue(receivedMessagesAfterPause.hasAckedMessage(specifiedMessage)); @@ -417,7 +353,7 @@ public void testUnPauseOnSkipEntries(SkipType skipType) throws Exception { final String specifiedMessage2 = "9876543211"; p1.send(specifiedMessage2); - ReceivedMessages receivedMessagesAfterPause = ackAllMessages(c1); + ReceivedMessages receivedMessagesAfterPause = ackAllMessages(c1); Assert.assertTrue(receivedMessagesAfterPause.hasReceivedMessage(specifiedMessage2)); Assert.assertTrue(receivedMessagesAfterPause.hasAckedMessage(specifiedMessage2)); @@ -520,7 +456,7 @@ public void testPauseOnAckStatPersistNotAffectReplayRead(SubscriptionType subscr messageIdsSent.add(messageId); } // Make ack holes. - ReceivedMessages receivedMessagesC1 = ackOddMessagesOnly(c1); + ReceivedMessages receivedMessagesC1 = ackOddMessagesOnly(c1); verifyAckHolesIsMuchThanLimit(tpName, subscription); cancelPendingRead(tpName, subscription); @@ -540,7 +476,7 @@ public void testPauseOnAckStatPersistNotAffectReplayRead(SubscriptionType subscr // Verify: close the previous consumer, the new one could receive all messages. c1.close(); - ReceivedMessages receivedMessagesC2 = ackAllMessages(c2); + ReceivedMessages receivedMessagesC2 = ackAllMessages(c2); int messageCountAckedByC1 = receivedMessagesC1.messagesAcked.size(); int messageCountAckedByC2 = receivedMessagesC2.messagesAcked.size(); Assert.assertEquals(messageCountAckedByC2, msgSendCount - messageCountAckedByC1 + specifiedMessageCount); @@ -577,7 +513,7 @@ public void testMultiConsumersPauseOnAckStatPersistNotAffectReplayRead(Subscript messageIdsSent.add(messageId); } // Make ack holes. - ReceivedMessages receivedMessagesC1AndC2 = ackOddMessagesOnly(c1, c2); + ReceivedMessages receivedMessagesC1AndC2 = ackOddMessagesOnly(c1, c2); verifyAckHolesIsMuchThanLimit(tpName, subscription); cancelPendingRead(tpName, subscription); @@ -601,7 +537,7 @@ public void testMultiConsumersPauseOnAckStatPersistNotAffectReplayRead(Subscript // Verify: close the previous consumer, the new one could receive all messages. c1.close(); c2.close(); - ReceivedMessages receivedMessagesC3AndC4 = ackAllMessages(c3, c4); + ReceivedMessages receivedMessagesC3AndC4 = ackAllMessages(c3, c4); int messageCountAckedByC1AndC2 = receivedMessagesC1AndC2.messagesAcked.size(); int messageCountAckedByC3AndC4 = receivedMessagesC3AndC4.messagesAcked.size(); Assert.assertEquals(messageCountAckedByC3AndC4, From 76f5a33447a78744ee8c52f4e3543fc70db252be Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Sat, 16 Mar 2024 19:57:08 +0800 Subject: [PATCH 09/12] guarantee the mechanism provided by #7105 --- ...PersistentDispatcherMultipleConsumers.java | 51 +++++++++++-------- ...tStickyKeyDispatcherMultipleConsumers.java | 25 ++++++++- .../client/api/KeySharedSubscriptionTest.java | 25 ++++++--- 3 files changed, 74 insertions(+), 27 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 650821de8b115..b441400dae11f 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 @@ -334,34 +334,25 @@ public synchronized void readMoreEntries() { } NavigableSet messagesToReplayNow = getMessagesToReplayNow(messagesToRead); - - if (!messagesToReplayNow.isEmpty()) { - // In Key_Shared mode: It is possible that all resend messages will be discarded: - // - Partly because the target consumer does not have enough permits - // - Partly because of mechanism “recentJoinedConsumers”. - NavigableSet messagesToReplayNowFiltered = - filterOutMessagesWillBeDiscarded(messagesToReplayNow); - if (messagesToReplayNowFiltered.isEmpty()) { - // No messages can be delivery now. - return; - } + NavigableSet messagesToReplayFiltered = filterOutEntriesWillBeDiscarded(messagesToReplayNow); + if (!messagesToReplayFiltered.isEmpty()) { if (log.isDebugEnabled()) { log.debug("[{}] Schedule replay of {} messages for {} consumers", name, - messagesToReplayNowFiltered.size(), consumerList.size()); + messagesToReplayFiltered.size(), consumerList.size()); } havePendingReplayRead = true; minReplayedPosition = messagesToReplayNow.first(); Set deletedMessages = topic.isDelayedDeliveryEnabled() - ? asyncReplayEntriesInOrder(messagesToReplayNowFiltered) - : asyncReplayEntries(messagesToReplayNowFiltered); + ? asyncReplayEntriesInOrder(messagesToReplayFiltered) + : asyncReplayEntries(messagesToReplayFiltered); // clear already acked positions from replay bucket deletedMessages.forEach(position -> redeliveryMessages.remove(((PositionImpl) position).getLedgerId(), ((PositionImpl) position).getEntryId())); // if all the entries are acked-entries and cleared up from redeliveryMessages, try to read // next entries as readCompletedEntries-callback was never called - if ((messagesToReplayNowFiltered.size() - deletedMessages.size()) == 0) { + if ((messagesToReplayFiltered.size() - deletedMessages.size()) == 0) { havePendingReplayRead = false; readMoreEntriesAsync(); } @@ -370,7 +361,7 @@ public synchronized void readMoreEntries() { log.debug("[{}] Dispatcher read is blocked due to unackMessages {} reached to max {}", name, totalUnackedMessages, topic.getMaxUnackedMessagesOnSubscription()); } - } else if (!havePendingRead) { + } else if (!havePendingRead && hasConsumersNeededNormalRead()) { if (shouldPauseOnAckStatePersist(ReadType.Normal)) { if (log.isDebugEnabled()) { log.debug("[{}] [{}] Skipping read for the topic, Due to blocked on ack state persistent.", @@ -406,7 +397,16 @@ public synchronized void readMoreEntries() { topic.getMaxReadPosition()); } } else { - log.debug("[{}] Cannot schedule next read until previous one is done", name); + if (log.isDebugEnabled()) { + if (!messagesToReplayNow.isEmpty()) { + log.debug("[{}] [{}] Skipping read for the topic: because all entries in replay queue were" + + " filtered out due to the mechanism of Key_Shared mode, and the left consumers have" + + " no permits now", + topic.getName(), getSubscriptionName()); + } else { + log.debug("[{}] Cannot schedule next read until previous one is done", name); + } + } } } else { if (log.isDebugEnabled()) { @@ -1190,15 +1190,26 @@ protected synchronized NavigableSet getMessagesToReplayNow(int max } /** + * This is a mode method designed for Key_Shared mode. + * Filter out the entries that will be discarded due to the order guarantee mechanism of Key_Shared mode. * This method is in order to avoid the scenario below: - * - Read entries from the Replay queue. - * - The Key_Shared anti-ordering mechanism filtered out all of the entries. + * - Get positions from the Replay queue. + * - Read entries from BK. + * - The order guarantee mechanism of Key_Shared mode filtered out all the entries. * - Delivery non entry to the client, but we did a BK read. */ - protected NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { + protected NavigableSet filterOutEntriesWillBeDiscarded(NavigableSet src) { return src; } + /** + * This is a mode method designed for Key_Shared mode, to avoid unnecessary stuck. + * See detail {@link PersistentStickyKeyDispatcherMultipleConsumers#hasConsumersNeededNormalRead}. + */ + protected boolean hasConsumersNeededNormalRead() { + return true; + } + protected synchronized boolean shouldPauseDeliveryForDelayTracker() { return delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().shouldPauseAllDeliveries(); } 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 5e0acdbd0eadc..c24e9383be661 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 @@ -447,7 +447,7 @@ private int getAvailablePermits(Consumer c) { } @Override - protected NavigableSet filterOutMessagesWillBeDiscarded(NavigableSet src) { + protected NavigableSet filterOutEntriesWillBeDiscarded(NavigableSet src) { // Remove invalid items. removeConsumersFromRecentJoinedConsumers(); if (MapUtils.isEmpty(recentlyJoinedConsumers) || src.isEmpty()) { @@ -489,6 +489,29 @@ protected NavigableSet filterOutMessagesWillBeDiscarded(NavigableS return res; } + /** + * In Key_Shared mode, the consumer will not receive any entries from a normal reading if it is included in + * {@link #recentlyJoinedConsumers}, they can only receive entries from replay reads. + * If all entries in {@link #redeliveryMessages} have been filtered out due to the order guarantee mechanism, + * Broker need a normal read to make the consumers not included in @link #recentlyJoinedConsumers} will not be + * stuck. See https://github.com/apache/pulsar/pull/7105. + */ + @Override + protected boolean hasConsumersNeededNormalRead() { + for (Consumer consumer : consumerList) { + if (consumer == null || consumer.isBlocked()) { + continue; + } + if (recentlyJoinedConsumers.containsKey(consumer)) { + continue; + } + if (consumer.getAvailablePermits() > 0) { + return true; + } + } + return false; + } + @Override public SubType getType() { return SubType.Key_Shared; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 97aed6f12a3bc..1720ddfbaf954 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -1757,7 +1757,7 @@ public void testNoRepeatedReadAndDiscard() throws Exception { */ @Test(timeOut = 180 * 1000) // the test will be finished in 60s. public void testRecentJoinedPosWillNotStuckOtherConsumer() throws Exception { - final int messagesSentCount = 100; + final int messagesSentPerTime = 100; final Set totalReceivedMessages = new TreeSet<>(); final String topic = BrokerTestUtil.newUniqueName("persistent://public/default/tp"); final String subName = "my-sub"; @@ -1767,14 +1767,13 @@ public void testRecentJoinedPosWillNotStuckOtherConsumer() throws Exception { // Send messages. @Cleanup Producer producer = pulsarClient.newProducer(Schema.INT32).topic(topic).enableBatching(false).create(); - for (int i = 0; i < messagesSentCount; i++) { + for (int i = 0; i < messagesSentPerTime; i++) { MessageId messageId = producer.newMessage() .key(String.valueOf(random.nextInt(NUMBER_OF_KEYS))) .value(100 + i) .send(); - log.info("Published delayed message :{}", messageId); + log.info("Published message :{}", messageId); } - producer.close(); // 1. Start 3 consumers and make ack holes. // - one consumer will be closed and trigger a messages redeliver. @@ -1849,6 +1848,17 @@ public void testRecentJoinedPosWillNotStuckOtherConsumer() throws Exception { .subscribe(); consumerWillBeClose.close(); + Thread.sleep(2000); + + for (int i = messagesSentPerTime; i < messagesSentPerTime * 2; i++) { + MessageId messageId = producer.newMessage() + .key(String.valueOf(random.nextInt(NUMBER_OF_KEYS))) + .value(100 + i) + .send(); + log.info("Published message :{}", messageId); + } + + // Send messages again. // Verify: "consumerAlwaysAck" can receive messages util the cursor.readerPosition is larger than LAC. while (true) { Message msg = consumerAlwaysAck.receive(2, TimeUnit.SECONDS); @@ -1865,22 +1875,25 @@ public void testRecentJoinedPosWillNotStuckOtherConsumer() throws Exception { log.info("cursor_readPosition {}, LAC {}", cursor.getReadPosition(), managedLedger.getLastConfirmedEntry()); assertTrue(((PositionImpl) cursor.getReadPosition()) .compareTo((PositionImpl) managedLedger.getLastConfirmedEntry()) > 0); + + // Make all consumers to start to read and acknowledge messages. // Verify: no repeated Read-and-discard. Thread.sleep(5 * 1000); - int maxReplayCount = messagesSentCount * 2; + int maxReplayCount = messagesSentPerTime * 2; log.info("Reply read count: {}", replyReadCounter.get()); assertTrue(replyReadCounter.get() < maxReplayCount); // Verify: at last, all messages will be received. ReceivedMessages receivedMessages = ackAllMessages(consumerAlwaysAck, consumerStuck, consumer4); totalReceivedMessages.addAll(receivedMessages.messagesReceived.stream().map(p -> p.getRight()).collect( Collectors.toList())); - assertEquals(totalReceivedMessages.size(), messagesSentCount); + assertEquals(totalReceivedMessages.size(), messagesSentPerTime * 2); // cleanup. consumer1.close(); consumer2.close(); consumer3.close(); consumer4.close(); + producer.close(); admin.topics().delete(topic, false); } } From 0cb5077f725004342df5bd883bf735d978e45fda Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 29 Mar 2024 02:12:12 +0800 Subject: [PATCH 10/12] - --- .../org/apache/pulsar/client/api/KeySharedSubscriptionTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 1720ddfbaf954..7219555050839 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -1681,7 +1681,7 @@ public void testNoRepeatedReadAndDiscard() throws Exception { .key(String.valueOf(random.nextInt(NUMBER_OF_KEYS))) .value(100 + i) .send(); - log.info("Published delayed message :{}", messageId); + log.info("Published message :{}", messageId); } producer.close(); From 7fb0a8970f4720157b86b13d5484edb4957627b0 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 29 Mar 2024 03:54:20 +0800 Subject: [PATCH 11/12] reuse getRestrictedMaxEntriesForConsumer --- ...tStickyKeyDispatcherMultipleConsumers.java | 57 ++++++++++--------- 1 file changed, 30 insertions(+), 27 deletions(-) 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 c24e9383be661..da23e4c1eb702 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 @@ -29,11 +29,11 @@ import java.util.List; import java.util.Map; import java.util.NavigableSet; -import java.util.NoSuchElementException; import java.util.Set; import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.Position; @@ -168,6 +168,14 @@ protected Map> initialValue() throws Exception { } }; + private static final FastThreadLocal>> localGroupedPositions = + new FastThreadLocal>>() { + @Override + protected Map> initialValue() throws Exception { + return new HashMap<>(); + } + }; + @Override protected synchronized boolean trySendMessagesToConsumers(ReadType readType, List entries) { long totalMessagesSent = 0; @@ -252,8 +260,8 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis List entriesWithSameKey = current.getValue(); int entriesWithSameKeyCount = entriesWithSameKey.size(); int availablePermits = getAvailablePermits(consumer); - int maxMessagesForC = Math.min(entriesWithSameKeyCount, availablePermits); - int messagesForC = getRestrictedMaxEntriesForConsumer(consumer, entriesWithSameKey, maxMessagesForC, + int messagesForC = getRestrictedMaxEntriesForConsumer(consumer, + entriesWithSameKey.stream().map(Entry::getPosition).collect(Collectors.toList()), availablePermits, readType, consumerStickyKeyHashesMap.get(consumer)); if (log.isDebugEnabled()) { log.debug("[{}] select consumer {} with messages num {}, read type is {}", @@ -328,8 +336,9 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis return false; } - private int getRestrictedMaxEntriesForConsumer(Consumer consumer, List entries, int maxMessages, - ReadType readType, Set stickyKeyHashes) { + private int getRestrictedMaxEntriesForConsumer(Consumer consumer, List entries, + int availablePermits, ReadType readType, Set stickyKeyHashes) { + int maxMessages = Math.min(entries.size(), availablePermits); if (maxMessages == 0) { return 0; } @@ -374,7 +383,7 @@ private int getRestrictedMaxEntriesForConsumer(Consumer consumer, List en // Here, the consumer is one that has recently joined, so we can only send messages that were // published before it has joined. for (int i = 0; i < maxMessages; i++) { - if (((PositionImpl) entries.get(i).getPosition()).compareTo(maxReadPosition) >= 0) { + if (((PositionImpl) entries.get(i)).compareTo(maxReadPosition) >= 0) { // We have already crossed the divider line. All messages in the list are now // newer than what we can currently dispatch to this consumer return i; @@ -447,20 +456,14 @@ private int getAvailablePermits(Consumer c) { } @Override - protected NavigableSet filterOutEntriesWillBeDiscarded(NavigableSet src) { - // Remove invalid items. - removeConsumersFromRecentJoinedConsumers(); - if (MapUtils.isEmpty(recentlyJoinedConsumers) || src.isEmpty()) { - return src; - } - PositionImpl firstMaxReadPos = null; - try { - firstMaxReadPos = recentlyJoinedConsumers.values().iterator().next(); - } catch (NoSuchElementException noSuchElementException) { - // Avoid error due to concurrent modifying. + protected synchronized NavigableSet filterOutEntriesWillBeDiscarded(NavigableSet src) { + if (src.isEmpty()) { return src; } NavigableSet res = new TreeSet<>(); + // Group positions. + final Map> groupedPositions = localGroupedPositions.get(); + groupedPositions.clear(); for (PositionImpl pos : src) { Long stickyKeyHash = redeliveryMessages.getHash(pos.getLedgerId(), pos.getEntryId()); if (stickyKeyHash == null) { @@ -472,18 +475,18 @@ protected NavigableSet filterOutEntriesWillBeDiscarded(NavigableSe // Maybe using HashRangeExclusiveStickyKeyConsumerSelector. continue; } - if (getAvailablePermits(c) == 0) { + groupedPositions.computeIfAbsent(c, k -> new ArrayList<>()).add(pos); + } + // Filter positions by the Recently Joined Position rule. + for (Map.Entry> item : groupedPositions.entrySet()) { + int availablePermits = getAvailablePermits(item.getKey()); + if (availablePermits == 0) { continue; } - PositionImpl currentMaxReadPosition = recentlyJoinedConsumers.get(c); - if (currentMaxReadPosition == null) { - res.add(pos); - } else if (pos.compareTo(firstMaxReadPos) < 0 && pos.compareTo(currentMaxReadPosition) < 0) { - res.add(pos); - } else { - // The subsequent positions will also be larger than "firstMaxReadPos", why need to check continuous. - // Because the subsequent positions may be delivered to a consumer which does not in the - // "recentlyJoinedConsumers". + int posCountToRead = getRestrictedMaxEntriesForConsumer(item.getKey(), item.getValue(), availablePermits, + ReadType.Replay, null); + for (int i = 0; i < posCountToRead; i++) { + res.add(item.getValue().get(i)); } } return res; From 2b7c3e6b541bf33adb5bba826a160eae1a8c2b3f Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 29 Mar 2024 10:18:20 +0800 Subject: [PATCH 12/12] address comment --- .../PersistentStickyKeyDispatcherMultipleConsumers.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 da23e4c1eb702..ee2ebd7ca867e 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 @@ -485,8 +485,8 @@ protected synchronized NavigableSet filterOutEntriesWillBeDiscarde } int posCountToRead = getRestrictedMaxEntriesForConsumer(item.getKey(), item.getValue(), availablePermits, ReadType.Replay, null); - for (int i = 0; i < posCountToRead; i++) { - res.add(item.getValue().get(i)); + if (posCountToRead > 0) { + res.addAll(item.getValue().subList(0, posCountToRead)); } } return res;