This is an automated email from the ASF dual-hosted git repository.

lhotari pushed a commit to branch branch-4.2
in repository https://gitbox.apache.org/repos/asf/pulsar.git

commit ba906cbaff5fb449f38cb4fd8e7ef765310d535f
Author: 道君- Tao Jiuming <[email protected]>
AuthorDate: Tue Jun 23 15:26:13 2026 +0800

    [improve][broker] Trim orphaned bucket snapshots when ledgers are deleted 
(#25984)
    
    (cherry picked from commit e45b425f927c63b4de4093ee6f9f0218b188b701)
---
 .../bucket/BucketDelayedDeliveryTracker.java       | 110 +++++++++--
 .../bucket/BucketDelayedDeliveryTrackerTest.java   | 213 +++++++++++++++++++++
 2 files changed, 312 insertions(+), 11 deletions(-)

diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java
index 69edf70dad7..7a1b95a3a4a 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java
@@ -38,6 +38,7 @@ import java.util.NavigableSet;
 import java.util.Optional;
 import java.util.TreeSet;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.TimeUnit;
@@ -48,6 +49,7 @@ import javax.annotation.concurrent.ThreadSafe;
 import lombok.Getter;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.bookkeeper.mledger.ManagedCursor;
+import org.apache.bookkeeper.mledger.ManagedLedger;
 import org.apache.bookkeeper.mledger.Position;
 import org.apache.bookkeeper.mledger.PositionFactory;
 import org.apache.commons.collections4.CollectionUtils;
@@ -112,6 +114,8 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
 
     private CompletableFuture<Void> pendingLoad = null;
 
+    private volatile CompletableFuture<Void> trimFuture;
+
     public 
BucketDelayedDeliveryTracker(AbstractPersistentDispatcherMultipleConsumers 
dispatcher,
                                         Timer timer, long tickTimeMillis,
                                         boolean 
isDelayedDeliveryDeliverAtTimeStrict,
@@ -387,8 +391,15 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
             afterCreateImmutableBucket(immutableBucketDelayedIndexPair, 
createStartTime);
             lastMutableBucket.resetLastMutableBucketRange();
 
-            if (maxNumBuckets > 0 && immutableBuckets.asMapOfRanges().size() > 
maxNumBuckets) {
-                asyncMergeBucketSnapshot();
+            if (maxNumBuckets > 0 && immutableBuckets.asMapOfRanges().size() > 
maxNumBuckets
+                    && (trimFuture == null || trimFuture.isDone())) {
+                trimFuture = asyncTrimImmutableBuckets()
+                        .thenCompose(ignore -> asyncMergeBucketSnapshot())
+                        .whenComplete((ignore, t) -> {
+                            if (t != null) {
+                                log.warn("Failed to trim or merge bucket 
snapshots", t);
+                            }
+                        });
             }
         }
 
@@ -452,6 +463,10 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
 
     private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot() {
         List<ImmutableBucket> immutableBucketList = 
immutableBuckets.asMapOfRanges().values().stream().toList();
+        if (maxNumBuckets <= 0 || immutableBucketList.size() <= maxNumBuckets) 
{
+            return CompletableFuture.completedFuture(null);
+        }
+
         List<ImmutableBucket> toBeMergeImmutableBuckets = 
selectMergedBuckets(immutableBucketList, MAX_MERGE_NUM);
 
         if (toBeMergeImmutableBuckets.isEmpty()) {
@@ -605,6 +620,7 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
         }
 
         long cutoffTime = getCutoffTime();
+        Long firstLiveLedgerId = firstActiveLedgerId();
 
         lastMutableBucket.moveScheduledMessageToSharedQueue(cutoffTime, 
sharedBucketPriorityQueue);
 
@@ -613,13 +629,19 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
 
         while (n > 0 && !sharedBucketPriorityQueue.isEmpty()) {
             long timestamp = sharedBucketPriorityQueue.peekN1();
+            long ledgerId = sharedBucketPriorityQueue.peekN2();
+            long entryId = sharedBucketPriorityQueue.peekN3();
+            if (firstLiveLedgerId != null && ledgerId < firstLiveLedgerId) {
+                sharedBucketPriorityQueue.pop();
+                if (removeIndexBit(ledgerId, entryId)) {
+                    numberDelayedMessages.decrementAndGet();
+                }
+                continue;
+            }
             if (timestamp > cutoffTime) {
                 break;
             }
 
-            long ledgerId = sharedBucketPriorityQueue.peekN2();
-            long entryId = sharedBucketPriorityQueue.peekN3();
-
             SnapshotKey snapshotKey = new SnapshotKey(ledgerId, entryId);
 
             ImmutableBucket bucket = 
snapshotSegmentLastIndexMap.get(snapshotKey);
@@ -719,12 +741,26 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
 
     @Override
     public synchronized CompletableFuture<Void> clear() {
-        CompletableFuture<Void> future = cleanImmutableBuckets();
-        sharedBucketPriorityQueue.clear();
-        lastMutableBucket.clear();
-        snapshotSegmentLastIndexMap.clear();
-        numberDelayedMessages.set(0);
-        return future;
+        // Wait for any in-flight trim+merge to settle, then clear.
+        // Reuse trimFuture to block new triggers until the clear chain 
completes.
+        CompletableFuture<Void> before = trimFuture != null && 
!trimFuture.isDone()
+                ? trimFuture : CompletableFuture.completedFuture(null);
+        trimFuture = before
+                .exceptionally(t -> {
+                    log.warn("Trim/merge buckets failed, but still clear", t);
+                    return null;
+                })
+                .thenCompose(__ -> {
+                    synchronized (BucketDelayedDeliveryTracker.this) {
+                        CompletableFuture<Void> future = 
cleanImmutableBuckets();
+                        sharedBucketPriorityQueue.clear();
+                        lastMutableBucket.clear();
+                        snapshotSegmentLastIndexMap.clear();
+                        numberDelayedMessages.set(0);
+                        return future;
+                    }
+                });
+        return trimFuture;
     }
 
     @Override
@@ -791,4 +827,56 @@ public class BucketDelayedDeliveryTracker extends 
AbstractDelayedDeliveryTracker
         stats.recordBucketSnapshotSizeBytes(totalSnapshotLength.longValue());
         return stats.genTopicMetricMap();
     }
+
+    /**
+     * Delete orphaned bucket snapshots whose ledger range lies entirely 
before the earliest
+     * surviving ledger. Buckets are deleted sequentially; the chain stops on 
first failure
+     * to avoid wasted work when storage is unavailable.
+     */
+    private synchronized CompletableFuture<Void> asyncTrimImmutableBuckets() {
+        Long firstLedgerId = firstActiveLedgerId();
+        if (null == firstLedgerId) {
+            return CompletableFuture.completedFuture(null);
+        }
+        ManagedLedger ledger = dispatcher.getCursor().getManagedLedger();
+
+        Map<Range<Long>, ImmutableBucket> toBeDeletedBuckets =
+                new 
HashMap<>(immutableBuckets.subRangeMap(Range.lessThan(firstLedgerId)).asMapOfRanges());
+
+        if (toBeDeletedBuckets.isEmpty()) {
+            return CompletableFuture.completedFuture(null);
+        }
+
+        String ledgerName = ledger.getName();
+        CompletableFuture<Void> chain = 
CompletableFuture.completedFuture(null);
+        for (Map.Entry<Range<Long>, ImmutableBucket> entry : 
toBeDeletedBuckets.entrySet()) {
+            chain = chain.thenCompose(__ ->
+                    deleteBucketSnapshot(ledgerName, entry.getKey(), 
entry.getValue()));
+        }
+        return chain;
+    }
+
+    private CompletableFuture<Void> deleteBucketSnapshot(String ledgerName,
+                                                          Range<Long> range, 
ImmutableBucket bucket) {
+        return bucket.asyncDeleteBucketSnapshot(stats)
+                .handle((__, t) -> {
+                    if (t != null) {
+                        log.warn("Failed to delete bucket snapshot, 
LedgerName: {}, BucketKey: {}",
+                                ledgerName, bucket.bucketKey());
+                        throw new CompletionException(t);
+                    }
+                    synchronized (this) {
+                        snapshotSegmentLastIndexMap.entrySet().removeIf(entry 
-> entry.getValue() == bucket);
+                        immutableBuckets.remove(range);
+                        
numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages());
+                    }
+                    return null;
+                });
+    }
+
+    private Long firstActiveLedgerId() {
+        ManagedCursor cursor = dispatcher.getCursor();
+        Position mdp = cursor.getMarkDeletedPosition();
+        return mdp == null ? null : mdp.getLedgerId();
+    }
 }
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java
index 6ff98fa7f70..ee0ec0d9294 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java
@@ -40,13 +40,16 @@ import java.util.NavigableSet;
 import java.util.Set;
 import java.util.TreeMap;
 import java.util.TreeSet;
+import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicLong;
 import org.apache.bookkeeper.mledger.ManagedCursor;
+import org.apache.bookkeeper.mledger.ManagedLedger;
 import org.apache.bookkeeper.mledger.Position;
 import org.apache.bookkeeper.mledger.PositionFactory;
+import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo.LedgerInfo;
 import org.apache.commons.lang3.mutable.MutableLong;
 import org.apache.pulsar.broker.delayed.AbstractDeliveryTrackerTest;
 import org.apache.pulsar.broker.delayed.MockBucketSnapshotStorage;
@@ -463,4 +466,214 @@ public class BucketDelayedDeliveryTrackerTest extends 
AbstractDeliveryTrackerTes
 
       tracker.close();
     }
+
+    private static class TrackerWithStorage {
+        final BucketDelayedDeliveryTracker tracker;
+        final MockBucketSnapshotStorage storage;
+        final AtomicLong clockTime;
+
+        TrackerWithStorage(BucketDelayedDeliveryTracker tracker, 
MockBucketSnapshotStorage storage,
+                           AtomicLong clockTime) {
+            this.tracker = tracker;
+            this.storage = storage;
+            this.clockTime = clockTime;
+        }
+
+        void close() throws Exception {
+            tracker.close();
+            storage.close();
+        }
+    }
+
+    private static class BlockingDeleteStorage extends 
MockBucketSnapshotStorage {
+        final CompletableFuture<Void> firstDeleteFuture = new 
CompletableFuture<>();
+        final AtomicLong deleteCalls = new AtomicLong();
+
+        @Override
+        public CompletableFuture<Void> deleteBucketSnapshot(long bucketId) {
+            if (deleteCalls.incrementAndGet() <= 4) {
+                return firstDeleteFuture;
+            }
+            return super.deleteBucketSnapshot(bucketId);
+        }
+    }
+
+    private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, 
int maxNumBuckets)
+            throws Exception {
+        return createTrackerWithMockLedger(firstLedgerId, maxNumBuckets, new 
MockBucketSnapshotStorage());
+    }
+
+    private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, 
int maxNumBuckets,
+                                                          
MockBucketSnapshotStorage storage)
+            throws Exception {
+        storage.start();
+
+        ManagedLedger mockLedger = mock(ManagedLedger.class);
+        NavigableMap<Long, LedgerInfo> ledgerInfo = new TreeMap<>();
+        ledgerInfo.put(firstLedgerId, mock(LedgerInfo.class));
+        when(mockLedger.getLedgersInfo()).thenReturn(ledgerInfo);
+        when(mockLedger.getName()).thenReturn("test-ledger");
+
+        ManagedCursor mockCursor = new MockManagedCursor("test-cursor") {
+            @Override
+            public ManagedLedger getManagedLedger() {
+                return mockLedger;
+            }
+
+            @Override
+            public Position getMarkDeletedPosition() {
+                return PositionFactory.create(firstLedgerId, -1);
+            }
+        };
+
+        AbstractPersistentDispatcherMultipleConsumers disp =
+                mock(AbstractPersistentDispatcherMultipleConsumers.class);
+        Clock mockClock = mock(Clock.class);
+        AtomicLong mockClockTime = new AtomicLong();
+        when(mockClock.millis()).then(x -> mockClockTime.get());
+        doReturn(mockCursor).when(disp).getCursor();
+        doReturn("persistent://public/default/testDelay" + " / " + 
mockCursor.getName()).when(disp).getName();
+
+        BucketDelayedDeliveryTracker tracker = new 
BucketDelayedDeliveryTracker(disp, mock(Timer.class),
+                100000, mockClock, true, storage, 5, 
TimeUnit.MILLISECONDS.toMillis(10), -1, maxNumBuckets);
+        return new TrackerWithStorage(tracker, storage, mockClockTime);
+    }
+
+    @Test
+    public void testTrimRemovesOrphanedBuckets() throws Exception {
+        long firstLedgerId = 31L;
+        int messageCount = 36;
+        TrackerWithStorage ts = createTrackerWithMockLedger(firstLedgerId, 5);
+
+        for (int i = 1; i <= messageCount; i++) {
+            ts.tracker.addMessage(i, i, i * 10);
+        }
+        Awaitility.await().untilAsserted(() ->
+                
Assert.assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream()
+                        .noneMatch(x -> x.merging)));
+
+        int bucketCount = 
ts.tracker.getImmutableBuckets().asMapOfRanges().size();
+        assertTrue(bucketCount <= 5,
+                "Bucket count " + bucketCount + " should be <= maxNumBuckets=5 
after trim+merge");
+
+        ts.tracker.getImmutableBuckets().asMapOfRanges().forEach((range, 
bucket) ->
+                assertTrue(range.lowerEndpoint() >= firstLedgerId,
+                        "Remaining bucket range " + range + " should be >= " + 
firstLedgerId));
+
+        long messagesAfterTrim = ts.tracker.getNumberOfDelayedMessages();
+        ts.clockTime.set(messageCount * 10);
+        NavigableSet<Position> scheduledMessages = 
ts.tracker.getScheduledMessages(1);
+        assertTrue(scheduledMessages.stream().noneMatch(position -> 
position.getLedgerId() < firstLedgerId),
+                "Trimmed ledgers should not be returned from the loaded shared 
queue");
+        assertEquals(ts.tracker.getNumberOfDelayedMessages(), 
messagesAfterTrim - scheduledMessages.size());
+
+        ts.close();
+    }
+
+    @Test
+    public void testTrimHandlesDeleteFailure() throws Exception {
+        long firstLedgerId = 50L;
+        int messageCount = 31;
+        TrackerWithStorage ts = createTrackerWithMockLedger(firstLedgerId, 5);
+
+        // MaxRetryTimes=3 means the first trim delete attempt plus 3 retries 
= 4 exceptions consumed.
+        for (int i = 0; i < 4; i++) {
+            ts.storage.injectDeleteException(
+                    new BucketSnapshotPersistenceException("Delete failed"));
+        }
+
+        for (int i = 1; i <= messageCount; i++) {
+            ts.tracker.addMessage(i, i, i * 10);
+        }
+        Awaitility.await().untilAsserted(() ->
+                
Assert.assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream()
+                        .noneMatch(x -> x.merging)));
+
+        Awaitility.await().untilAsserted(() ->
+                assertTrue(ts.storage.deleteExceptionQueue.isEmpty(),
+                        "Delete exception should have been consumed"));
+
+        // Trim failed on the first orphaned bucket; the sequential chain 
stopped, so all
+        // 6 orphaned buckets remain in immutableBuckets.
+        assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().size() > 0,
+                "Orphaned buckets should remain when trim delete fails");
+        ts.tracker.getImmutableBuckets().asMapOfRanges().forEach((range, 
bucket) ->
+                assertTrue(range.upperEndpoint() < firstLedgerId,
+                        "Remaining bucket " + range + " should be an orphaned 
bucket"));
+
+        // numberDelayedMessages is unchanged because failed deletes do not 
decrement the count.
+        assertEquals(ts.tracker.getNumberOfDelayedMessages(), messageCount);
+
+        ts.close();
+    }
+
+    @Test
+    public void testClearRunsAfterInFlightTrimFailure() throws Exception {
+        long firstLedgerId = 50L;
+        int messageCount = 31;
+        BlockingDeleteStorage storage = new BlockingDeleteStorage();
+        TrackerWithStorage ts = createTrackerWithMockLedger(firstLedgerId, 5, 
storage);
+
+        for (int i = 1; i <= messageCount; i++) {
+            ts.tracker.addMessage(i, i, i * 10);
+        }
+        Awaitility.await().untilAsserted(() ->
+                assertTrue(storage.deleteCalls.get() > 0, "Trim delete should 
be in flight"));
+
+        CompletableFuture<Void> clearFuture = ts.tracker.clear();
+        storage.firstDeleteFuture.completeExceptionally(new 
BucketSnapshotPersistenceException("Delete failed"));
+
+        clearFuture.get(1, TimeUnit.MINUTES);
+        assertEquals(ts.tracker.getNumberOfDelayedMessages(), 0);
+        assertEquals(ts.tracker.getImmutableBuckets().asMapOfRanges().size(), 
0);
+        assertEquals(ts.tracker.getLastMutableBucket().size(), 0);
+        assertEquals(ts.tracker.getSharedBucketPriorityQueue().size(), 0);
+
+        ts.close();
+    }
+
+    @Test
+    public void testTrimWithNoOrphanedBuckets() throws Exception {
+        TrackerWithStorage ts = createTrackerWithMockLedger(0L, 5);
+
+        for (int i = 1; i <= 31; i++) {
+            ts.tracker.addMessage(i, i, i * 10);
+        }
+        Awaitility.await().untilAsserted(() ->
+                
Assert.assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream()
+                        .noneMatch(x -> x.merging)));
+
+        int bucketCount = 
ts.tracker.getImmutableBuckets().asMapOfRanges().size();
+        assertTrue(bucketCount <= 5,
+                "Bucket count " + bucketCount + " should be <= 
maxNumBuckets=5");
+        assertTrue(bucketCount > 0, "Should have at least one bucket after 
merge");
+
+        ts.close();
+    }
+
+    @Test
+    public void testMergeEarlyReturnWhenWithinLimit() throws Exception {
+        TrackerWithStorage ts = createTrackerWithMockLedger(0L, 50);
+
+        for (int i = 1; i <= 30; i++) {
+            ts.tracker.addMessage(i, i, i * 10);
+        }
+        Awaitility.await().untilAsserted(() ->
+                
Assert.assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream()
+                        .noneMatch(x -> x.merging)));
+
+        int bucketCount = 
ts.tracker.getImmutableBuckets().asMapOfRanges().size();
+        assertTrue(bucketCount < 50,
+                "Bucket count " + bucketCount + " should be well below 
maxNumBuckets=50");
+
+        long msgsBefore = ts.tracker.getNumberOfDelayedMessages();
+        ts.tracker.addMessage(200, 200, 200 * 10);
+        Awaitility.await().untilAsserted(() ->
+                
Assert.assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream()
+                        .noneMatch(x -> x.merging)));
+
+        assertEquals(ts.tracker.getNumberOfDelayedMessages(), msgsBefore + 1);
+
+        ts.close();
+    }
 }

Reply via email to