From 9861ea331b4a14e3e402058a10b08f91500e6aa5 Mon Sep 17 00:00:00 2001 From: coderzc Date: Thu, 29 Dec 2022 14:59:58 +0800 Subject: [PATCH 1/7] Implement bucket snapshot merge/delete --- .../pulsar/broker/delayed/bucket/Bucket.java | 5 + .../bucket/BucketDelayedDeliveryTracker.java | 79 +++++++++- .../CombinedSegmentDelayedIndexQueue.java | 108 +++++++++++++ .../delayed/bucket/DelayedIndexQueue.java | 42 +++++ .../delayed/bucket/ImmutableBucket.java | 17 +- .../broker/delayed/bucket/MutableBucket.java | 31 ++-- .../TripleLongPriorityDelayedIndexQueue.java | 57 +++++++ .../BucketDelayedDeliveryTrackerTest.java | 17 ++ .../broker/delayed/DelayedIndexQueueTest.java | 145 ++++++++++++++++++ 9 files changed, 477 insertions(+), 24 deletions(-) create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java index c094d1dee7b54..5d2a556337a6e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java @@ -143,4 +143,9 @@ private CompletableFuture putBucketKeyId(String bucketKey, Long bucketId) return executeWithRetry(() -> cursor.putCursorProperty(bucketKey, String.valueOf(bucketId)), ManagedLedgerException.BadVersionException.class, MaxRetryTimes); } + + protected CompletableFuture removeBucketCursorProperty(String bucketKey) { + return executeWithRetry(() -> cursor.removeCursorProperty(bucketKey), + ManagedLedgerException.BadVersionException.class, MaxRetryTimes); + } } 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 b402b51ce071c..0f28bff46b6f6 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 @@ -21,6 +21,7 @@ import static com.google.common.base.Preconditions.checkArgument; import static org.apache.pulsar.broker.delayed.bucket.Bucket.DELAYED_BUCKET_KEY_PREFIX; import static org.apache.pulsar.broker.delayed.bucket.Bucket.DELIMITER; +import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.HashBasedTable; import com.google.common.collect.Range; import com.google.common.collect.RangeMap; @@ -40,6 +41,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import javax.annotation.concurrent.ThreadSafe; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; @@ -71,6 +73,8 @@ public class BucketDelayedDeliveryTracker extends AbstractDelayedDeliveryTracker private final TripleLongPriorityQueue sharedBucketPriorityQueue; + @Getter + @VisibleForTesting private final RangeMap immutableBuckets; private final Table snapshotSegmentLastIndexTable; @@ -126,6 +130,12 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { CompletableFuture future = immutableBucket.asyncRecoverBucketSnapshotEntry(this::getCutoffTime).thenAccept(indexList -> { if (CollectionUtils.isEmpty(indexList)) { + // Delete bucket snapshot if indexList is empty + synchronized (immutableBuckets) { + immutableBuckets.remove( + Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId)); + } + immutableBucket.asyncDeleteBucketSnapshot(); return; } DelayedIndex lastDelayedIndex = indexList.get(indexList.size() - 1); @@ -179,10 +189,7 @@ private Optional findImmutableBucket(long ledgerId) { return Optional.ofNullable(immutableBuckets.get(ledgerId)); } - private void sealBucket() { - Pair immutableBucketDelayedIndexPair = - lastMutableBucket.sealBucketAndAsyncPersistent(this.timeStepPerBucketSnapshotSegment, - this.sharedBucketPriorityQueue); + private void afterCreateImmutableBucket(Pair immutableBucketDelayedIndexPair) { if (immutableBucketDelayedIndexPair != null) { ImmutableBucket immutableBucket = immutableBucketDelayedIndexPair.getLeft(); immutableBuckets.put(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId), @@ -214,11 +221,18 @@ public synchronized boolean addMessage(long ledgerId, long entryId, long deliver if (!existBucket && ledgerId > lastMutableBucket.endLedgerId && lastMutableBucket.size() >= minIndexCountPerBucket && !lastMutableBucket.isEmpty()) { - sealBucket(); + Pair immutableBucketDelayedIndexPair = + lastMutableBucket.sealBucketAndAsyncPersistent(this.timeStepPerBucketSnapshotSegment, + this.sharedBucketPriorityQueue); + afterCreateImmutableBucket(immutableBucketDelayedIndexPair); lastMutableBucket.resetLastMutableBucketRange(); if (immutableBuckets.asMapOfRanges().size() > maxNumBuckets) { - // TODO merge bucket snapshot (synchronize operate) + try { + asyncMergeBucketSnapshot().get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); + } catch (InterruptedException | ExecutionException | TimeoutException e) { + throw new RuntimeException(e); + } } } @@ -243,6 +257,53 @@ public synchronized boolean addMessage(long ledgerId, long entryId, long deliver return true; } + private synchronized CompletableFuture asyncMergeBucketSnapshot() { + List values = immutableBuckets.asMapOfRanges().values().stream().toList(); + long minNumberMessages = Long.MAX_VALUE; + int minIndex = -1; + for (int i = 0; i + 1 < values.size(); i++) { + ImmutableBucket bucketL = values.get(i); + ImmutableBucket bucketR = values.get(i + 1); + long numberMessages = bucketL.numberBucketDelayedMessages + bucketR.numberBucketDelayedMessages; + if (numberMessages < minNumberMessages) { + minNumberMessages = (int) numberMessages; + minIndex = i; + } + } + return asyncMergeBucketSnapshot(values.get(minIndex), values.get(minIndex + 1)); + } + + private synchronized CompletableFuture asyncMergeBucketSnapshot(ImmutableBucket bucketA, + ImmutableBucket bucketB) { + immutableBuckets.remove(Range.closed(bucketA.startLedgerId, bucketA.endLedgerId)); + immutableBuckets.remove(Range.closed(bucketB.startLedgerId, bucketB.endLedgerId)); + + CompletableFuture snapshotCreateFutureA = + bucketA.getSnapshotCreateFuture().orElse(CompletableFuture.completedFuture(null)); + CompletableFuture snapshotCreateFutureB = + bucketB.getSnapshotCreateFuture().orElse(CompletableFuture.completedFuture(null)); + + return CompletableFuture.allOf(snapshotCreateFutureA, snapshotCreateFutureB).thenCompose(__ -> { + CompletableFuture> futureA = + bucketA.getRemainSnapshotSegment(); + CompletableFuture> futureB = + bucketB.getRemainSnapshotSegment(); + return futureA.thenCombine(futureB, CombinedSegmentDelayedIndexQueue::wrap) + .thenCompose(combinedDelayedIndexQueue -> { + CompletableFuture removeAFuture = bucketA.asyncDeleteBucketSnapshot(); + CompletableFuture removeBFuture = bucketB.asyncDeleteBucketSnapshot(); + + return CompletableFuture.allOf(removeAFuture, removeBFuture).thenRun(() -> { + Pair immutableBucketDelayedIndexPair = + lastMutableBucket.createImmutableBucketAndAsyncPersistent( + timeStepPerBucketSnapshotSegment, sharedBucketPriorityQueue, + combinedDelayedIndexQueue, bucketA.startLedgerId, bucketB.endLedgerId); + afterCreateImmutableBucket(immutableBucketDelayedIndexPair); + }); + }); + }); + } + @Override public synchronized boolean hasMessageAvailable() { long cutoffTime = getCutoffTime(); @@ -299,7 +360,7 @@ public synchronized NavigableSet getScheduledMessages(int maxMessa removeIndexBit(ledgerId, entryId); ImmutableBucket bucket = snapshotSegmentLastIndexTable.remove(ledgerId, entryId); - if (bucket != null) { + if (bucket != null && immutableBuckets.asMapOfRanges().containsValue(bucket)) { if (log.isDebugEnabled()) { log.debug("[{}] Load next snapshot segment, bucket: {}", dispatcher.getName(), bucket); } @@ -308,6 +369,10 @@ public synchronized NavigableSet getScheduledMessages(int maxMessa try { bucket.asyncLoadNextBucketSnapshotEntry().thenAccept(indexList -> { if (CollectionUtils.isEmpty(indexList)) { + synchronized (immutableBuckets) { + immutableBuckets.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId)); + } + bucket.asyncDeleteBucketSnapshot(); return; } DelayedMessageIndexBucketSnapshotFormat.DelayedIndex diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java new file mode 100644 index 0000000000000..fd7b7884cd69a --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java @@ -0,0 +1,108 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.delayed.bucket; + +import java.util.List; +import javax.annotation.concurrent.NotThreadSafe; +import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.DelayedIndex; +import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.SnapshotSegment; + +@NotThreadSafe +public class CombinedSegmentDelayedIndexQueue implements DelayedIndexQueue { + + private final List segmentListA; + private final List segmentListB; + + private int segmentListACursor = 0; + private int segmentListBCursor = 0; + private int segmentACursor = 0; + private int segmentBCursor = 0; + + private CombinedSegmentDelayedIndexQueue(List segmentListA, + List segmentListB) { + this.segmentListA = segmentListA; + this.segmentListB = segmentListB; + } + + public static CombinedSegmentDelayedIndexQueue wrap( + List segmentListA, + List segmentListB) { + return new CombinedSegmentDelayedIndexQueue(segmentListA, segmentListB); + } + + @Override + public boolean isEmpty() { + return segmentListACursor >= segmentListA.size() && segmentListBCursor >= segmentListB.size(); + } + + @Override + public DelayedIndex peek() { + return getValue(false); + } + + @Override + public DelayedIndex pop() { + return getValue(true); + } + + private DelayedIndex getValue(boolean needAdvanceCursor) { + while (segmentListACursor < segmentListA.size() + && segmentACursor >= segmentListA.get(segmentListACursor).getIndexesCount()) { + segmentListACursor++; + } + while (segmentListBCursor < segmentListB.size() + && segmentBCursor >= segmentListB.get(segmentListBCursor).getIndexesCount()) { + segmentListBCursor++; + } + + DelayedIndex delayedIndexA = null; + DelayedIndex delayedIndexB = null; + if (segmentListACursor >= segmentListA.size()) { + delayedIndexB = segmentListB.get(segmentListBCursor).getIndexes(segmentBCursor); + } else if (segmentListBCursor >= segmentListB.size()) { + delayedIndexA = segmentListA.get(segmentListACursor).getIndexes(segmentACursor); + } else { + delayedIndexA = segmentListA.get(segmentListACursor).getIndexes(segmentACursor); + delayedIndexB = segmentListB.get(segmentListBCursor).getIndexes(segmentBCursor); + } + + DelayedIndex resultValue; + if (delayedIndexB == null || (delayedIndexA != null && COMPARATOR.compare(delayedIndexA, delayedIndexB) < 0)) { + resultValue = delayedIndexA; + if (needAdvanceCursor) { + if (++segmentACursor >= segmentListA.get(segmentListACursor).getIndexesCount()) { + segmentListA.set(segmentListACursor, null); + ++segmentListACursor; + segmentACursor = 0; + } + } + } else { + resultValue = delayedIndexB; + if (needAdvanceCursor) { + if (++segmentBCursor >= segmentListB.get(segmentListBCursor).getIndexesCount()) { + segmentListB.set(segmentListBCursor, null); + ++segmentListBCursor; + segmentBCursor = 0; + } + } + } + + return resultValue; + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java new file mode 100644 index 0000000000000..df80c21767b12 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java @@ -0,0 +1,42 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.delayed.bucket; + +import java.util.Comparator; +import java.util.Objects; +import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat; + +public interface DelayedIndexQueue { + + Comparator COMPARATOR = (o1, o2) -> { + if (!Objects.equals(o1.getTimestamp(), o2.getTimestamp())) { + return Long.compare(o1.getTimestamp(), o2.getTimestamp()); + } else if (!Objects.equals(o1.getLedgerId(), o2.getLedgerId())) { + return Long.compare(o1.getLedgerId(), o2.getLedgerId()); + } else { + return Long.compare(o1.getEntryId(), o2.getEntryId()); + } + }; + + boolean isEmpty(); + + DelayedMessageIndexBucketSnapshotFormat.DelayedIndex peek(); + + DelayedMessageIndexBucketSnapshotFormat.DelayedIndex pop(); +} 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 913893b175350..dd23620aceb4f 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 @@ -90,7 +90,6 @@ private CompletableFuture> asyncLoadNextBucketSnapshotEntry(b return loadMetaDataFuture.thenCompose(nextSegmentEntryId -> { if (nextSegmentEntryId > lastSegmentEntryId) { - // TODO Delete bucket snapshot return CompletableFuture.completedFuture(null); } @@ -132,12 +131,26 @@ private void recoverDelayedIndexBitMapAndNumber(int startSnapshotIndex, this.setNumberBucketDelayedMessages(numberMessages.getValue()); } + CompletableFuture> getRemainSnapshotSegment() { + return bucketSnapshotStorage.getBucketSnapshotSegment(getAndUpdateBucketId(), currentSegmentEntryId, + lastSegmentEntryId); + } + + CompletableFuture asyncDeleteBucketSnapshot() { + return removeBucketCursorProperty(bucketKey()).thenCompose(__ -> + bucketSnapshotStorage.deleteBucketSnapshot(getAndUpdateBucketId())); + } + void clear(boolean delete) { delayedIndexBitMap.clear(); getSnapshotCreateFuture().ifPresent(snapshotGenerateFuture -> { if (delete) { snapshotGenerateFuture.cancel(true); - // TODO delete bucket snapshot + try { + asyncDeleteBucketSnapshot().get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); + } catch (Exception e) { + throw new RuntimeException(e); + } } else { try { snapshotGenerateFuture.get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java index 36026298269d7..e70ccf2e9831d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java @@ -51,14 +51,19 @@ class MutableBucket extends Bucket implements AutoCloseable { Pair sealBucketAndAsyncPersistent( long timeStepPerBucketSnapshotSegment, TripleLongPriorityQueue sharedQueue) { - if (priorityQueue.isEmpty()) { + return createImmutableBucketAndAsyncPersistent(timeStepPerBucketSnapshotSegment, sharedQueue, + TripleLongPriorityDelayedIndexQueue.wrap(priorityQueue), startLedgerId, endLedgerId); + } + + Pair createImmutableBucketAndAsyncPersistent( + final long timeStepPerBucketSnapshotSegment, + TripleLongPriorityQueue sharedQueue, DelayedIndexQueue delayedIndexQueue, final long startLedgerId, + final long endLedgerId) { + if (delayedIndexQueue.isEmpty()) { return null; } long numMessages = 0; - final long startLedgerId = getStartLedgerId(); - final long endLedgerId = getEndLedgerId(); - List bucketSnapshotSegments = new ArrayList<>(); List segmentMetadataList = new ArrayList<>(); Map bitMap = new HashMap<>(); @@ -66,14 +71,15 @@ Pair sealBucketAndAsyncPersistent( SnapshotSegmentMetadata.Builder segmentMetadataBuilder = SnapshotSegmentMetadata.newBuilder(); long currentTimestampUpperLimit = 0; - while (!priorityQueue.isEmpty()) { - long timestamp = priorityQueue.peekN1(); + while (!delayedIndexQueue.isEmpty()) { + DelayedIndex delayedIndex = delayedIndexQueue.peek(); + long timestamp = delayedIndex.getTimestamp(); if (currentTimestampUpperLimit == 0) { currentTimestampUpperLimit = timestamp + timeStepPerBucketSnapshotSegment - 1; } - long ledgerId = priorityQueue.peekN2(); - long entryId = priorityQueue.peekN3(); + long ledgerId = delayedIndex.getLedgerId(); + long entryId = delayedIndex.getEntryId(); checkArgument(ledgerId >= startLedgerId && ledgerId <= endLedgerId); @@ -82,19 +88,14 @@ Pair sealBucketAndAsyncPersistent( sharedQueue.add(timestamp, ledgerId, entryId); } - priorityQueue.pop(); + delayedIndexQueue.pop(); numMessages++; - DelayedIndex delayedIndex = DelayedIndex.newBuilder() - .setTimestamp(timestamp) - .setLedgerId(ledgerId) - .setEntryId(entryId).build(); - bitMap.computeIfAbsent(ledgerId, k -> new RoaringBitmap()).add(entryId, entryId + 1); snapshotSegmentBuilder.addIndexes(delayedIndex); - if (priorityQueue.isEmpty() || priorityQueue.peekN1() > currentTimestampUpperLimit) { + if (delayedIndexQueue.isEmpty() || delayedIndexQueue.peek().getTimestamp() > currentTimestampUpperLimit) { segmentMetadataBuilder.setMaxScheduleTimestamp(timestamp); currentTimestampUpperLimit = 0; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java new file mode 100644 index 0000000000000..97ff958188295 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java @@ -0,0 +1,57 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.delayed.bucket; + +import javax.annotation.concurrent.NotThreadSafe; +import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat; +import org.apache.pulsar.common.util.collections.TripleLongPriorityQueue; + +@NotThreadSafe +public class TripleLongPriorityDelayedIndexQueue implements DelayedIndexQueue { + + private final TripleLongPriorityQueue queue; + + private TripleLongPriorityDelayedIndexQueue(TripleLongPriorityQueue queue) { + this.queue = queue; + } + + public static TripleLongPriorityDelayedIndexQueue wrap(TripleLongPriorityQueue queue) { + return new TripleLongPriorityDelayedIndexQueue(queue); + } + + @Override + public boolean isEmpty() { + return queue.isEmpty(); + } + + @Override + public DelayedMessageIndexBucketSnapshotFormat.DelayedIndex peek() { + DelayedMessageIndexBucketSnapshotFormat.DelayedIndex delayedIndex = + DelayedMessageIndexBucketSnapshotFormat.DelayedIndex.newBuilder().setTimestamp(queue.peekN1()) + .setLedgerId(queue.peekN2()).setEntryId(queue.peekN3()).build(); + return delayedIndex; + } + + @Override + public DelayedMessageIndexBucketSnapshotFormat.DelayedIndex pop() { + DelayedMessageIndexBucketSnapshotFormat.DelayedIndex peek = peek(); + queue.pop(); + return peek; + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java index abcde7902d84e..5be37ec3ef5e9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java @@ -133,6 +133,10 @@ public Object[][] provider(Method method) throws Exception { new BucketDelayedDeliveryTracker(dispatcher, timer, 500, clock, true, bucketSnapshotStorage, 5, TimeUnit.MILLISECONDS.toMillis(10), 50) }}; + case "testMergeSnapshot" -> new Object[][]{{ + new BucketDelayedDeliveryTracker(dispatcher, timer, 100000, clock, + true, bucketSnapshotStorage, 5, TimeUnit.MILLISECONDS.toMillis(10), 10) + }}; default -> new Object[][]{{ new BucketDelayedDeliveryTracker(dispatcher, timer, 1, clock, true, bucketSnapshotStorage, 1000, TimeUnit.MILLISECONDS.toMillis(100), 50) @@ -235,4 +239,17 @@ public void testRoaringBitmapSerialize() { assertTrue(Arrays.equals(array, array2)); assertNotSame(array, array2); } + + @Test(dataProvider = "delayedTracker") + public void testMergeSnapshot(BucketDelayedDeliveryTracker tracker) { + for (int i = 1; i <= 110; i++) { + tracker.addMessage(i, i, i * 10); + } + + assertEquals(110, tracker.getNumberOfDelayedMessages()); + + int size = tracker.getImmutableBuckets().asMapOfRanges().size(); + + assertEquals(10, size); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java new file mode 100644 index 0000000000000..357d14fad7438 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java @@ -0,0 +1,145 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.delayed; + +import static org.apache.pulsar.broker.delayed.bucket.DelayedIndexQueue.COMPARATOR; +import java.util.ArrayList; +import java.util.List; +import lombok.Cleanup; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.delayed.bucket.CombinedSegmentDelayedIndexQueue; +import org.apache.pulsar.broker.delayed.bucket.TripleLongPriorityDelayedIndexQueue; +import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.DelayedIndex; +import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.SnapshotSegment; +import org.apache.pulsar.common.util.collections.TripleLongPriorityQueue; +import org.testng.Assert; +import org.testng.annotations.Test; + +@Slf4j +public class DelayedIndexQueueTest { + + @Test + public void testCompare() { + DelayedIndex delayedIndex = + DelayedIndex.newBuilder().setTimestamp(1).setLedgerId(1L).setEntryId(1L) + .build(); + DelayedIndex delayedIndex2 = + DelayedIndex.newBuilder().setTimestamp(2).setLedgerId(2L).setEntryId(2L) + .build(); + Assert.assertTrue(COMPARATOR.compare(delayedIndex, delayedIndex2) < 0); + + delayedIndex = + DelayedIndex.newBuilder().setTimestamp(1).setLedgerId(1L).setEntryId(1L) + .build(); + delayedIndex2 = + DelayedIndex.newBuilder().setTimestamp(1).setLedgerId(2L).setEntryId(2L) + .build(); + Assert.assertTrue(COMPARATOR.compare(delayedIndex, delayedIndex2) < 0); + + delayedIndex = + DelayedIndex.newBuilder().setTimestamp(1).setLedgerId(1L).setEntryId(1L) + .build(); + delayedIndex2 = + DelayedIndex.newBuilder().setTimestamp(1).setLedgerId(1L).setEntryId(2L) + .build(); + Assert.assertTrue(COMPARATOR.compare(delayedIndex, delayedIndex2) < 0); + } + + @Test + public void testCombinedSegmentDelayedIndexQueue() { + List listA = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + DelayedIndex delayedIndex = + DelayedIndex.newBuilder().setTimestamp(i).setLedgerId(1L).setEntryId(1L) + .build(); + listA.add(delayedIndex); + } + SnapshotSegment snapshotSegmentA1 = SnapshotSegment.newBuilder().addAllIndexes(listA).build(); + + List listA2 = new ArrayList<>(); + for (int i = 10; i < 20; i++) { + DelayedIndex delayedIndex = + DelayedIndex.newBuilder().setTimestamp(i).setLedgerId(1L).setEntryId(1L) + .build(); + listA2.add(delayedIndex); + } + SnapshotSegment snapshotSegmentA2 = SnapshotSegment.newBuilder().addAllIndexes(listA2).build(); + + List segmentListA = new ArrayList<>(); + segmentListA.add(snapshotSegmentA1); + segmentListA.add(snapshotSegmentA2); + + List listB = new ArrayList<>(); + for (int i = 0; i < 9; i++) { + DelayedIndex delayedIndex = + DelayedIndex.newBuilder().setTimestamp(i).setLedgerId(2L).setEntryId(1L) + .build(); + + DelayedIndex delayedIndex2 = + DelayedIndex.newBuilder().setTimestamp(i).setLedgerId(2L).setEntryId(2L) + .build(); + listB.add(delayedIndex); + listB.add(delayedIndex2); + } + + SnapshotSegment snapshotSegmentB = SnapshotSegment.newBuilder().addAllIndexes(listB).build(); + List segmentListB = new ArrayList<>(); + segmentListB.add(snapshotSegmentB); + segmentListB.add(SnapshotSegment.newBuilder().build()); + + CombinedSegmentDelayedIndexQueue delayedIndexQueue = + CombinedSegmentDelayedIndexQueue.wrap(segmentListA, segmentListB); + + int count = 0; + while (!delayedIndexQueue.isEmpty()) { + DelayedIndex pop = delayedIndexQueue.pop(); + log.info("{} , {}, {}", pop.getTimestamp(), pop.getLedgerId(), pop.getEntryId()); + count++; + if (!delayedIndexQueue.isEmpty()) { + DelayedIndex peek = delayedIndexQueue.peek(); + Assert.assertTrue(COMPARATOR.compare(peek, pop) >= 0); + } + } + Assert.assertEquals(38, count); + } + + @Test + public void TripleLongPriorityDelayedIndexQueueTest() { + + @Cleanup + TripleLongPriorityQueue queue = new TripleLongPriorityQueue(); + for (int i = 0; i < 10; i++) { + queue.add(i, 1, 1); + } + + TripleLongPriorityDelayedIndexQueue delayedIndexQueue = TripleLongPriorityDelayedIndexQueue.wrap(queue); + + int count = 0; + while (!delayedIndexQueue.isEmpty()) { + DelayedIndex pop = delayedIndexQueue.pop(); + count++; + if (!delayedIndexQueue.isEmpty()) { + DelayedIndex peek = delayedIndexQueue.peek(); + Assert.assertTrue(COMPARATOR.compare(peek, pop) >= 0); + } + } + + Assert.assertEquals(10, count); + } +} From c27a1943dcdd76e7a1dd081ec410a14ab2f0b7fa Mon Sep 17 00:00:00 2001 From: coderzc Date: Thu, 5 Jan 2023 16:16:50 +0800 Subject: [PATCH 2/7] fix license --- .../broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java | 2 +- .../apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java | 2 +- .../delayed/bucket/TripleLongPriorityDelayedIndexQueue.java | 2 +- .../org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java index fd7b7884cd69a..30acdc39f312d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java @@ -1,4 +1,4 @@ -/** +/* * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java index df80c21767b12..318da118a1f68 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java @@ -1,4 +1,4 @@ -/** +/* * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java index 97ff958188295..1aef3baf0adda 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java @@ -1,4 +1,4 @@ -/** +/* * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java index 357d14fad7438..64a6bb7118452 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java @@ -1,4 +1,4 @@ -/** +/* * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information From 3d96a699c9b1391071bbf471965e21d108aa9ee2 Mon Sep 17 00:00:00 2001 From: coderzc Date: Thu, 5 Jan 2023 18:16:37 +0800 Subject: [PATCH 3/7] fix ConcurrentModificationException --- .../bucket/BucketDelayedDeliveryTracker.java | 20 ++++++++++++------- 1 file changed, 13 insertions(+), 7 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 0f28bff46b6f6..a59c405290d14 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 @@ -33,10 +33,12 @@ import java.util.ArrayList; import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.NavigableSet; import java.util.Optional; import java.util.TreeSet; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; @@ -126,16 +128,15 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { } List> futures = new ArrayList<>(immutableBuckets.asMapOfRanges().size()); - for (ImmutableBucket immutableBucket : immutableBuckets.asMapOfRanges().values()) { + Map, ImmutableBucket> toBeDeletedBucketMap = new ConcurrentHashMap<>(); + for (Map.Entry, ImmutableBucket> entry :immutableBuckets.asMapOfRanges().entrySet()) { + Range key = entry.getKey(); + ImmutableBucket immutableBucket = entry.getValue(); CompletableFuture future = immutableBucket.asyncRecoverBucketSnapshotEntry(this::getCutoffTime).thenAccept(indexList -> { if (CollectionUtils.isEmpty(indexList)) { // Delete bucket snapshot if indexList is empty - synchronized (immutableBuckets) { - immutableBuckets.remove( - Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId)); - } - immutableBucket.asyncDeleteBucketSnapshot(); + toBeDeletedBucketMap.put(key, immutableBucket); return; } DelayedIndex lastDelayedIndex = indexList.get(indexList.size() - 1); @@ -154,7 +155,12 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { } try { - FutureUtil.waitForAll(futures).get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); + FutureUtil.waitForAll(futures).whenComplete((__, ex) -> { + toBeDeletedBucketMap.forEach((k, immutableBucket) -> { + immutableBuckets.remove(k); + immutableBucket.asyncDeleteBucketSnapshot(); + }); + }).get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); } catch (InterruptedException | ExecutionException | TimeoutException e) { throw new RuntimeException(e); } From 621a8e82ad2a31926938593db993820ebe402f42 Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 11 Jan 2023 20:33:13 +0800 Subject: [PATCH 4/7] Address comment --- .../delayed/bucket/CombinedSegmentDelayedIndexQueue.java | 5 +++-- .../pulsar/broker/delayed/bucket/DelayedIndexQueue.java | 3 +-- .../bucket/TripleLongPriorityDelayedIndexQueue.java | 2 +- .../{ => bucket}/BucketDelayedDeliveryTrackerTest.java | 8 +++++--- .../delayed/{ => bucket}/DelayedIndexQueueTest.java | 4 +--- 5 files changed, 11 insertions(+), 11 deletions(-) rename pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/{ => bucket}/BucketDelayedDeliveryTrackerTest.java (97%) rename pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/{ => bucket}/DelayedIndexQueueTest.java (96%) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java index 30acdc39f312d..9b5a07364b3d2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java @@ -62,12 +62,13 @@ public DelayedIndex pop() { } private DelayedIndex getValue(boolean needAdvanceCursor) { + // skip empty segment while (segmentListACursor < segmentListA.size() - && segmentACursor >= segmentListA.get(segmentListACursor).getIndexesCount()) { + && segmentListA.get(segmentListACursor).getIndexesCount() == 0) { segmentListACursor++; } while (segmentListBCursor < segmentListB.size() - && segmentBCursor >= segmentListB.get(segmentListBCursor).getIndexesCount()) { + && segmentListB.get(segmentListBCursor).getIndexesCount() == 0) { segmentListBCursor++; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java index 318da118a1f68..dee476c376ec9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueue.java @@ -22,8 +22,7 @@ import java.util.Objects; import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat; -public interface DelayedIndexQueue { - +interface DelayedIndexQueue { Comparator COMPARATOR = (o1, o2) -> { if (!Objects.equals(o1.getTimestamp(), o2.getTimestamp())) { return Long.compare(o1.getTimestamp(), o2.getTimestamp()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java index 1aef3baf0adda..b8d54bd78b428 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/TripleLongPriorityDelayedIndexQueue.java @@ -23,7 +23,7 @@ import org.apache.pulsar.common.util.collections.TripleLongPriorityQueue; @NotThreadSafe -public class TripleLongPriorityDelayedIndexQueue implements DelayedIndexQueue { +class TripleLongPriorityDelayedIndexQueue implements DelayedIndexQueue { private final TripleLongPriorityQueue queue; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java similarity index 97% rename from pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java rename to pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java index 5be37ec3ef5e9..0a2a76ec339eb 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BucketDelayedDeliveryTrackerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java @@ -16,7 +16,7 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.delayed; +package org.apache.pulsar.broker.delayed.bucket; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; @@ -42,8 +42,10 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.pulsar.broker.delayed.bucket.BucketDelayedDeliveryTracker; -import org.apache.pulsar.broker.delayed.bucket.BucketSnapshotStorage; +import org.apache.pulsar.broker.delayed.AbstractDeliveryTrackerTest; +import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; +import org.apache.pulsar.broker.delayed.MockBucketSnapshotStorage; +import org.apache.pulsar.broker.delayed.MockManagedCursor; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.roaringbitmap.RoaringBitmap; import org.roaringbitmap.buffer.ImmutableRoaringBitmap; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueueTest.java similarity index 96% rename from pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java rename to pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueueTest.java index 64a6bb7118452..865ccb6934ab7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/DelayedIndexQueueTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/DelayedIndexQueueTest.java @@ -16,15 +16,13 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.delayed; +package org.apache.pulsar.broker.delayed.bucket; import static org.apache.pulsar.broker.delayed.bucket.DelayedIndexQueue.COMPARATOR; import java.util.ArrayList; import java.util.List; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; -import org.apache.pulsar.broker.delayed.bucket.CombinedSegmentDelayedIndexQueue; -import org.apache.pulsar.broker.delayed.bucket.TripleLongPriorityDelayedIndexQueue; import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.DelayedIndex; import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.SnapshotSegment; import org.apache.pulsar.common.util.collections.TripleLongPriorityQueue; From 14b5ed2fcae5192fcc0d4d4fc9ff6b99b9f5a7c9 Mon Sep 17 00:00:00 2001 From: coderzc Date: Fri, 13 Jan 2023 19:07:23 +0800 Subject: [PATCH 5/7] Address comment --- .../bucket/BucketDelayedDeliveryTracker.java | 7 ++++--- .../broker/delayed/bucket/ImmutableBucket.java | 18 +++++++++++++++--- 2 files changed, 19 insertions(+), 6 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 a59c405290d14..5c22db2e1e177 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 @@ -237,6 +237,9 @@ public synchronized boolean addMessage(long ledgerId, long entryId, long deliver try { asyncMergeBucketSnapshot().get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); } catch (InterruptedException | ExecutionException | TimeoutException e) { + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } throw new RuntimeException(e); } } @@ -375,9 +378,7 @@ public synchronized NavigableSet getScheduledMessages(int maxMessa try { bucket.asyncLoadNextBucketSnapshotEntry().thenAccept(indexList -> { if (CollectionUtils.isEmpty(indexList)) { - synchronized (immutableBuckets) { - immutableBuckets.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId)); - } + immutableBuckets.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId)); bucket.asyncDeleteBucketSnapshot(); return; } 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 dd23620aceb4f..8348b4999ed80 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 @@ -24,7 +24,9 @@ import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.function.Supplier; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedCursor; @@ -137,8 +139,15 @@ CompletableFuture> } CompletableFuture asyncDeleteBucketSnapshot() { - return removeBucketCursorProperty(bucketKey()).thenCompose(__ -> - bucketSnapshotStorage.deleteBucketSnapshot(getAndUpdateBucketId())); + String bucketKey = bucketKey(); + long bucketId = getAndUpdateBucketId(); + return removeBucketCursorProperty(bucketKey).thenCompose(__ -> + bucketSnapshotStorage.deleteBucketSnapshot(bucketId)).whenComplete((__, ex) -> { + if (ex != null) { + log.warn("Failed to delete bucket snapshot, bucketId: {}, bucketKey: {}", + bucketId, bucketKey, ex); + } + }); } void clear(boolean delete) { @@ -148,7 +157,10 @@ void clear(boolean delete) { snapshotGenerateFuture.cancel(true); try { asyncDeleteBucketSnapshot().get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); - } catch (Exception e) { + } catch (InterruptedException | ExecutionException | TimeoutException e) { + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } throw new RuntimeException(e); } } else { From 90c6690f8923f33be90e51c7fb72ab3430ddff24 Mon Sep 17 00:00:00 2001 From: coderzc Date: Mon, 16 Jan 2023 16:23:18 +0800 Subject: [PATCH 6/7] remove public modifier --- .../broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java index 9b5a07364b3d2..3f89cc9fdfb15 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/CombinedSegmentDelayedIndexQueue.java @@ -24,7 +24,7 @@ import org.apache.pulsar.broker.delayed.proto.DelayedMessageIndexBucketSnapshotFormat.SnapshotSegment; @NotThreadSafe -public class CombinedSegmentDelayedIndexQueue implements DelayedIndexQueue { +class CombinedSegmentDelayedIndexQueue implements DelayedIndexQueue { private final List segmentListA; private final List segmentListB; From eb33c21069aeeddf407844584b1f7ce6987b4f4a Mon Sep 17 00:00:00 2001 From: coderzc Date: Mon, 16 Jan 2023 22:45:34 +0800 Subject: [PATCH 7/7] clean up old overlap buckets when recover buckets --- .../bucket/BucketDelayedDeliveryTracker.java | 52 +++++++++++++------ .../broker/delayed/bucket/MutableBucket.java | 4 +- 2 files changed, 38 insertions(+), 18 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 5c22db2e1e177..715123487d55c 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 @@ -111,6 +111,7 @@ public BucketDelayedDeliveryTracker(PersistentDispatcherMultipleConsumers dispat private synchronized long recoverBucketSnapshot() throws RuntimeException { ManagedCursor cursor = this.lastMutableBucket.cursor; + Map, ImmutableBucket> toBeDeletedBucketMap = new ConcurrentHashMap<>(); cursor.getCursorProperties().keySet().forEach(key -> { if (key.startsWith(DELAYED_BUCKET_KEY_PREFIX)) { String[] keys = key.split(DELIMITER); @@ -118,8 +119,8 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { ImmutableBucket immutableBucket = new ImmutableBucket(cursor, this.lastMutableBucket.bucketSnapshotStorage, Long.parseLong(keys[1]), Long.parseLong(keys[2])); - immutableBuckets.put(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId), - immutableBucket); + putAndCleanOverlapRange(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId), + immutableBucket, toBeDeletedBucketMap); } }); @@ -128,7 +129,6 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { } List> futures = new ArrayList<>(immutableBuckets.asMapOfRanges().size()); - Map, ImmutableBucket> toBeDeletedBucketMap = new ConcurrentHashMap<>(); for (Map.Entry, ImmutableBucket> entry :immutableBuckets.asMapOfRanges().entrySet()) { Range key = entry.getKey(); ImmutableBucket immutableBucket = entry.getValue(); @@ -157,7 +157,7 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { try { FutureUtil.waitForAll(futures).whenComplete((__, ex) -> { toBeDeletedBucketMap.forEach((k, immutableBucket) -> { - immutableBuckets.remove(k); + immutableBuckets.asMapOfRanges().remove(k); immutableBucket.asyncDeleteBucketSnapshot(); }); }).get(AsyncOperationTimeoutSeconds, TimeUnit.SECONDS); @@ -176,6 +176,26 @@ private synchronized long recoverBucketSnapshot() throws RuntimeException { return numberDelayedMessages.getValue(); } + private synchronized void putAndCleanOverlapRange(Range range, ImmutableBucket immutableBucket, + Map, ImmutableBucket> toBeDeletedBucketMap) { + RangeMap subRangeMap = immutableBuckets.subRangeMap(range); + boolean canPut = false; + if (!subRangeMap.asMapOfRanges().isEmpty()) { + for (Map.Entry, ImmutableBucket> rangeEntry : subRangeMap.asMapOfRanges().entrySet()) { + if (range.encloses(rangeEntry.getKey())) { + toBeDeletedBucketMap.put(rangeEntry.getKey(), rangeEntry.getValue()); + canPut = true; + } + } + } else { + canPut = true; + } + + if (canPut) { + immutableBuckets.put(range, immutableBucket); + } + } + @Override public void run(Timeout timeout) throws Exception { synchronized (this) { @@ -298,17 +318,19 @@ private synchronized CompletableFuture asyncMergeBucketSnapshot(ImmutableB CompletableFuture> futureB = bucketB.getRemainSnapshotSegment(); return futureA.thenCombine(futureB, CombinedSegmentDelayedIndexQueue::wrap) - .thenCompose(combinedDelayedIndexQueue -> { - CompletableFuture removeAFuture = bucketA.asyncDeleteBucketSnapshot(); - CompletableFuture removeBFuture = bucketB.asyncDeleteBucketSnapshot(); - - return CompletableFuture.allOf(removeAFuture, removeBFuture).thenRun(() -> { - Pair immutableBucketDelayedIndexPair = - lastMutableBucket.createImmutableBucketAndAsyncPersistent( - timeStepPerBucketSnapshotSegment, sharedBucketPriorityQueue, - combinedDelayedIndexQueue, bucketA.startLedgerId, bucketB.endLedgerId); - afterCreateImmutableBucket(immutableBucketDelayedIndexPair); - }); + .thenAccept(combinedDelayedIndexQueue -> { + Pair immutableBucketDelayedIndexPair = + lastMutableBucket.createImmutableBucketAndAsyncPersistent( + timeStepPerBucketSnapshotSegment, sharedBucketPriorityQueue, + combinedDelayedIndexQueue, bucketA.startLedgerId, bucketB.endLedgerId); + afterCreateImmutableBucket(immutableBucketDelayedIndexPair); + + immutableBucketDelayedIndexPair.getLeft().getSnapshotCreateFuture() + .orElse(CompletableFuture.completedFuture(null)).thenCompose(___ -> { + CompletableFuture removeAFuture = bucketA.asyncDeleteBucketSnapshot(); + CompletableFuture removeBFuture = bucketB.asyncDeleteBucketSnapshot(); + return CompletableFuture.allOf(removeAFuture, removeBFuture); + }); }); }); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java index e70ccf2e9831d..ad457329c427f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java @@ -137,9 +137,7 @@ Pair createImmutableBucketAndAsyncPersistent( bucketSnapshotMetadata, bucketSnapshotSegments); bucket.setSnapshotCreateFuture(future); future.whenComplete((__, ex) -> { - if (ex == null) { - bucket.setSnapshotCreateFuture(null); - } else { + if (ex != null) { //TODO Record create snapshot failed log.error("Failed to create snapshot: ", ex); }