From 4113f073a055a8c57fa738ad892d9e6b4063fc9e Mon Sep 17 00:00:00 2001 From: yunmaoQu <2643354262@qq.com> Date: Sun, 2 Feb 2025 18:43:01 +0800 Subject: [PATCH] [improve]Optimize InMemoryDelayedDeliveryTracker by maintaining state --- .../InMemoryDelayedDeliveryTracker.java | 19 ++++++++++++++++++- ...PersistentDispatcherMultipleConsumers.java | 11 ++++------- 2 files changed, 22 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java index 5796fcbd78550..5f2537c128714 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java @@ -218,9 +218,26 @@ public NavigableSet getScheduledMessages(int maxMessages) { return positions; } + public boolean shouldSkipMessage(long ledgerId, long entryId) { + for (Long2ObjectMap ledgerMap : delayedMessageMap.values()) { + Roaring64Bitmap entryIds = ledgerMap.get(ledgerId); + if (entryIds != null && entryIds.contains(entryId)) { + return true; + } + } + return false; + } + @Override public CompletableFuture clear() { - this.delayedMessageMap.clear(); + long cutoffTime = getCutoffTime(); + delayedMessageMap.headMap(cutoffTime).clear(); + + if (log.isDebugEnabled()) { + log.debug("[{}] Cleared expired delayed messages before {}, remaining messages: {}", + dispatcher.getName(), cutoffTime, getNumberOfDelayedMessages()); + } + return CompletableFuture.completedFuture(null); } 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 3ceb703c26cd8..f383c39161d8c 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 @@ -193,12 +193,6 @@ public synchronized CompletableFuture addConsumer(Consumer consumer) { shouldRewindBeforeReadingOrReplaying = false; } redeliveryMessages.clear(); - delayedDeliveryTracker.ifPresent(tracker -> { - // Don't clean up BucketDelayedDeliveryTracker, otherwise we will lose the bucket snapshot - if (tracker instanceof InMemoryDelayedDeliveryTracker) { - tracker.clear(); - } - }); } if (isConsumersExceededOnSubscription()) { @@ -448,7 +442,10 @@ protected Predicate createReadEntriesSkipConditionForNormalRead() { // Filter out and skip read delayed messages exist in DelayedDeliveryTracker if (delayedDeliveryTracker.isPresent()) { final DelayedDeliveryTracker deliveryTracker = delayedDeliveryTracker.get(); - if (deliveryTracker instanceof BucketDelayedDeliveryTracker) { + if (deliveryTracker instanceof InMemoryDelayedDeliveryTracker) { + skipCondition = position -> ((InMemoryDelayedDeliveryTracker) deliveryTracker) + .shouldSkipMessage(position.getLedgerId(), position.getEntryId()); + } else if (deliveryTracker instanceof BucketDelayedDeliveryTracker) { skipCondition = position -> ((BucketDelayedDeliveryTracker) deliveryTracker) .containsMessage(position.getLedgerId(), position.getEntryId()); }