From 1e2e7ea66b76c6379ce4e089157086d944b8e41b Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Thu, 23 Apr 2026 12:06:56 +0800 Subject: [PATCH] [improve][broker] Centralize mark-delete position change handling --- ...PersistentDispatcherMultipleConsumers.java | 31 ++++++++++------ .../PersistentMessageExpiryMonitor.java | 4 +- ...tStickyKeyDispatcherMultipleConsumers.java | 2 + .../persistent/PersistentSubscription.java | 37 ++++--------------- 4 files changed, 30 insertions(+), 44 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 ab848064074fa..bb7d4d270a3ea 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 @@ -65,6 +65,7 @@ import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; import org.apache.pulsar.broker.service.EntryBatchSizes; import org.apache.pulsar.broker.service.InMemoryRedeliveryTracker; +import org.apache.pulsar.broker.service.PendingAcksMap; import org.apache.pulsar.broker.service.RedeliveryTracker; import org.apache.pulsar.broker.service.RedeliveryTrackerDisabled; import org.apache.pulsar.broker.service.SendMessageInfo; @@ -145,7 +146,6 @@ public class PersistentDispatcherMultipleConsumers extends AbstractPersistentDis protected enum ReadType { Normal, Replay } - private Position lastMarkDeletePositionBeforeReadMoreEntries; private volatile long readMoreEntriesCallCount; public PersistentDispatcherMultipleConsumers(PersistentTopic topic, ManagedCursor cursor, @@ -339,6 +339,24 @@ public void readMoreEntriesAsync() { } } + @Override + public void markDeletePositionMoveForward() { + // When the mark-delete position advances (due to ack or TTL expiry), remove stale entries that are + // now below the new mark-delete position from the redelivery tracker and each consumer's pending acks. + synchronized (PersistentDispatcherMultipleConsumers.this) { + Position mdp = cursor.getMarkDeletedPosition(); + if (mdp != null) { + redeliveryMessages.removeAllUpTo(mdp.getLedgerId(), mdp.getEntryId()); + for (Consumer consumer : consumerList) { + PendingAcksMap pendingAcks = consumer.getPendingAcks(); + if (pendingAcks != null) { + pendingAcks.removeAllUpTo(mdp.getLedgerId(), mdp.getEntryId()); + } + } + } + } + } + public synchronized void readMoreEntries() { if (cursor.isClosed()) { log.debug("Cursor is already closed, skipping read more entries"); @@ -362,17 +380,6 @@ public synchronized void readMoreEntries() { // increment the counter for readMoreEntries calls, to track the number of times readMoreEntries is called readMoreEntriesCallCount++; - // remove possible expired messages from redelivery tracker and pending acks - Position markDeletePosition = cursor.getMarkDeletedPosition(); - if (lastMarkDeletePositionBeforeReadMoreEntries != markDeletePosition) { - redeliveryMessages.removeAllUpTo(markDeletePosition.getLedgerId(), markDeletePosition.getEntryId()); - for (Consumer consumer : consumerList) { - consumer.getPendingAcks() - .removeAllUpTo(markDeletePosition.getLedgerId(), markDeletePosition.getEntryId()); - } - lastMarkDeletePositionBeforeReadMoreEntries = markDeletePosition; - } - // totalAvailablePermits may be updated by other threads int firstAvailableConsumerPermits = getFirstAvailableConsumerPermits(); int currentTotalAvailablePermits = Math.max(totalAvailablePermits, firstAvailableConsumerPermits); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentMessageExpiryMonitor.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentMessageExpiryMonitor.java index bfc8951be4960..0447f8978dc48 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentMessageExpiryMonitor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentMessageExpiryMonitor.java @@ -40,7 +40,6 @@ import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo; import org.apache.pulsar.broker.service.MessageExpirer; import org.apache.pulsar.client.impl.MessageImpl; -import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType; import org.apache.pulsar.common.stats.Rate; import org.jspecify.annotations.Nullable; /** @@ -211,8 +210,7 @@ public void markDeleteComplete(Object ctx) { long numMessagesExpired = (long) ctx - cursor.getNumberOfEntriesInBacklog(false); msgExpired.recordMultipleEvents(numMessagesExpired, 0 /* no value stats */); totalMsgExpired.add(numMessagesExpired); - // If the subscription is a Key_Shared subscription, we should to trigger message dispatch. - if (subscription != null && subscription.getType() == SubType.Key_Shared) { + if (subscription != null && subscription.getDispatcher() != null) { subscription.getDispatcher().markDeletePositionMoveForward(); } expirationCheckInProgress = FALSE; 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 2549e4a34a911..c061711e00f73 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 @@ -642,6 +642,8 @@ protected int getStickyKeyHash(Entry entry) { @Override public void markDeletePositionMoveForward() { + // clean up stale entries from the redelivery tracker and pending acks + super.markDeletePositionMoveForward(); // reschedule a read with a backoff after moving the mark-delete position forward since there might have // been consumers that were blocked by hash and couldn't make progress reScheduleReadWithKeySharedUnblockingInterval(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index cb60a70a0c32e..6b719f831bbe7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -515,12 +515,7 @@ public void markDeleteComplete(Object ctx) { .attr("newMD", newMD) .attr("oldMD", oldMD) .log("Mark deleted messages to position from position"); - // Signal the dispatchers to give chance to take extra actions - if (dispatcher != null) { - dispatcher.afterAckMessages(null, ctx); - } - // Signal the dispatchers to give chance to take extra actions - notifyTheMarkDeletePositionChanged(oldMD); + notifyTheMarkDeletePositionChanged(oldMD); } @Override @@ -530,10 +525,6 @@ public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { .attr("ctx", ctx) .exceptionMessage(exception) .log("Failed to mark delete for position"); - // Signal the dispatchers to give chance to take extra actions - if (dispatcher != null) { - dispatcher.afterAckMessages(null, ctx); - } } }; @@ -544,10 +535,6 @@ public void deleteComplete(Object context) { log.debug() .attr("context", context) .log("Deleted message"); - // Signal the dispatchers to give chance to take extra actions - if (dispatcher != null) { - dispatcher.afterAckMessages(null, context); - } notifyTheMarkDeletePositionChanged((Position) context); } @@ -557,10 +544,6 @@ public void deleteFailed(ManagedLedgerException exception, Object ctx) { .attr("ctx", ctx) .exceptionMessage(exception) .log("Failed to delete message"); - // Signal the dispatchers to give chance to take extra actions - if (dispatcher != null) { - dispatcher.afterAckMessages(exception, ctx); - } } }; @@ -579,6 +562,7 @@ private void notifyTheMarkDeletePositionChanged(Position oldPosition) { handleReplicatedSubscriptionsUpdate(newMD); if (dispatcher != null) { + dispatcher.afterAckMessages(null, null); dispatcher.markDeletePositionMoveForward(); } } @@ -784,6 +768,7 @@ public CompletableFuture clearBacklog() { .attr("entriesInBacklog", cursor.getNumberOfEntriesInBacklog(false)) .log("Backlog size before clearing"); + final Position oldMarkDeletePosition = cursor.getMarkDeletedPosition(); cursor.asyncClearBacklog(new ClearBacklogCallback() { @Override public void clearBacklogComplete(Object ctx) { @@ -798,10 +783,10 @@ public void clearBacklogComplete(Object ctx) { future.complete(null); } }); - dispatcher.afterAckMessages(null, ctx); } else { future.complete(null); } + notifyTheMarkDeletePositionChanged(oldMarkDeletePosition); } @Override @@ -810,9 +795,6 @@ public void clearBacklogFailed(ManagedLedgerException exception, Object ctx) { .exception(exception) .log("Failed to clear backlog"); future.completeExceptionally(exception); - if (dispatcher != null) { - dispatcher.afterAckMessages(exception, ctx); - } } }, null); @@ -827,6 +809,7 @@ public CompletableFuture skipMessages(int numMessagesToSkip) { .attr("numMessagesToSkip", numMessagesToSkip) .attr("entriesInBacklog", cursor.getNumberOfEntriesInBacklog(false)) .log("Skipping messages"); + final Position oldMarkDeletePosition = cursor.getMarkDeletedPosition(); cursor.asyncSkipEntries(numMessagesToSkip, IndividualDeletedEntries.Exclude, new AsyncCallbacks.SkipEntriesCallback() { @Override @@ -836,9 +819,7 @@ public void skipEntriesComplete(Object ctx) { .attr("entriesInBacklog", cursor.getNumberOfEntriesInBacklog(false)) .log("Skipped messages"); future.complete(null); - if (dispatcher != null) { - dispatcher.afterAckMessages(null, ctx); - } + notifyTheMarkDeletePositionChanged(oldMarkDeletePosition); } @Override @@ -848,9 +829,6 @@ public void skipEntriesFailed(ManagedLedgerException exception, Object ctx) { .exception(exception) .log("Failed to skip messages"); future.completeExceptionally(exception); - if (dispatcher != null) { - dispatcher.afterAckMessages(exception, ctx); - } } }, null); @@ -983,6 +961,7 @@ private CompletableFuture resetCursorInternal(Position finalPosition, Comp } forceReset.thenAccept(forceResetValue -> { + final Position oldMarkDeletePosition = cursor.getMarkDeletedPosition(); cursor.asyncResetCursor(finalPosition, forceResetValue, new AsyncCallbacks.ResetCursorCallback() { @Override public void resetComplete(Object ctx) { @@ -991,8 +970,8 @@ public void resetComplete(Object ctx) { .log("Successfully reset subscription to position"); if (dispatcher != null) { dispatcher.cursorIsReset(); - dispatcher.afterAckMessages(null, finalPosition); } + notifyTheMarkDeletePositionChanged(oldMarkDeletePosition); IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); inProgressResetCursorFuture = null; future.complete(null);