From 7f22e5653731c108f7949e8d433d9858b61f5f0c Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Mon, 24 Aug 2026 12:56:19 +0800 Subject: [PATCH 1/3] [fix][broker] Manage in-flight bucket work during clear --- .../bucket/BucketDelayedDeliveryTracker.java | 70 +++- .../BucketDelayedDeliveryTrackerTest.java | 379 +++++++++++++++++- 2 files changed, 426 insertions(+), 23 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 bd7c596561cc5..f2c8cb2865c63 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 @@ -37,6 +37,7 @@ import java.util.List; 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.CompletionException; @@ -130,6 +131,8 @@ public static record SnapshotKey(long ledgerId, long entryId) {} private CompletableFuture pendingLoad = null; + private final Set> pendingDeletes = ConcurrentHashMap.newKeySet(); + private volatile CompletableFuture trimFuture; public BucketDelayedDeliveryTracker(AbstractPersistentDispatcherMultipleConsumers dispatcher, @@ -255,8 +258,7 @@ private synchronized long recoverBucketSnapshot() throws RecoverDelayedDeliveryT Range key = mapEntry.getKey(); ImmutableBucket immutableBucket = mapEntry.getValue(); removeBucket(key); - // delete asynchronously without waiting for completion - immutableBucket.asyncDeleteBucketSnapshot(stats); + trackDelete(immutableBucket.asyncDeleteBucketSnapshot(stats)); } long totalLength = 0; @@ -572,7 +574,7 @@ private synchronized CompletableFuture asyncMergeBucketSnapshot(List { + .thenCompose(combinedDelayedIndexQueue -> { synchronized (BucketDelayedDeliveryTracker.this) { long createStartTime = System.currentTimeMillis(); stats.recordTriggerEvent(BucketDelayedMessageIndexStats.Type.create); @@ -604,17 +606,16 @@ private synchronized CompletableFuture asyncMergeBucketSnapshot(List { - List> removeFutures = - buckets.stream().map(bucket -> bucket.asyncDeleteBucketSnapshot(stats)) - .toList(); - return FutureUtil.waitForAll(removeFutures); - }); - for (ImmutableBucket bucket : buckets) { removeBucket(Range.closed(bucket.getStartLedgerId(), bucket.getEndLedgerId())); } + + return immutableBucketDelayedIndexPair.getLeft().getSnapshotCreateFuture() + .orElse(NULL_LONG_PROMISE) + .thenCompose(___ -> FutureUtil.waitForAll( + buckets.stream() + .map(bucket -> bucket.asyncDeleteBucketSnapshot(stats)) + .toList())); } }); }); @@ -715,13 +716,12 @@ public synchronized NavigableSet getScheduledMessages(int maxMessages) long loadStartTime = System.currentTimeMillis(); stats.recordTriggerEvent(BucketDelayedMessageIndexStats.Type.load); CompletableFuture loadFuture = pendingLoad = bucket.asyncLoadNextBucketSnapshotEntry() - .thenAccept(indexList -> { + .thenCompose(indexList -> { synchronized (BucketDelayedDeliveryTracker.this) { this.snapshotSegmentLastIndexMap.remove(snapshotKey); if (CollectionUtils.isEmpty(indexList)) { removeBucket(Range.closed(bucket.getStartLedgerId(), bucket.getEndLedgerId())); - bucket.asyncDeleteBucketSnapshot(stats); - return; + return bucket.asyncDeleteBucketSnapshot(stats); } DelayedIndex lastDelayedIndex = indexList.get(indexList.size() - 1); @@ -732,6 +732,7 @@ public synchronized NavigableSet getScheduledMessages(int maxMessages) sharedBucketPriorityQueue.add(index.getTimestamp(), index.getLedgerId(), index.getEntryId()); } + return CompletableFuture.completedFuture(null); } }).whenComplete((__, ex) -> { if (ex != null) { @@ -786,6 +787,17 @@ private synchronized boolean checkPendingLoadDone() { return false; } + private void trackDelete(CompletableFuture deleteFuture) { + synchronized (this) { + pendingDeletes.add(deleteFuture); + } + deleteFuture.whenComplete((__, ex) -> { + synchronized (BucketDelayedDeliveryTracker.this) { + pendingDeletes.remove(deleteFuture); + } + }); + } + @Override public boolean shouldPauseAllDeliveries() { return false; @@ -793,8 +805,7 @@ public boolean shouldPauseAllDeliveries() { @Override public synchronized CompletableFuture clear() { - // Wait for any in-flight trim+merge to settle, then clear. - // Reuse trimFuture to block new triggers until the clear chain completes. + // Wait for in-flight trim/merge and pending load/delete work before resetting state. CompletableFuture before = trimFuture != null && !trimFuture.isDone() ? trimFuture : CompletableFuture.completedFuture(null); trimFuture = before @@ -803,14 +814,29 @@ public synchronized CompletableFuture clear() { return null; }) .thenCompose(__ -> { + List> pending = new ArrayList<>(); synchronized (BucketDelayedDeliveryTracker.this) { - CompletableFuture future = cleanImmutableBuckets(); - sharedBucketPriorityQueue.clear(); - index.clear(); - lastMutableBucket.clear(); - snapshotSegmentLastIndexMap.clear(); - return future; + if (pendingLoad != null) { + pending.add(pendingLoad); + } + pending.addAll(pendingDeletes); } + return FutureUtil.waitForAll(pending) + .exceptionally(t -> { + log.warn().exception(t) + .log("Failed to wait for pending delayed delivery work, but still clear"); + return null; + }) + .thenCompose(ignore -> { + synchronized (BucketDelayedDeliveryTracker.this) { + CompletableFuture future = cleanImmutableBuckets(); + sharedBucketPriorityQueue.clear(); + index.clear(); + lastMutableBucket.clear(); + snapshotSegmentLastIndexMap.clear(); + return future; + } + }); }); return trimFuture; } 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 54140cb2bfdf7..2ca0f41481aea 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 @@ -18,6 +18,7 @@ */ package org.apache.pulsar.broker.delayed.bucket; +import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.Mockito.doReturn; @@ -35,6 +36,7 @@ import java.lang.reflect.Method; import java.nio.ByteBuffer; import java.time.Clock; +import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; import java.util.List; @@ -48,6 +50,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedCursor; @@ -62,6 +65,7 @@ import org.apache.pulsar.broker.delayed.proto.SnapshotMetadata; import org.apache.pulsar.broker.delayed.proto.SnapshotSegment; import org.apache.pulsar.broker.service.persistent.AbstractPersistentDispatcherMultipleConsumers; +import org.apache.pulsar.common.util.FutureUtil; import org.awaitility.Awaitility; import org.roaringbitmap.RoaringBitmap; import org.roaringbitmap.buffer.ImmutableRoaringBitmap; @@ -623,6 +627,112 @@ public CompletableFuture createBucketSnapshot(SnapshotMetadata snapshotMet } } + private static class GatedSegmentLoadStorage extends MockBucketSnapshotStorage { + volatile CompletableFuture segmentLoadGate; + + @Override + public CompletableFuture> getBucketSnapshotSegment(long bucketId, + long firstSegmentEntryId, + long lastSegmentEntryId) { + CompletableFuture> future = + super.getBucketSnapshotSegment(bucketId, firstSegmentEntryId, lastSegmentEntryId); + CompletableFuture gate = segmentLoadGate; + if (gate == null) { + return future; + } + return gate.thenCompose(__ -> future); + } + } + + private static class GatedMergeLoadStorage extends MockBucketSnapshotStorage { + final CompletableFuture mergeLoadGate = new CompletableFuture<>(); + final AtomicLong mergeLoadCalls = new AtomicLong(); + + @Override + public CompletableFuture> getBucketSnapshotSegment(long bucketId, + long firstSegmentEntryId, + long lastSegmentEntryId) { + mergeLoadCalls.incrementAndGet(); + return mergeLoadGate.thenCompose(__ -> + super.getBucketSnapshotSegment(bucketId, firstSegmentEntryId, lastSegmentEntryId)); + } + } + + private static class GatedMergeCreateStorage extends MockBucketSnapshotStorage { + final CompletableFuture mergeCreateGate = new CompletableFuture<>(); + final AtomicLong createCalls = new AtomicLong(); + + @Override + public CompletableFuture createBucketSnapshot(SnapshotMetadata snapshotMetadata, + List bucketSnapshotSegments, + String bucketKey, String topicName, String cursorName) { + if (createCalls.incrementAndGet() <= 2) { + return super.createBucketSnapshot(snapshotMetadata, bucketSnapshotSegments, bucketKey, + topicName, cursorName); + } + return mergeCreateGate.thenCompose(__ -> + super.createBucketSnapshot(snapshotMetadata, bucketSnapshotSegments, bucketKey, + topicName, cursorName)); + } + } + + private static class GatedDeleteStorage extends MockBucketSnapshotStorage { + final CompletableFuture deleteGate = new CompletableFuture<>(); + final AtomicLong deleteCalls = new AtomicLong(); + + @Override + public CompletableFuture deleteBucketSnapshot(long bucketId) { + if (deleteCalls.incrementAndGet() <= 2) { + return deleteGate; + } + return super.deleteBucketSnapshot(bucketId); + } + } + + private static class FailingMergeDeleteStorage extends MockBucketSnapshotStorage { + final AtomicBoolean firstMergeDeleteStarted = new AtomicBoolean(); + private final AtomicLong failedBucketId = new AtomicLong(-1); + + @Override + public CompletableFuture deleteBucketSnapshot(long bucketId) { + if (firstMergeDeleteStarted.compareAndSet(false, true)) { + failedBucketId.set(bucketId); + } + if (bucketId == failedBucketId.get()) { + return FutureUtil.failedFuture(new BucketSnapshotPersistenceException("Merge delete failed")); + } + return super.deleteBucketSnapshot(bucketId); + } + } + + /** + * Fails every load of snapshot segments after the first one, and gates the first snapshot delete. + */ + private static class FailingSegmentLoadGatedDeleteStorage extends MockBucketSnapshotStorage { + final CompletableFuture firstDeleteGate = new CompletableFuture<>(); + final AtomicLong failedSegmentLoadCalls = new AtomicLong(); + final AtomicLong deleteCalls = new AtomicLong(); + + @Override + public CompletableFuture> getBucketSnapshotSegment(long bucketId, + long firstSegmentEntryId, + long lastSegmentEntryId) { + if (firstSegmentEntryId >= 2) { + failedSegmentLoadCalls.incrementAndGet(); + return FutureUtil.failedFuture(new BucketSnapshotPersistenceException("Load failed")); + } + return super.getBucketSnapshotSegment(bucketId, firstSegmentEntryId, lastSegmentEntryId); + } + + @Override + public CompletableFuture deleteBucketSnapshot(long bucketId) { + if (deleteCalls.incrementAndGet() == 1) { + return firstDeleteGate; + } + return super.deleteBucketSnapshot(bucketId); + } + } + private ImmutableBucket createMergeableBucket(TrackerWithStorage trackerWithStorage, long startLedgerId, long endLedgerId, List firstScheduleTimestamps) { ImmutableBucket bucket = new ImmutableBucket(trackerWithStorage.tracker.getCtx(), startLedgerId, endLedgerId); @@ -641,6 +751,13 @@ private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, int m private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, int maxNumBuckets, MockBucketSnapshotStorage storage) throws Exception { + return createTrackerWithMockLedger(firstLedgerId, maxNumBuckets, storage, -1); + } + + private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, int maxNumBuckets, + MockBucketSnapshotStorage storage, + int maxIndexesPerSegment) + throws Exception { storage.start(); ManagedLedger mockLedger = mock(ManagedLedger.class); @@ -670,7 +787,8 @@ public Position getMarkDeletedPosition() { 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); + 100000, mockClock, true, storage, 5, TimeUnit.MILLISECONDS.toMillis(10), + maxIndexesPerSegment, maxNumBuckets); return new TrackerWithStorage(tracker, storage, mockClockTime); } @@ -779,6 +897,265 @@ public void testTrimWaitsForInFlightSnapshotCreation() throws Exception { } } + @Test + public void testClearWaitsForInFlightSegmentLoad() throws Exception { + AbstractPersistentDispatcherMultipleConsumers testDispatcher = + mock(AbstractPersistentDispatcherMultipleConsumers.class); + Clock testClock = mock(Clock.class); + AtomicLong testClockTime = new AtomicLong(); + when(testClock.millis()).then(x -> testClockTime.get()); + + GatedSegmentLoadStorage storage = new GatedSegmentLoadStorage(); + storage.start(); + + ManagedCursor cursor = new MockManagedCursor("test_clear_load_cursor"); + doReturn(cursor).when(testDispatcher).getCursor(); + doReturn("persistent://public/default/testClearLoad / " + cursor.getName()) + .when(testDispatcher).getName(); + + BucketDelayedDeliveryTracker tracker = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 1000, testClock, true, storage, + 4, TimeUnit.MILLISECONDS.toMillis(10), 2, 50); + try { + // Two indexes per segment: reaching the first segment boundary triggers a load of the + // next snapshot segment. + for (int i = 1; i <= 6; i++) { + tracker.addMessage(i, i, i * 100); + } + Awaitility.await().untilAsserted(() -> + assertTrue(tracker.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + + storage.segmentLoadGate = new CompletableFuture<>(); + testClockTime.set(600); + tracker.getScheduledMessages(10); + + CompletableFuture clearFuture = tracker.clear(); + assertFalse("clear() should wait for the in-flight segment load", clearFuture.isDone()); + + storage.segmentLoadGate.complete(null); + + assertThat(clearFuture).succeedsWithin(Duration.ofSeconds(3)); + + assertEquals(tracker.getNumberOfDelayedMessages(), 0); + assertEquals(tracker.getImmutableBuckets().asMapOfRanges().size(), 0); + assertEquals(tracker.getLastMutableBucket().size(), 0); + assertEquals(tracker.getSharedBucketPriorityQueue().size(), 0); + } finally { + tracker.close(); + storage.clean(); + } + } + + @Test + public void testClearStillWaitsForPendingDeletesWhenLoadFails() throws Exception { + FailingSegmentLoadGatedDeleteStorage storage = new FailingSegmentLoadGatedDeleteStorage(); + storage.start(); + + ManagedCursor cursor = new MockManagedCursor("test_failed_load_cursor"); + AbstractPersistentDispatcherMultipleConsumers testDispatcher = + mock(AbstractPersistentDispatcherMultipleConsumers.class); + Clock testClock = mock(Clock.class); + AtomicLong testClockTime = new AtomicLong(); + when(testClock.millis()).then(x -> testClockTime.get()); + doReturn(cursor).when(testDispatcher).getCursor(); + doReturn("persistent://public/default/testFailedLoad / " + cursor.getName()) + .when(testDispatcher).getName(); + + // Bucket [1..5] with already-expired timestamps and bucket [6..10] with future ones. + BucketDelayedDeliveryTracker producer = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + for (int i = 1; i <= 11; i++) { + producer.addMessage(i, i, i <= 5 ? 10L * i : 100L * i); + } + Awaitility.await().untilAsserted(() -> + assertTrue(producer.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + } finally { + producer.close(); + } + + // Recovery with cutoff 50: [1..5] is fully expired, so its recovery delete (gated) is an + // in-flight tracked delete; [6..10] recovers with its first segment loaded. + testClockTime.set(50); + BucketDelayedDeliveryTracker tracker = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + testClockTime.set(1000); + tracker.getScheduledMessages(10); + Awaitility.await().untilAsserted(() -> + assertEquals(storage.failedSegmentLoadCalls.get(), 4, + "The load should have failed after initial attempt plus retries")); + + CompletableFuture clearFuture = tracker.clear(); + assertFalse("clear() must still wait for the in-flight tracked delete although the " + + "pending load already failed", clearFuture.isDone()); + + storage.firstDeleteGate.complete(null); + + assertThat(clearFuture).succeedsWithin(Duration.ofSeconds(3)); + assertEquals(tracker.getNumberOfDelayedMessages(), 0); + assertEquals(tracker.getImmutableBuckets().asMapOfRanges().size(), 0); + } finally { + tracker.close(); + storage.clean(); + } + } + + @Test + public void testClearWaitsForTerminalSegmentLoadDelete() throws Exception { + GatedDeleteStorage storage = new GatedDeleteStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(0L, 50, storage, 2); + try { + for (int i = 1; i <= 6; i++) { + assertTrue(ts.tracker.addMessage(i, i, i * 100)); + } + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> + assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + + ts.clockTime.set(600); + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + ts.tracker.getScheduledMessages(10); + assertEquals(storage.deleteCalls.get(), 1, + "The terminal segment load should start the snapshot delete"); + }); + + CompletableFuture clearFuture = ts.tracker.clear(); + assertFalse("clear() must wait for the terminal delete chained to segment load", + clearFuture.isDone()); + + storage.deleteGate.complete(null); + assertThat(clearFuture).succeedsWithin(Duration.ofSeconds(3)); + assertEquals(ts.tracker.getImmutableBuckets().asMapOfRanges().size(), 0); + } finally { + ts.close(); + } + } + + @Test + public void testClearCompletesWhenTerminalSegmentLoadDeleteFails() throws Exception { + TrackerWithStorage ts = createTrackerWithMockLedger(0L, 50, new MockBucketSnapshotStorage(), 2); + try { + for (int i = 1; i <= 6; i++) { + assertTrue(ts.tracker.addMessage(i, i, 10L * i)); + } + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> + assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + + for (int i = 0; i < 4; i++) { + ts.storage.injectDeleteException( + new BucketSnapshotPersistenceException("Terminal delete failed")); + } + + ts.clockTime.set(1000); + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + ts.tracker.getScheduledMessages(10); + assertTrue(ts.storage.deleteExceptionQueue.isEmpty(), + "The terminal delete should have consumed all injected failures"); + }); + + assertThat(ts.tracker.clear()).succeedsWithin(Duration.ofSeconds(3)); + assertEquals(ts.tracker.getNumberOfDelayedMessages(), 0); + assertEquals(ts.tracker.getImmutableBuckets().asMapOfRanges().size(), 0); + assertEquals(ts.tracker.getSharedBucketPriorityQueue().size(), 0); + } finally { + ts.close(); + } + } + + @Test + public void testClearWaitsForMergeSourceBucketDelete() throws Exception { + GatedDeleteStorage storage = new GatedDeleteStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1, storage); + try { + for (int i = 1; i <= 11; i++) { + assertTrue(ts.tracker.addMessage(i, i, i * 10)); + } + + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted( + () -> assertEquals(storage.deleteCalls.get(), 2, + "Merge should start deleting both source buckets")); + + CompletableFuture clearFuture = ts.tracker.clear(); + assertFalse("clear() must wait for merge's in-flight source-bucket deletes", + clearFuture.isDone()); + + storage.deleteGate.complete(null); + assertThat(clearFuture).succeedsWithin(Duration.ofSeconds(3)); + assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().size() <= 1); + } finally { + ts.close(); + } + } + + @Test + public void testClearCompletesWhenMergeSourceBucketDeleteFails() throws Exception { + FailingMergeDeleteStorage storage = new FailingMergeDeleteStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1, storage); + try { + for (int i = 1; i <= 11; i++) { + assertTrue(ts.tracker.addMessage(i, i, i * 10)); + } + + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> + assertTrue(storage.firstMergeDeleteStarted.get(), + "The first merge source-bucket delete should have started")); + + assertThat(ts.tracker.clear()).succeedsWithin(Duration.ofSeconds(3)); + } finally { + ts.close(); + } + } + + @Test + public void testClearWaitsForMergeSegmentLoad() throws Exception { + GatedMergeLoadStorage storage = new GatedMergeLoadStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1, storage); + try { + for (int i = 1; i <= 11; i++) { + assertTrue(ts.tracker.addMessage(i, i, i * 10)); + } + + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted( + () -> assertEquals(storage.mergeLoadCalls.get(), 2)); + + CompletableFuture clearFuture = ts.tracker.clear(); + assertFalse(clearFuture.isDone()); + + storage.mergeLoadGate.complete(null); + assertThat(clearFuture).succeedsWithin(Duration.ofSeconds(3)); + } finally { + ts.close(); + } + } + + @Test + public void testClearWaitsForMergeSnapshotCreation() throws Exception { + GatedMergeCreateStorage storage = new GatedMergeCreateStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1, storage); + try { + for (int i = 1; i <= 11; i++) { + assertTrue(ts.tracker.addMessage(i, i, i * 10)); + } + + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted( + () -> assertEquals(storage.createCalls.get(), 3)); + + CompletableFuture clearFuture = ts.tracker.clear(); + assertFalse(clearFuture.isDone()); + + storage.mergeCreateGate.complete(null); + assertThat(clearFuture).succeedsWithin(Duration.ofSeconds(3)); + } finally { + ts.close(); + } + } + @Test public void testSelectMergedBucketsSupportsTwoBuckets() throws Exception { TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1); From e3c264aaddb4ded4f67afa3e1eaf206cca8f0e53 Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Mon, 24 Aug 2026 21:25:15 +0800 Subject: [PATCH 2/3] [fix][broker] Clean up stale bucket cursor properties on recovery The BucketNotExist branch in handleRecoverBucketSnapshotEntry never matched: dependent future stages wrap the checked exception in a CompletionException, so a stale cursor property pointing to a deleted snapshot failed the whole recovery instead of being cleaned up. Unwrap before the check so the bucket is deleted and the property is removed, like the path was designed to do. Also make MockBucketSnapshotStorage re-readable (parse from a slice, same as the metadata read) and add recreate-tracker tests for both cleanup-failure shapes: consistent snapshot+property residue loads normally, property-only residue is removed on recovery. --- .../bucket/BucketDelayedDeliveryTracker.java | 4 +- .../delayed/MockBucketSnapshotStorage.java | 4 +- .../BucketDelayedDeliveryTrackerTest.java | 152 ++++++++++++++++++ 3 files changed, 158 insertions(+), 2 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 f2c8cb2865c63..0655efc39208d 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 @@ -292,7 +292,9 @@ private CompletableFuture> handleRecoverBucketSnapshotEntry(I if (e == null) { f.complete(v); } else { - if (e instanceof BucketNotExistException) { + // Dependent stages wrap checked exceptions in CompletionException, so unwrap + // before matching the not-exist case. + if (FutureUtil.unwrapCompletionException(e) instanceof BucketNotExistException) { // If the bucket does not exist, return an empty list, // the bucket will be deleted from `immutableBuckets` in the next step. f.complete(Collections.emptyList()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockBucketSnapshotStorage.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockBucketSnapshotStorage.java index 28d25ca00366c..26cff4fca23b8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockBucketSnapshotStorage.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockBucketSnapshotStorage.java @@ -136,7 +136,9 @@ public CompletableFuture> getBucketSnapshotSegment(long bu for (int i = (int) firstSegmentEntryId; i <= lastEntryId; i++) { ByteBuf byteBuf = this.bucketSnapshots.get(bucketId).get(i); SnapshotSegment snapshotSegment = new SnapshotSegment(); - snapshotSegment.parseFrom(byteBuf, byteBuf.readableBytes()); + // Parse from a slice so the stored buffer stays re-readable, like a real storage. + ByteBuf slice = byteBuf.slice(); + snapshotSegment.parseFrom(slice, slice.readableBytes()); snapshotSegments.add(snapshotSegment); } return snapshotSegments; 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 2ca0f41481aea..05dd0f46ecbfd 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 @@ -55,6 +55,7 @@ import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionFactory; import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo.LedgerInfo; @@ -733,6 +734,22 @@ public CompletableFuture deleteBucketSnapshot(long bucketId) { } } + private static class MissingBucketAsNotExistStorage extends MockBucketSnapshotStorage { + @Override + public CompletableFuture getBucketSnapshotMetadata(long bucketId) { + CompletableFuture future = new CompletableFuture<>(); + super.getBucketSnapshotMetadata(bucketId).whenComplete((metadata, ex) -> { + if (ex != null) { + future.completeExceptionally( + new BucketNotExistException("Bucket " + bucketId + " does not exist")); + } else { + future.complete(metadata); + } + }); + return future; + } + } + private ImmutableBucket createMergeableBucket(TrackerWithStorage trackerWithStorage, long startLedgerId, long endLedgerId, List firstScheduleTimestamps) { ImmutableBucket bucket = new ImmutableBucket(trackerWithStorage.tracker.getCtx(), startLedgerId, endLedgerId); @@ -1068,6 +1085,141 @@ public void testClearCompletesWhenTerminalSegmentLoadDeleteFails() throws Except } } + @Test + public void testRecoveryLoadsBucketWhenTerminalDeleteFailed() throws Exception { + MockBucketSnapshotStorage storage = new MockBucketSnapshotStorage(); + storage.start(); + + ManagedCursor cursor = new MockManagedCursor("test_residual_snapshot_cursor"); + AbstractPersistentDispatcherMultipleConsumers testDispatcher = + mock(AbstractPersistentDispatcherMultipleConsumers.class); + Clock testClock = mock(Clock.class); + AtomicLong testClockTime = new AtomicLong(); + when(testClock.millis()).then(x -> testClockTime.get()); + doReturn(cursor).when(testDispatcher).getCursor(); + doReturn("persistent://public/default/testResidualSnapshot / " + cursor.getName()) + .when(testDispatcher).getName(); + + BucketDelayedDeliveryTracker producer = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + for (int i = 1; i <= 6; i++) { + producer.addMessage(i, i, 10L * i); + } + Awaitility.await().untilAsserted(() -> + assertTrue(producer.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + } finally { + producer.close(); + } + + // The terminal delete exhausts its retries, so the snapshot and the cursor property survive + // together as a consistent pair. + BucketDelayedDeliveryTracker tracker = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + for (int i = 0; i < 4; i++) { + storage.injectDeleteException( + new BucketSnapshotPersistenceException("Terminal delete failed")); + } + testClockTime.set(1000); + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + tracker.getScheduledMessages(10); + assertTrue(storage.deleteExceptionQueue.isEmpty()); + }); + } finally { + tracker.close(); + } + + // A tracker recreated with the same cursor/storage recovers the bucket normally. + testClockTime.set(0); + BucketDelayedDeliveryTracker recovered = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + assertTrue(recovered.getImmutableBuckets().asMapOfRanges() + .containsKey(Range.closed(1L, 5L))); + assertEquals(recovered.getNumberOfDelayedMessages(), 5); + } finally { + recovered.close(); + } + storage.close(); + } + + @Test + public void testRecoveryRemovesResidualPropertyWhenSnapshotDeleted() throws Exception { + MissingBucketAsNotExistStorage storage = new MissingBucketAsNotExistStorage(); + storage.start(); + + AtomicLong propertyRemoveCalls = new AtomicLong(); + ManagedCursor cursor = new MockManagedCursor("test_residual_property_cursor") { + @Override + public CompletableFuture removeCursorProperty(String key) { + if (propertyRemoveCalls.incrementAndGet() <= 4) { + return FutureUtil.failedFuture( + new ManagedLedgerException.BadVersionException("version conflict")); + } + return super.removeCursorProperty(key); + } + }; + AbstractPersistentDispatcherMultipleConsumers testDispatcher = + mock(AbstractPersistentDispatcherMultipleConsumers.class); + Clock testClock = mock(Clock.class); + AtomicLong testClockTime = new AtomicLong(); + when(testClock.millis()).then(x -> testClockTime.get()); + doReturn(cursor).when(testDispatcher).getCursor(); + doReturn("persistent://public/default/testResidualProperty / " + cursor.getName()) + .when(testDispatcher).getName(); + + BucketDelayedDeliveryTracker producer = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + for (int i = 1; i <= 6; i++) { + producer.addMessage(i, i, 10L * i); + } + Awaitility.await().untilAsserted(() -> + assertTrue(producer.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + } finally { + producer.close(); + } + + // The snapshot delete succeeds but every property-removal attempt hits BadVersion, so only + // the cursor property is left behind. + String bucketKey = BucketDelayedDeliveryTracker.DELAYED_BUCKET_KEY_PREFIX + "_1_5"; + BucketDelayedDeliveryTracker tracker = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + testClockTime.set(1000); + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + tracker.getScheduledMessages(10); + assertTrue(propertyRemoveCalls.get() >= 4); + }); + assertTrue(cursor.getCursorProperties().containsKey(bucketKey), + "The property should survive the failed removal"); + } finally { + tracker.close(); + } + + // A tracker recreated with the same cursor/storage resolves the stale property, hits + // BucketNotExist, and removes it. + BucketDelayedDeliveryTracker recovered = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> + assertFalse(cursor.getCursorProperties().containsKey(bucketKey))); + assertEquals(recovered.getImmutableBuckets().asMapOfRanges().size(), 0); + } finally { + recovered.close(); + } + storage.close(); + } + @Test public void testClearWaitsForMergeSourceBucketDelete() throws Exception { GatedDeleteStorage storage = new GatedDeleteStorage(); From 97817549a80cb5a9c38d99c1a04f34b73def11f9 Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Mon, 24 Aug 2026 21:54:20 +0800 Subject: [PATCH 3/3] [fix][broker] Keep merge source snapshots when the merged creation fails A failed merged-bucket creation completes with the INVALID_BUCKET_ID sentinel, so the source-delete chain still ran and deleted the source snapshots although their ledger ranges were already removed from immutableBuckets - after a restart the only durable delayed-index state for those ranges was gone. Skip the source deletes when the merged creation failed, so the source buckets stay recoverable from their cursor properties. --- .../bucket/BucketDelayedDeliveryTracker.java | 6 +- .../BucketDelayedDeliveryTrackerTest.java | 65 +++++++++++++++++++ 2 files changed, 70 insertions(+), 1 deletion(-) 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 0655efc39208d..28487b402934a 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 @@ -612,9 +612,13 @@ private synchronized CompletableFuture asyncMergeBucketSnapshot(List FutureUtil.waitForAll( + .thenCompose(mergedBucketId -> INVALID_BUCKET_ID.equals(mergedBucketId) + ? CompletableFuture.completedFuture(null) + : FutureUtil.waitForAll( buckets.stream() .map(bucket -> bucket.asyncDeleteBucketSnapshot(stats)) .toList())); 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 05dd0f46ecbfd..0b86bb40af351 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 @@ -1264,6 +1264,71 @@ public void testClearCompletesWhenMergeSourceBucketDeleteFails() throws Exceptio } } + @Test + public void testMergeCreateFailureKeepsSourceBucketsRecoverable() throws Exception { + MockBucketSnapshotStorage storage = new MockBucketSnapshotStorage() { + final AtomicLong mergedCreateFailures = new AtomicLong(); + + @Override + public CompletableFuture createBucketSnapshot(SnapshotMetadata snapshotMetadata, + List bucketSnapshotSegments, + String bucketKey, String topicName, + String cursorName) { + if (bucketKey.endsWith("_1_10") && mergedCreateFailures.incrementAndGet() > 0) { + return FutureUtil.failedFuture( + new BucketSnapshotPersistenceException("Merged create failed")); + } + return super.createBucketSnapshot(snapshotMetadata, bucketSnapshotSegments, + bucketKey, topicName, cursorName); + } + }; + storage.start(); + + ManagedCursor cursor = new MockManagedCursor("test_merge_fail_cursor"); + AbstractPersistentDispatcherMultipleConsumers testDispatcher = + mock(AbstractPersistentDispatcherMultipleConsumers.class); + Clock testClock = mock(Clock.class); + AtomicLong testClockTime = new AtomicLong(); + when(testClock.millis()).then(x -> testClockTime.get()); + doReturn(cursor).when(testDispatcher).getCursor(); + doReturn("persistent://public/default/testMergeFail / " + cursor.getName()) + .when(testDispatcher).getName(); + + BucketDelayedDeliveryTracker tracker = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 1); + try { + for (int i = 1; i <= 11; i++) { + assertTrue(tracker.addMessage(i, i, 10L * i)); + } + + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> + assertTrue(tracker.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging))); + // The source cursor properties survive the failed merged creation. + Map properties = cursor.getCursorProperties(); + assertTrue(properties.containsKey(BucketDelayedDeliveryTracker.DELAYED_BUCKET_KEY_PREFIX + "_1_5")); + assertTrue(properties.containsKey(BucketDelayedDeliveryTracker.DELAYED_BUCKET_KEY_PREFIX + "_6_10")); + } finally { + tracker.close(); + } + + // A tracker recreated with the same cursor/storage recovers both source buckets. + BucketDelayedDeliveryTracker recovered = new BucketDelayedDeliveryTracker( + testDispatcher, timer, 100000, testClock, true, storage, + 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50); + try { + assertTrue(recovered.getImmutableBuckets().asMapOfRanges() + .containsKey(Range.closed(1L, 5L))); + assertTrue(recovered.getImmutableBuckets().asMapOfRanges() + .containsKey(Range.closed(6L, 10L))); + assertEquals(recovered.getNumberOfDelayedMessages(), 10); + } finally { + recovered.close(); + } + storage.close(); + } + @Test public void testClearWaitsForMergeSegmentLoad() throws Exception { GatedMergeLoadStorage storage = new GatedMergeLoadStorage();