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..0ab1f87c40bd6 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 @@ -84,8 +84,6 @@ public static record SnapshotKey(long ledgerId, long entryId) {} static final int AsyncOperationTimeoutSeconds = 60; - private static final Long INVALID_BUCKET_ID = -1L; - private static final int MAX_MERGE_NUM = 4; private final long minIndexCountPerBucket; @@ -354,7 +352,7 @@ private void afterCreateImmutableBucket(Pair immu immutableBucket); immutableBucket.getSnapshotCreateFuture().ifPresent(createFuture -> { - CompletableFuture future = createFuture.handle((bucketId, ex) -> { + CompletableFuture future = createFuture.whenComplete((bucketId, ex) -> { if (ex == null) { immutableBucket.setSnapshotSegments(null); immutableBucket.asyncUpdateSnapshotLength() @@ -369,8 +367,7 @@ private void afterCreateImmutableBucket(Pair immu stats.recordSuccessEvent(BucketDelayedMessageIndexStats.Type.create, System.currentTimeMillis() - startTime); - - return bucketId; + return; } log.error() @@ -397,7 +394,6 @@ private void afterCreateImmutableBucket(Pair immu snapshotSegmentLastIndexMap.remove( new SnapshotKey(lastDelayedIndex.getLedgerId(), lastDelayedIndex.getEntryId())); } - return INVALID_BUCKET_ID; }); immutableBucket.setSnapshotCreateFuture(future); }); @@ -559,16 +555,12 @@ private synchronized CompletableFuture asyncMergeBucketSnapshot(List bucket.getSnapshotCreateFuture().orElse(NULL_LONG_PROMISE)) .toList(); - return FutureUtil.waitForAll(createFutures).thenCompose(bucketId -> { - if (createFutures.stream().anyMatch(future -> INVALID_BUCKET_ID.equals(future.join()))) { - return FutureUtil.failedFuture(new RuntimeException("Can't merge buckets due to bucket create failed")); - } - + return FutureUtil.waitForAll(createFutures).thenCompose(__ -> { List>> getAllSnapshotFutures = buckets.stream().map(ImmutableBucket::getAllSnapshotSegments).toList(); return FutureUtil.waitForAll(getAllSnapshotFutures) - .thenApply(__ -> { + .thenApply(ignore -> { return CombinedSegmentDelayedIndexQueue.wrap( getAllSnapshotFutures.stream().map(CompletableFuture::join).toList()); }) @@ -675,21 +667,20 @@ public synchronized NavigableSet getScheduledMessages(int maxMessages) while (n > 0 && !sharedBucketPriorityQueue.isEmpty()) { long timestamp = sharedBucketPriorityQueue.peekN1(); - long ledgerId = sharedBucketPriorityQueue.peekN2(); - long entryId = sharedBucketPriorityQueue.peekN3(); - if (firstLiveLedgerId != null && ledgerId < firstLiveLedgerId) { - sharedBucketPriorityQueue.pop(); - removeIndexBit(ledgerId, entryId); - continue; - } if (timestamp > cutoffTime) { break; } - SnapshotKey snapshotKey = new SnapshotKey(ledgerId, entryId); + long ledgerId = sharedBucketPriorityQueue.peekN2(); + long entryId = sharedBucketPriorityQueue.peekN3(); + boolean orphaned = firstLiveLedgerId != null && ledgerId < firstLiveLedgerId; + SnapshotKey snapshotKey = new SnapshotKey(ledgerId, entryId); ImmutableBucket bucket = snapshotSegmentLastIndexMap.get(snapshotKey); - if (bucket != null && immutableBuckets.asMapOfRanges().containsValue(bucket)) { + boolean shouldLoadNextSegment = bucket != null + && immutableBuckets.asMapOfRanges().containsValue(bucket) + && (!orphaned || bucket.getEndLedgerId() >= firstLiveLedgerId); + if (shouldLoadNextSegment) { // All message of current snapshot segment are scheduled, try load next snapshot segment if (bucket.merging) { log.info() @@ -765,6 +756,12 @@ public synchronized NavigableSet getScheduledMessages(int maxMessages) } sharedBucketPriorityQueue.pop(); + + if (orphaned) { + removeIndexBit(ledgerId, entryId); + continue; + } + // Dedup: queue may carry the same position twice (initial seal + merge); only the // first delivery of each position decrements the counter via removeIndexBit. if (removeIndexBit(ledgerId, entryId)) { @@ -906,13 +903,18 @@ private CompletableFuture deleteBucketSnapshot(String ledgerName, // creation before deleting. When the creation failed the bucket has already been removed and // downgraded to memory mode, so there is no snapshot left to delete. return bucket.getSnapshotCreateFuture().orElse(NULL_LONG_PROMISE) - .thenCompose(bucketId -> INVALID_BUCKET_ID.equals(bucketId) + .thenCompose(bucketId -> bucketId == null ? CompletableFuture.completedFuture(null) : doDeleteBucketSnapshot(ledgerName, range, bucket)); } private CompletableFuture doDeleteBucketSnapshot(String ledgerName, Range range, ImmutableBucket bucket) { + synchronized (this) { + if (!isCurrentBucket(range, bucket)) { + return CompletableFuture.completedFuture(null); + } + } return bucket.asyncDeleteBucketSnapshot(stats) .handle((__, t) -> { if (t != null) { @@ -922,6 +924,9 @@ private CompletableFuture doDeleteBucketSnapshot(String ledgerName, throw new CompletionException(t); } synchronized (this) { + if (!isCurrentBucket(range, bucket)) { + return null; + } snapshotSegmentLastIndexMap.entrySet().removeIf(entry -> entry.getValue() == bucket); removeBucket(range); bucket.getDelayedIndexBitMap().forEach((ledgerId, bitmap) -> @@ -931,6 +936,11 @@ private CompletableFuture doDeleteBucketSnapshot(String ledgerName, }); } + private boolean isCurrentBucket(Range range, ImmutableBucket expectedBucket) { + ImmutableBucket currentBucket = immutableBuckets.asMapOfRanges().get(range); + return currentBucket == expectedBucket; + } + private Long firstActiveLedgerId() { ManagedCursor cursor = context.getCursor(); Position mdp = cursor.getMarkDeletedPosition(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/ImmutableBucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/ImmutableBucket.java index 925c4756019c7..b229e174c8602 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/ImmutableBucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/ImmutableBucket.java @@ -319,32 +319,55 @@ CompletableFuture asyncDeleteBucketSnapshot(BucketDelayedMessageIndexStats long deleteStartTime = System.currentTimeMillis(); stats.recordTriggerEvent(BucketDelayedMessageIndexStats.Type.delete); String bucketKey = bucketKey(); - long bucketId = getAndUpdateBucketId(); + CompletableFuture> bucketIdFuture = resolveBucketIdForDelete(); - return executeWithRetry(() -> ctx.bucketSnapshotStorage().deleteBucketSnapshot(bucketId), - BucketSnapshotPersistenceException.class, MaxRetryTimes) - .whenComplete((__, ex) -> { - if (ex != null) { - log.error() - .attr("dispatcher", ctx.dispatcherName()) - .attr("bucketId", bucketId) - .attr("bucketKey", bucketKey) - .exception(ex) - .log("Failed to delete bucket snapshot"); + return bucketIdFuture.thenCompose(optionalBucketId -> { + if (optionalBucketId.isEmpty()) { + return CompletableFuture.completedFuture(null); + } + long bucketId = optionalBucketId.get(); + return executeWithRetry(() -> ctx.bucketSnapshotStorage().deleteBucketSnapshot(bucketId), + BucketSnapshotPersistenceException.class, MaxRetryTimes) + .whenComplete((__, ex) -> { + if (ex != null) { + log.error() + .attr("dispatcher", ctx.dispatcherName()) + .attr("bucketId", bucketId) + .attr("bucketKey", bucketKey) + .exception(ex) + .log("Failed to delete bucket snapshot"); - stats.recordFailEvent(BucketDelayedMessageIndexStats.Type.delete); - } else { - log.info() - .attr("dispatcher", ctx.dispatcherName()) - .attr("bucketId", bucketId) - .attr("bucketKey", bucketKey) - .log("Delete bucket snapshot finish"); + stats.recordFailEvent(BucketDelayedMessageIndexStats.Type.delete); + } else { + log.info() + .attr("dispatcher", ctx.dispatcherName()) + .attr("bucketId", bucketId) + .attr("bucketKey", bucketKey) + .log("Delete bucket snapshot finish"); + + stats.recordSuccessEvent(BucketDelayedMessageIndexStats.Type.delete, + System.currentTimeMillis() - deleteStartTime); + } + }); + }).thenCompose(__ -> removeBucketCursorProperty(bucketKey)); + } - stats.recordSuccessEvent(BucketDelayedMessageIndexStats.Type.delete, - System.currentTimeMillis() - deleteStartTime); - } - }) - .thenCompose(__ -> removeBucketCursorProperty(bucketKey)); + private CompletableFuture> resolveBucketIdForDelete() { + Optional> snapshotCreateFuture = getSnapshotCreateFuture(); + if (snapshotCreateFuture.isPresent()) { + return snapshotCreateFuture.get().handle((bucketId, ex) -> { + if (ex != null) { + return Optional.empty(); + } + setBucketId(bucketId); + return Optional.of(bucketId); + }); + } + try { + return CompletableFuture.completedFuture(Optional.of(getAndUpdateBucketId())); + } catch (Exception e) { + return FutureUtil.failedFuture(e); + } } CompletableFuture clear(BucketDelayedMessageIndexStats stats) { 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 9ba3a4be50e28..8485310bd69fb 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 @@ -1002,67 +1002,89 @@ private CompletableFuture resetCursorInternal(Position finalPosition, Comp log.info() .log("Successfully disconnected consumers from subscription, proceeding with cursor reset"); - CompletableFuture forceReset = new CompletableFuture<>(); - if (topic.getTopicCompactionService() == null) { - forceReset.complete(false); - } else { - topic.getTopicCompactionService().getLastCompactedPosition().thenAccept(lastCompactedPosition -> { - Position resetTo = finalPosition; - if (lastCompactedPosition != null && resetTo.compareTo(lastCompactedPosition.getLedgerId(), - lastCompactedPosition.getEntryId()) <= 0) { - forceReset.complete(true); - } else { - forceReset.complete(false); - } - }).exceptionally(ex -> { - forceReset.completeExceptionally(ex); - return null; - }); + CompletableFuture clearDelayedMessagesFuture; + try { + clearDelayedMessagesFuture = dispatcher != null + ? dispatcher.clearDelayedMessages() + : CompletableFuture.completedFuture(null); + } catch (Throwable t) { + clearDelayedMessagesFuture = FutureUtil.failedFuture(t); } - forceReset.thenAccept(forceResetValue -> { - cursor.asyncResetCursor(finalPosition, forceResetValue, new AsyncCallbacks.ResetCursorCallback() { - @Override - public void resetComplete(Object ctx) { - log.debug() - .attr("finalPosition", finalPosition) - .log("Successfully reset subscription to position"); - if (dispatcher != null) { - dispatcher.cursorIsReset(); - dispatcher.afterAckMessages(null, finalPosition); - } - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - inProgressResetCursorFuture = null; - future.complete(null); - } + clearDelayedMessagesFuture.whenComplete((__, clearEx) -> { + if (clearEx != null) { + log.error() + .exception(clearEx) + .log("Error while clearing delayed messages during cursor reset"); + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + inProgressResetCursorFuture = null; + future.completeExceptionally(new BrokerServiceException(clearEx)); + return; + } - @Override - public void resetFailed(ManagedLedgerException exception, Object ctx) { - log.error() - .attr("finalPosition", finalPosition) - .exception(exception) - .log("Failed to reset subscription to position"); - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - inProgressResetCursorFuture = null; - // todo - retry on InvalidCursorPositionException - // or should we just ask user to retry one more time? - if (exception instanceof InvalidCursorPositionException) { - future.completeExceptionally(new SubscriptionInvalidCursorPosition(exception.getMessage())); - } else if (exception instanceof ConcurrentFindCursorPositionException) { - future.completeExceptionally(new SubscriptionBusyException(exception.getMessage())); + CompletableFuture forceReset = new CompletableFuture<>(); + if (topic.getTopicCompactionService() == null) { + forceReset.complete(false); + } else { + topic.getTopicCompactionService().getLastCompactedPosition().thenAccept(lastCompactedPosition -> { + Position resetTo = finalPosition; + if (lastCompactedPosition != null && resetTo.compareTo(lastCompactedPosition.getLedgerId(), + lastCompactedPosition.getEntryId()) <= 0) { + forceReset.complete(true); } else { - future.completeExceptionally(new BrokerServiceException(exception)); + forceReset.complete(false); } - } + }).exceptionally(ex -> { + forceReset.completeExceptionally(ex); + return null; + }); + } + + forceReset.thenAccept(forceResetValue -> { + cursor.asyncResetCursor(finalPosition, forceResetValue, new AsyncCallbacks.ResetCursorCallback() { + @Override + public void resetComplete(Object ctx) { + log.debug() + .attr("finalPosition", finalPosition) + .log("Successfully reset subscription to position"); + if (dispatcher != null) { + dispatcher.cursorIsReset(); + dispatcher.afterAckMessages(null, finalPosition); + } + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + inProgressResetCursorFuture = null; + future.complete(null); + } + + @Override + public void resetFailed(ManagedLedgerException exception, Object ctx) { + log.error() + .attr("finalPosition", finalPosition) + .exception(exception) + .log("Failed to reset subscription to position"); + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + inProgressResetCursorFuture = null; + // todo - retry on InvalidCursorPositionException + // or should we just ask user to retry one more time? + if (exception instanceof InvalidCursorPositionException) { + future.completeExceptionally( + new SubscriptionInvalidCursorPosition(exception.getMessage())); + } else if (exception instanceof ConcurrentFindCursorPositionException) { + future.completeExceptionally(new SubscriptionBusyException(exception.getMessage())); + } else { + future.completeExceptionally(new BrokerServiceException(exception)); + } + } + }); + }).exceptionally((e) -> { + log.error() + .exception(e) + .log("Error while resetting cursor"); + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + inProgressResetCursorFuture = null; + future.completeExceptionally(new BrokerServiceException(e)); + return null; }); - }).exceptionally((e) -> { - log.error() - .exception(e) - .log("Error while resetting cursor"); - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - inProgressResetCursorFuture = null; - future.completeExceptionally(new BrokerServiceException(e)); - return null; }); }); return future; 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..61ebba70b441a 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 @@ -20,9 +20,13 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.atLeastOnce; +import static org.mockito.Mockito.clearInvocations; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotSame; @@ -167,7 +171,9 @@ public Object[][] provider(Method method) throws Exception { true, bucketSnapshotStorage, 1000, TimeUnit.MILLISECONDS.toMillis(100), -1, 50) }}; case "testExpiredTrackedMessageReturnsFalse", "testRecoverThenExpireAddMessage", - "testExpiredTrackedMessageDecrementsCount" -> new Object[][]{{ + "testExpiredTrackedMessageDecrementsCount", + "testLoadsNextSnapshotSegmentAfterCutoff", + "testDoesNotLoadNextSnapshotSegmentBeforeCutoff" -> new Object[][]{{ new BucketDelayedDeliveryTracker(dispatcher, timer, 1, clock, true, bucketSnapshotStorage, 5, TimeUnit.MILLISECONDS.toMillis(10), -1, 50) }}; @@ -247,6 +253,69 @@ public void testRecoverThenExpireAddMessage(BucketDelayedDeliveryTracker tracker tracker2.close(); } + @Test(dataProvider = "delayedTracker") + public void testDoesNotLoadNextSnapshotSegmentBeforeCutoff(BucketDelayedDeliveryTracker tracker) + throws Exception { + for (int i = 1; i <= 6; i++) { + tracker.addMessage(i, i, i * 10); + } + + Awaitility.await().untilAsserted(() -> + assertTrue(tracker.getImmutableBuckets().asMapOfRanges().values().stream() + .noneMatch(x -> x.merging || !x.getSnapshotCreateFuture().get().isDone()))); + + tracker.close(); + + MockManagedCursor recoveredCursor = spy((MockManagedCursor) dispatcher.getCursor()); + doReturn(PositionFactory.create(3, 0)).when(recoveredCursor).getMarkDeletedPosition(); + doReturn(recoveredCursor).when(dispatcher).getCursor(); + + clockTime.set(1); + MockBucketSnapshotStorage localStorage = spy((MockBucketSnapshotStorage) bucketSnapshotStorage); + @Cleanup + BucketDelayedDeliveryTracker tracker2 = new BucketDelayedDeliveryTracker( + dispatcher, timer, 1000, clock, true, localStorage, 4, TimeUnit.MILLISECONDS.toMillis(10), 2, 50); + + // Recovery loads segment 1; reset so the assertion observes only getScheduledMessages-driven loads. + clearInvocations(localStorage); + assertTrue(tracker2.getScheduledMessages(10).isEmpty()); + + Awaitility.await().during(3, TimeUnit.SECONDS).untilAsserted(() -> + verify(localStorage, never()).getBucketSnapshotSegment(anyLong(), anyLong(), anyLong())); + } + + @Test(dataProvider = "delayedTracker") + public void testLoadsNextSnapshotSegmentAfterCutoff(BucketDelayedDeliveryTracker tracker) + throws Exception { + 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()))); + + tracker.close(); + + MockManagedCursor recoveredCursor = spy((MockManagedCursor) dispatcher.getCursor()); + doReturn(PositionFactory.create(3, 0)).when(recoveredCursor).getMarkDeletedPosition(); + doReturn(recoveredCursor).when(dispatcher).getCursor(); + + clockTime.set(0); + MockBucketSnapshotStorage localStorage = spy((MockBucketSnapshotStorage) bucketSnapshotStorage); + @Cleanup + BucketDelayedDeliveryTracker tracker2 = new BucketDelayedDeliveryTracker( + dispatcher, timer, 1000, clock, true, localStorage, 4, TimeUnit.MILLISECONDS.toMillis(10), 2, 50); + + // Trigger the segment loading + clockTime.set(600); + tracker2.getScheduledMessages(10); + Awaitility.await().untilAsserted(() -> { + assertTrue(tracker2.getScheduledMessages(10).contains(PositionFactory.create(4, 4))); + verify(localStorage, atLeastOnce()).getBucketSnapshotSegment(anyLong(), anyLong(), anyLong()); + }); + } + @Test(dataProvider = "delayedTracker") public void testExpiredTrackedMessageDecrementsCount(BucketDelayedDeliveryTracker tracker) { clockTime.set(1000);