From 505779f4d0adf0d9145f4ac13b2835ce87b077f9 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sun, 6 Nov 2022 23:38:32 +0800 Subject: [PATCH 1/9] [improve][broker] Reduce object allocation when iterate on `LongPairRangeSet` (#18356) --- .../mledger/impl/ManagedCursorImpl.java | 36 ++++++----- .../mledger/impl/RangeSetWrapper.java | 7 +++ .../ConcurrentOpenLongPairRangeSet.java | 37 ++++++++++- .../util/collections/LongPairRangeSet.java | 63 ++++++++++++++++++- 4 files changed, 120 insertions(+), 23 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 5e761a4fae275..efb55afc547aa 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -2834,25 +2834,27 @@ private List buildIndividualDeletedMessageRanges() { return Collections.emptyList(); } - MLDataFormats.NestedPositionInfo.Builder nestedPositionBuilder = MLDataFormats.NestedPositionInfo - .newBuilder(); - MLDataFormats.MessageRange.Builder messageRangeBuilder = MLDataFormats.MessageRange.newBuilder(); AtomicInteger acksSerializedSize = new AtomicInteger(0); List rangeList = new ArrayList<>(); - individualDeletedMessages.forEach((positionRange) -> { - PositionImpl p = positionRange.lowerEndpoint(); - nestedPositionBuilder.setLedgerId(p.getLedgerId()); - nestedPositionBuilder.setEntryId(p.getEntryId()); - messageRangeBuilder.setLowerEndpoint(nestedPositionBuilder.build()); - p = positionRange.upperEndpoint(); - nestedPositionBuilder.setLedgerId(p.getLedgerId()); - nestedPositionBuilder.setEntryId(p.getEntryId()); - messageRangeBuilder.setUpperEndpoint(nestedPositionBuilder.build()); - MessageRange messageRange = messageRangeBuilder.build(); - acksSerializedSize.addAndGet(messageRange.getSerializedSize()); - rangeList.add(messageRange); - return rangeList.size() <= config.getMaxUnackedRangesToPersist(); - }); + + individualDeletedMessages.forEachWithRangeBoundMapper( + (ledgerId, entryId) -> MLDataFormats.NestedPositionInfo.newBuilder() + .setLedgerId(ledgerId) + .setEntryId(entryId) + .build(), + null, + (lowerBound, upperBound) -> { + MessageRange messageRange = MLDataFormats.MessageRange.newBuilder() + .setLowerEndpoint(lowerBound) + .setUpperEndpoint(upperBound).build(); + + acksSerializedSize.addAndGet(messageRange.getSerializedSize()); + rangeList.add(messageRange); + + return rangeList.size() <= config.getMaxUnackedRangesToPersist(); + + }); + this.individualDeletedMessagesSerializedSize = acksSerializedSize.get(); individualDeletedMessages.resetDirtyKeys(); return rangeList; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java index 274c625e7459e..0ed35df3de8a8 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java @@ -118,6 +118,13 @@ public void forEach(RangeProcessor action, LongPairConsumer cons rangeSet.forEach(action, consumer); } + @Override + public void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundMapper, + RangeBoundConvertFunction rangeBoundMapper, + RangeBoundBiConsumer action) { + rangeSet.forEachWithRangeBoundMapper(rawRangeBoundMapper, rangeBoundMapper, action); + } + @Override public int size() { return rangeSet.size(); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java index 6faa61d3b374d..f2107305c2582 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java @@ -209,8 +209,10 @@ public void forEach(RangeProcessor action, LongPairConsumer cons int currentClosedMark = first; while (currentClosedMark != -1 && currentClosedMark <= last) { int nextOpenMark = set.nextClearBit(currentClosedMark); - Range range = Range.openClosed(consumer.apply(key, currentClosedMark - 1), - consumer.apply(key, nextOpenMark - 1)); + Range range = Range.openClosed( + consumer.apply(key, currentClosedMark - 1), + consumer.apply(key, nextOpenMark - 1) + ); if (!action.process(range)) { completed.set(true); break; @@ -220,6 +222,35 @@ public void forEach(RangeProcessor action, LongPairConsumer cons }); } + @Override + public void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundMapper, + RangeBoundConvertFunction __, + RangeBoundBiConsumer action) { + AtomicBoolean completed = new AtomicBoolean(false); + rangeBitSetMap.forEach((key, set) -> { + if (completed.get()) { + return; + } + if (set.isEmpty()) { + return; + } + int first = set.nextSetBit(0); + int last = set.previousSetBit(set.size()); + int currentClosedMark = first; + while (currentClosedMark != -1 && currentClosedMark <= last) { + int nextOpenMark = set.nextClearBit(currentClosedMark); + O lower = rawRangeBoundMapper.apply(key, currentClosedMark - 1); + O upper = rawRangeBoundMapper.apply(key, nextOpenMark - 1); + if (!action.process(lower, upper)) { + completed.set(true); + break; + } + currentClosedMark = set.nextSetBit(nextOpenMark); + } + }); + } + + @Override public Range firstRange() { if (rangeBitSetMap.isEmpty()) { @@ -269,7 +300,7 @@ public int cardinality(long lowerKey, long lowerValue, long upperKey, long upper public int size() { if (updatedAfterCachedForSize) { MutableInt size = new MutableInt(0); - forEach((range) -> { + forEachWithRangeBoundMapper((ledgerId, entryId) -> 0, null, (ignored1, ignored2) -> { size.increment(); return true; }); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java index 70a589cf55166..61f35ca77bd66 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java @@ -104,6 +104,25 @@ public interface LongPairRangeSet> { */ void forEach(RangeProcessor action, LongPairConsumer consumer); + /** + * Performs the given action for each entry in this map until all entries have been processed + * or action returns "false". Unless otherwise specified by the implementing class, + * actions are performed in the order of entry set iteration (if an iteration order is specified.) + * + * This method is optimized for reduce `Range` and `PositionImpl` object creation. + * Caller of this method can use either {@param rawRangeBoundConsumer} + * or {@param rangeBoundConvertFunction} to apply conversion directly + * on the raw LongPair like (long, long) + * or on the `T` (PositionImpl). + * + * Those convert function will apply on both bound of the range, then the result will pass to + * {@param action} to do iteration jobs. + * + */ + void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundConsumer, + RangeBoundConvertFunction rangeBoundConvertFunction, + RangeBoundBiConsumer action); + /** * Returns total number of ranges into the set. * @@ -135,7 +154,7 @@ public interface LongPairRangeSet> { * * @param the type of the result. */ - public interface LongPairConsumer { + interface LongPairConsumer { T apply(long key, long value); } @@ -143,7 +162,7 @@ public interface LongPairConsumer { * The interface exposing a method for processing of ranges. * @param - The incoming type of data in the range object. */ - public interface RangeProcessor> { + interface RangeProcessor> { /** * * @param range @@ -152,6 +171,30 @@ public interface RangeProcessor> { boolean process(Range range); } + /** + * The interface exposing a method for conversion each bound of ranges. + * this will be use in `forEachWithRangeBoundMapper` when iteration. + * @param - The incoming type of data in the range object. + * @param - The output type of data after apply the conversion. + */ + interface RangeBoundConvertFunction { + O apply(T rangeBound); + } + + /** + * The interface exposing a method for do iteration jobs after apply + * user define range bound conversion function. + * @param the input type of parameter return by `RangeBoundConvertFunction` + */ + interface RangeBoundBiConsumer { + /** + * + * @param range + * @return false if there is no further processing required + */ + boolean process(O rangeLowerBound, O rangeUpperBound); + } + /** * This class is a simple key-value data structure. */ @@ -270,7 +313,7 @@ public void forEach(RangeProcessor action) { } @Override - public void forEach(RangeProcessor action, LongPairConsumer consumer) { + public void forEach(RangeProcessor action, LongPairConsumer __) { for (Range range : asRanges()) { if (!action.process(range)) { break; @@ -278,6 +321,20 @@ public void forEach(RangeProcessor action, LongPairConsumer cons } } + @Override + public void forEachWithRangeBoundMapper(LongPairConsumer __, + RangeBoundConvertFunction rangeBoundMapper, + RangeBoundBiConsumer action) { + for (Range range : asRanges()) { + O lower = rangeBoundMapper.apply(range.lowerEndpoint()); + O upper = rangeBoundMapper.apply(range.upperEndpoint()); + if (!action.process(lower, upper)) { + break; + } + } + } + + @Override public boolean contains(long key, long value) { return this.contains(consumer.apply(key, value)); From 775c7fe655adf62fd4da7f921573db0204fe95ff Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sun, 6 Nov 2022 23:57:20 +0800 Subject: [PATCH 2/9] fix npe --- .../apache/bookkeeper/mledger/impl/ManagedCursorImpl.java | 7 ++++++- .../util/collections/ConcurrentOpenLongPairRangeSet.java | 4 +++- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index efb55afc547aa..df7b75e60e7e3 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -2838,11 +2838,16 @@ private List buildIndividualDeletedMessageRanges() { List rangeList = new ArrayList<>(); individualDeletedMessages.forEachWithRangeBoundMapper( + // conversion function for `ConcurrentOpenLongPairRangeSet` (ledgerId, entryId) -> MLDataFormats.NestedPositionInfo.newBuilder() .setLedgerId(ledgerId) .setEntryId(entryId) .build(), - null, + // conversion function for `LongPairRangeSet.DefaultRangeSet` + (positionImpl) -> MLDataFormats.NestedPositionInfo.newBuilder() + .setLedgerId(positionImpl.getLedgerId()) + .setEntryId(positionImpl.getEntryId()) + .build(), (lowerBound, upperBound) -> { MessageRange messageRange = MLDataFormats.MessageRange.newBuilder() .setLowerEndpoint(lowerBound) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java index f2107305c2582..83d48439d9038 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java @@ -300,7 +300,9 @@ public int cardinality(long lowerKey, long lowerValue, long upperKey, long upper public int size() { if (updatedAfterCachedForSize) { MutableInt size = new MutableInt(0); - forEachWithRangeBoundMapper((ledgerId, entryId) -> 0, null, (ignored1, ignored2) -> { + + // ignore result because we just want to count + forEachWithRangeBoundMapper((ledgerId, entryId) -> 0, __ -> 0, (ignored1, ignored2) -> { size.increment(); return true; }); From fc18417e85abf69f14fbc39bd79638cba53c5c4f Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Mon, 7 Nov 2022 00:30:36 +0800 Subject: [PATCH 3/9] add unit test --- .../ConcurrentOpenLongPairRangeSetTest.java | 35 +++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java index 3037a9deba356..30a3b377ef28b 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -26,7 +26,9 @@ import java.util.ArrayList; import java.util.List; import java.util.Set; +import java.util.function.Function; +import org.apache.commons.lang.mutable.MutableInt; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPair; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; import org.testng.annotations.Test; @@ -480,4 +482,37 @@ public void testCardinality() { v = set.cardinality(1, 0, 3, 30); assertEquals(v, 80 + 31); } + + @Test + public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + set.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); + set.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); + set.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); + set.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); + set.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); + + + MutableInt size = new MutableInt(0); + + List forEachIterResult = new ArrayList<>(); + set.forEach((range) -> { + forEachIterResult.add(range.lowerEndpoint()); + forEachIterResult.add(range.upperEndpoint()); + + size.increment(); + return true; + }); + + List forEachIterWithRangeBoundMapperResult = new ArrayList<>(); + + set.forEachWithRangeBoundMapper(LongPair::new, (pair) -> pair, (rangeLowerBound, rangeUpperBound) -> { + forEachIterWithRangeBoundMapperResult.add(rangeLowerBound); + forEachIterWithRangeBoundMapperResult.add(rangeUpperBound); + return true; + }); + + assertEquals(forEachIterResult, forEachIterWithRangeBoundMapperResult); + assertEquals(size.intValue(), set.size()); + } } From 40a4d531fd5ce31fbff9fe6fe0be056fd14a2f36 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Mon, 7 Nov 2022 00:59:10 +0800 Subject: [PATCH 4/9] add unit test --- .../ConcurrentOpenLongPairRangeSetTest.java | 26 +++++++++++++++---- 1 file changed, 21 insertions(+), 5 deletions(-) diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java index 30a3b377ef28b..68f0ff2cb5f5e 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -486,11 +486,17 @@ public void testCardinality() { @Test public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); - set.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); - set.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); - set.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); - set.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); - set.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); + LongPairRangeSet.DefaultRangeSet defaultRangeSet = new LongPairRangeSet.DefaultRangeSet<>(consumer); + + set.addOpenClosed(1, 10, 1, 15); + set.addOpenClosed(2, 25, 2, 28); + set.addOpenClosed(3, 12, 3, 20); + set.addOpenClosed(4, 12, 4, 20); + + defaultRangeSet.addOpenClosed(1, 10, 1, 15); + defaultRangeSet.addOpenClosed(2, 25, 2, 28); + defaultRangeSet.addOpenClosed(3, 12, 3, 20); + defaultRangeSet.addOpenClosed(4, 12, 4, 20); MutableInt size = new MutableInt(0); @@ -504,8 +510,15 @@ public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { return true; }); + List defaultRangeSetResult = new ArrayList<>(); List forEachIterWithRangeBoundMapperResult = new ArrayList<>(); + defaultRangeSet.forEachWithRangeBoundMapper(LongPair::new, (pair) -> pair, (rangeLowerBound, rangeUpperBound) -> { + defaultRangeSetResult.add(rangeLowerBound); + defaultRangeSetResult.add(rangeUpperBound); + return true; + }); + set.forEachWithRangeBoundMapper(LongPair::new, (pair) -> pair, (rangeLowerBound, rangeUpperBound) -> { forEachIterWithRangeBoundMapperResult.add(rangeLowerBound); forEachIterWithRangeBoundMapperResult.add(rangeUpperBound); @@ -513,6 +526,9 @@ public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { }); assertEquals(forEachIterResult, forEachIterWithRangeBoundMapperResult); + assertEquals(forEachIterResult, defaultRangeSetResult); + assertEquals(size.intValue(), set.size()); + assertEquals(size.intValue(), defaultRangeSet.size()); } } From 94bb2751e480a15fcf73749d92d7393404b57d18 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Tue, 8 Nov 2022 02:30:16 +0800 Subject: [PATCH 5/9] reuse builder --- .../mledger/impl/ManagedCursorImpl.java | 43 ++++++++------- .../mledger/impl/RangeSetWrapper.java | 13 +++-- .../mledger/impl/RangeSetWrapperTest.java | 35 +++++++------ .../ConcurrentOpenLongPairRangeSet.java | 43 +++++---------- .../util/collections/LongPairRangeSet.java | 52 +++++++++---------- .../ConcurrentOpenLongPairRangeSetTest.java | 23 ++++---- .../util/collections/DefaultRangeSetTest.java | 4 +- 7 files changed, 97 insertions(+), 116 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index df7b75e60e7e3..7b962e46d8e6c 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -98,6 +98,7 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.common.util.collections.BitSetRecyclable; import org.apache.pulsar.common.util.collections.LongPairRangeSet; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeBoundConsumer; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; import org.apache.pulsar.metadata.api.Stat; import org.slf4j.Logger; @@ -181,6 +182,10 @@ public class ManagedCursorImpl implements ManagedCursor { private volatile ManagedCursorInfo managedCursorInfo; private static final LongPairConsumer positionRangeConverter = PositionImpl::new; + + private static final RangeBoundConsumer positionRangeReverseConverter = + (position) -> new LongPairRangeSet.LongPair(position.ledgerId, position.entryId); + private static final LongPairConsumer recyclePositionRangeConverter = (key, value) -> { PositionImplRecyclable position = PositionImplRecyclable.create(); position.ledgerId = key; @@ -294,7 +299,8 @@ public interface VoidCallback { this.config = config; this.ledger = ledger; this.name = cursorName; - this.individualDeletedMessages = new RangeSetWrapper<>(positionRangeConverter, this); + this.individualDeletedMessages = new RangeSetWrapper<>(positionRangeConverter, + positionRangeReverseConverter, this); if (config.isDeletionAtBatchIndexLevelEnabled()) { this.batchDeletedIndexes = new ConcurrentSkipListMap<>(); } else { @@ -2834,31 +2840,24 @@ private List buildIndividualDeletedMessageRanges() { return Collections.emptyList(); } + MLDataFormats.NestedPositionInfo.Builder nestedPositionBuilder = MLDataFormats.NestedPositionInfo + .newBuilder(); + + MLDataFormats.MessageRange.Builder messageRangeBuilder = MLDataFormats.MessageRange.newBuilder(); AtomicInteger acksSerializedSize = new AtomicInteger(0); List rangeList = new ArrayList<>(); - individualDeletedMessages.forEachWithRangeBoundMapper( - // conversion function for `ConcurrentOpenLongPairRangeSet` - (ledgerId, entryId) -> MLDataFormats.NestedPositionInfo.newBuilder() - .setLedgerId(ledgerId) - .setEntryId(entryId) - .build(), - // conversion function for `LongPairRangeSet.DefaultRangeSet` - (positionImpl) -> MLDataFormats.NestedPositionInfo.newBuilder() - .setLedgerId(positionImpl.getLedgerId()) - .setEntryId(positionImpl.getEntryId()) - .build(), - (lowerBound, upperBound) -> { - MessageRange messageRange = MLDataFormats.MessageRange.newBuilder() - .setLowerEndpoint(lowerBound) - .setUpperEndpoint(upperBound).build(); - - acksSerializedSize.addAndGet(messageRange.getSerializedSize()); - rangeList.add(messageRange); - - return rangeList.size() <= config.getMaxUnackedRangesToPersist(); + individualDeletedMessages.forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> { + MLDataFormats.NestedPositionInfo lowerPosition = nestedPositionBuilder.setLedgerId(lowerKey).setEntryId(lowerValue).build(); + MLDataFormats.NestedPositionInfo upperPosition = nestedPositionBuilder.setLedgerId(lowerKey).setEntryId(lowerValue).build(); - }); + MessageRange messageRange = messageRangeBuilder.setLowerEndpoint(lowerPosition).setUpperEndpoint(upperPosition).build(); + + acksSerializedSize.addAndGet(messageRange.getSerializedSize()); + rangeList.add(messageRange); + + return rangeList.size() <= config.getMaxUnackedRangesToPersist(); + }); this.individualDeletedMessagesSerializedSize = acksSerializedSize.get(); individualDeletedMessages.resetDirtyKeys(); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java index 0ed35df3de8a8..4d5f1e79d18bf 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java @@ -45,15 +45,16 @@ public class RangeSetWrapper> implements LongPairRangeSe * Record which Ledger is dirty. */ private final DefaultRangeSet dirtyLedgers = new LongPairRangeSet.DefaultRangeSet<>( - (LongPairConsumer) (key, value) -> key); + (LongPairConsumer) (key, value) -> key, + (RangeBoundConsumer) key -> new LongPair(key, 0)); - public RangeSetWrapper(LongPairConsumer rangeConverter, ManagedCursorImpl managedCursor) { + public RangeSetWrapper(LongPairConsumer rangeConverter, RangeBoundConsumer rangeBoundConsumer, ManagedCursorImpl managedCursor) { requireNonNull(managedCursor); this.config = managedCursor.getConfig(); this.rangeConverter = rangeConverter; this.rangeSet = config.isUnackedRangesOpenCacheSetEnabled() ? new ConcurrentOpenLongPairRangeSet<>(4096, rangeConverter) - : new LongPairRangeSet.DefaultRangeSet<>(rangeConverter); + : new LongPairRangeSet.DefaultRangeSet<>(rangeConverter, rangeBoundConsumer); this.enableMultiEntry = config.isPersistentUnackedRangesWithMultipleEntriesEnabled(); } @@ -119,10 +120,8 @@ public void forEach(RangeProcessor action, LongPairConsumer cons } @Override - public void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundMapper, - RangeBoundConvertFunction rangeBoundMapper, - RangeBoundBiConsumer action) { - rangeSet.forEachWithRangeBoundMapper(rawRangeBoundMapper, rangeBoundMapper, action); + public void forEachRawRange(RawRangeProcessor action) { + rangeSet.forEachRawRange(action); } @Override diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java index c374b034054da..03925fdeb3ab2 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java @@ -31,6 +31,7 @@ import java.util.List; import java.util.Set; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeBoundConsumer; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPair; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; import org.testng.annotations.AfterMethod; @@ -40,6 +41,8 @@ public class RangeSetWrapperTest { static final LongPairConsumer consumer = (key, value) -> new LongPair(key, value); + static final RangeBoundConsumer reverseConvert = (pair) -> pair; + ManagedLedgerImpl managedLedger; RangeSetWrapper set; ManagedLedgerConfig managedLedgerConfig; @@ -67,7 +70,7 @@ public void clean() throws Exception { @Test public void testDirtyLedger() { - RangeSetWrapper rangeSetWrapper = new RangeSetWrapper<>(consumer, managedCursor); + RangeSetWrapper rangeSetWrapper = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); // Test add range rangeSetWrapper.addOpenClosed(10, 0, 20, 0); assertEquals(rangeSetWrapper.size(), 1); @@ -93,7 +96,7 @@ public void testAddForSameKey() { } private void doTestAddForSameKey() { - set = new RangeSetWrapper(consumer, managedCursor); + set = new RangeSetWrapper(consumer, reverseConvert, managedCursor); // add 0 to 5 set.addOpenClosed(0, 0, 0, 5); // add 8,9,10 @@ -113,7 +116,7 @@ private void doTestAddForSameKey() { @Test public void testAddForDifferentKey() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); // [98,100],[(1,5),(1,5)],[(1,10,1,15)],[(1,20),(1,20)],[(2,0),(2,10)] set.addOpenClosed(0, 98, 0, 99); set.addOpenClosed(0, 100, 1, 5); @@ -131,7 +134,7 @@ public void testAddForDifferentKey() { @Test public void testAddForDifferentKey2() { managedLedgerConfig.setUnackedRangesOpenCacheSetEnabled(false); - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); // [98,100],[(1,5),(1,5)],[(1,10,1,15)],[(1,20),(1,20)],[(2,0),(2,10)] set.addOpenClosed(0, 98, 0, 99); set.addOpenClosed(0, 100, 1, 5); @@ -148,7 +151,7 @@ public void testAddForDifferentKey2() { @Test public void testAddCompareCompareWithGuava() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); // add 10K values for key 0 @@ -187,7 +190,7 @@ public void testAddCompareCompareWithGuava() { @Test public void testDeleteCompareWithGuava() throws Exception { - RangeSetWrapper set = new RangeSetWrapper<>(consumer, managedCursor); + RangeSetWrapper set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); // add 10K values for key 0 @@ -241,7 +244,7 @@ public void testDeleteCompareWithGuava() throws Exception { @Test public void testSpanWithGuava() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); set.addOpenClosed(0, 97, 0, 99); gSet.add(Range.openClosed(new LongPair(0, 97), new LongPair(0, 99))); @@ -266,7 +269,7 @@ public void testSpanWithGuava() { @Test public void testFirstRange() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); assertNull(set.firstRange()); set.addOpenClosed(0, 97, 0, 99); assertEquals(set.firstRange(), Range.openClosed(new LongPair(0, 97), new LongPair(0, 99))); @@ -281,7 +284,7 @@ public void testFirstRange() { @Test public void testLastRange() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); assertNull(set.lastRange()); Range range = Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)); set.addOpenClosed(0, 97, 0, 99); @@ -302,7 +305,7 @@ public void testLastRange() { @Test public void testToString() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.addOpenClosed(0, 97, 0, 99); assertEquals(set.toString(), "[(0:97..0:99]]"); set.addOpenClosed(0, 98, 0, 105); @@ -313,7 +316,7 @@ public void testToString() { @Test public void testDeleteForDifferentKey() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.addOpenClosed(0, 97, 0, 99); set.addOpenClosed(0, 99, 1, 5); set.addOpenClosed(1, 9, 1, 15); @@ -344,7 +347,7 @@ public void testDeleteForDifferentKey() { @Test public void testDeleteWithAtMost() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.addOpenClosed(0, 98, 0, 99); set.addOpenClosed(0, 100, 1, 5); set.addOpenClosed(1, 10, 1, 15); @@ -370,7 +373,7 @@ public void testDeleteWithAtMost() { @Test public void testDeleteWithAtMost2() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.addOpenClosed(0, 98, 0, 99); set.addOpenClosed(0, 100, 1, 5); set.addOpenClosed(1, 10, 1, 15); @@ -390,7 +393,7 @@ public void testDeleteWithAtMost2() { assertEquals(ranges.get(count), (Range.openClosed(new LongPair(2, 25), new LongPair(2, 28)))); managedLedgerConfig.setUnackedRangesOpenCacheSetEnabled(false); - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.addOpenClosed(0, 98, 0, 99); set.addOpenClosed(0, 100, 1, 5); set.addOpenClosed(1, 10, 1, 15); @@ -411,7 +414,7 @@ public void testDeleteWithAtMost2() { @Test public void testDeleteWithLeastMost() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.addOpenClosed(0, 98, 0, 99); set.addOpenClosed(0, 100, 1, 5); set.addOpenClosed(1, 10, 1, 15); @@ -439,7 +442,7 @@ public void testDeleteWithLeastMost() { @Test public void testRangeContaining() { - set = new RangeSetWrapper<>(consumer, managedCursor); + set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java index 83d48439d9038..ad5fbb9442a14 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java @@ -195,37 +195,18 @@ public void forEach(RangeProcessor action) { } @Override - public void forEach(RangeProcessor action, LongPairConsumer consumer) { - AtomicBoolean completed = new AtomicBoolean(false); - rangeBitSetMap.forEach((key, set) -> { - if (completed.get()) { - return; - } - if (set.isEmpty()) { - return; - } - int first = set.nextSetBit(0); - int last = set.previousSetBit(set.size()); - int currentClosedMark = first; - while (currentClosedMark != -1 && currentClosedMark <= last) { - int nextOpenMark = set.nextClearBit(currentClosedMark); - Range range = Range.openClosed( - consumer.apply(key, currentClosedMark - 1), - consumer.apply(key, nextOpenMark - 1) - ); - if (!action.process(range)) { - completed.set(true); - break; - } - currentClosedMark = set.nextSetBit(nextOpenMark); - } + public void forEach(RangeProcessor action, LongPairConsumer consumerParam) { + forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> { + Range range = Range.openClosed( + consumerParam.apply(lowerKey, lowerValue), + consumerParam.apply(upperKey, upperValue) + ); + return action.process(range); }); } @Override - public void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundMapper, - RangeBoundConvertFunction __, - RangeBoundBiConsumer action) { + public void forEachRawRange(RawRangeProcessor processor) { AtomicBoolean completed = new AtomicBoolean(false); rangeBitSetMap.forEach((key, set) -> { if (completed.get()) { @@ -239,9 +220,8 @@ public void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundMap int currentClosedMark = first; while (currentClosedMark != -1 && currentClosedMark <= last) { int nextOpenMark = set.nextClearBit(currentClosedMark); - O lower = rawRangeBoundMapper.apply(key, currentClosedMark - 1); - O upper = rawRangeBoundMapper.apply(key, nextOpenMark - 1); - if (!action.process(lower, upper)) { + if (!processor.processRawRange(key, currentClosedMark - 1, + key, nextOpenMark - 1)) { completed.set(true); break; } @@ -302,10 +282,11 @@ public int size() { MutableInt size = new MutableInt(0); // ignore result because we just want to count - forEachWithRangeBoundMapper((ledgerId, entryId) -> 0, __ -> 0, (ignored1, ignored2) -> { + forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> { size.increment(); return true; }); + cachedSize = size.intValue(); updatedAfterCachedForSize = false; } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java index 61f35ca77bd66..661decf946a4b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java @@ -119,9 +119,7 @@ public interface LongPairRangeSet> { * {@param action} to do iteration jobs. * */ - void forEachWithRangeBoundMapper(LongPairConsumer rawRangeBoundConsumer, - RangeBoundConvertFunction rangeBoundConvertFunction, - RangeBoundBiConsumer action); + void forEachRawRange(RawRangeProcessor action); /** * Returns total number of ranges into the set. @@ -158,6 +156,15 @@ interface LongPairConsumer { T apply(long key, long value); } + /** + * Represents a function that accepts result and produces a LongPair. + * Reverse ops of `LongPairConsumer` + * @param the type of the result. + */ + interface RangeBoundConsumer { + LongPair apply(T bound); + } + /** * The interface exposing a method for processing of ranges. * @param - The incoming type of data in the range object. @@ -172,27 +179,15 @@ interface RangeProcessor> { } /** - * The interface exposing a method for conversion each bound of ranges. - * this will be use in `forEachWithRangeBoundMapper` when iteration. + * The interface exposing a method for processing raw form of ranges. + * This method will omit the process to convert (long, long) to `T` + * create less object during the iteration. * @param - The incoming type of data in the range object. * @param - The output type of data after apply the conversion. */ - interface RangeBoundConvertFunction { - O apply(T rangeBound); - } - - /** - * The interface exposing a method for do iteration jobs after apply - * user define range bound conversion function. - * @param the input type of parameter return by `RangeBoundConvertFunction` - */ - interface RangeBoundBiConsumer { - /** - * - * @param range - * @return false if there is no further processing required - */ - boolean process(O rangeLowerBound, O rangeUpperBound); + interface RawRangeProcessor { + boolean processRawRange(long lowerKey, long lowerValue, + long upperKey, long upperValue); } /** @@ -251,9 +246,11 @@ class DefaultRangeSet> implements LongPairRangeSet { RangeSet set = TreeRangeSet.create(); private final LongPairConsumer consumer; + private final RangeBoundConsumer rangeEndPointConsumer; - public DefaultRangeSet(LongPairConsumer consumer) { + public DefaultRangeSet(LongPairConsumer consumer, RangeBoundConsumer reverseConsumer) { this.consumer = consumer; + this.rangeEndPointConsumer = reverseConsumer; } @Override @@ -322,13 +319,12 @@ public void forEach(RangeProcessor action, LongPairConsumer __) } @Override - public void forEachWithRangeBoundMapper(LongPairConsumer __, - RangeBoundConvertFunction rangeBoundMapper, - RangeBoundBiConsumer action) { + public void forEachRawRange(RawRangeProcessor action) { for (Range range : asRanges()) { - O lower = rangeBoundMapper.apply(range.lowerEndpoint()); - O upper = rangeBoundMapper.apply(range.upperEndpoint()); - if (!action.process(lower, upper)) { + LongPair lowerEndpoint = this.rangeEndPointConsumer.apply(range.lowerEndpoint()); + LongPair upperEndpoint = this.rangeEndPointConsumer.apply(range.upperEndpoint()); + if (!action.processRawRange(lowerEndpoint.key, lowerEndpoint.value, + upperEndpoint.key, upperEndpoint.value)) { break; } } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java index 68f0ff2cb5f5e..05465e961f152 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -26,10 +26,10 @@ import java.util.ArrayList; import java.util.List; import java.util.Set; -import java.util.function.Function; import org.apache.commons.lang.mutable.MutableInt; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPair; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeBoundConsumer; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; import org.testng.annotations.Test; @@ -39,7 +39,8 @@ public class ConcurrentOpenLongPairRangeSetTest { - static final LongPairConsumer consumer = (key, value) -> new LongPair(key, value); + static final LongPairConsumer consumer = LongPair::new; + static final RangeBoundConsumer reverseConsumer = pair -> pair; @Test public void testIsEmpty() { @@ -486,7 +487,7 @@ public void testCardinality() { @Test public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); - LongPairRangeSet.DefaultRangeSet defaultRangeSet = new LongPairRangeSet.DefaultRangeSet<>(consumer); + LongPairRangeSet.DefaultRangeSet defaultRangeSet = new LongPairRangeSet.DefaultRangeSet<>(consumer, reverseConsumer); set.addOpenClosed(1, 10, 1, 15); set.addOpenClosed(2, 25, 2, 28); @@ -511,21 +512,21 @@ public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { }); List defaultRangeSetResult = new ArrayList<>(); - List forEachIterWithRangeBoundMapperResult = new ArrayList<>(); + List forEachRawRangeResult = new ArrayList<>(); - defaultRangeSet.forEachWithRangeBoundMapper(LongPair::new, (pair) -> pair, (rangeLowerBound, rangeUpperBound) -> { - defaultRangeSetResult.add(rangeLowerBound); - defaultRangeSetResult.add(rangeUpperBound); + defaultRangeSet.forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> { + defaultRangeSetResult.add(new LongPair(lowerKey, lowerValue)); + defaultRangeSetResult.add(new LongPair(upperKey, upperValue)); return true; }); - set.forEachWithRangeBoundMapper(LongPair::new, (pair) -> pair, (rangeLowerBound, rangeUpperBound) -> { - forEachIterWithRangeBoundMapperResult.add(rangeLowerBound); - forEachIterWithRangeBoundMapperResult.add(rangeUpperBound); + set.forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> { + forEachRawRangeResult.add(new LongPair(lowerKey, lowerValue)); + forEachRawRangeResult.add(new LongPair(upperKey, upperValue)); return true; }); - assertEquals(forEachIterResult, forEachIterWithRangeBoundMapperResult); + assertEquals(forEachIterResult, forEachRawRangeResult); assertEquals(forEachIterResult, defaultRangeSetResult); assertEquals(size.intValue(), set.size()); diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/DefaultRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/DefaultRangeSetTest.java index 57303013a112c..f6103061a420c 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/DefaultRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/DefaultRangeSetTest.java @@ -27,11 +27,13 @@ public class DefaultRangeSetTest { static final LongPairRangeSet.LongPairConsumer consumer = LongPairRangeSet.LongPair::new; + static final LongPairRangeSet.RangeBoundConsumer reverseConsumer = + pair -> pair; @Test public void testBehavior() { LongPairRangeSet.DefaultRangeSet set = - new LongPairRangeSet.DefaultRangeSet<>(consumer); + new LongPairRangeSet.DefaultRangeSet<>(consumer, reverseConsumer); ConcurrentOpenLongPairRangeSet rangeSet = new ConcurrentOpenLongPairRangeSet<>(consumer); From 265dec63f37639d84896a2b8c98c704e472f3403 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 12 Nov 2022 22:58:16 +0800 Subject: [PATCH 6/9] fix checkstyle and javadoc --- .../mledger/impl/ManagedCursorImpl.java | 22 ++++++++++++++----- .../mledger/impl/RangeSetWrapper.java | 4 +++- .../mledger/impl/RangeSetWrapperTest.java | 4 +++- .../util/collections/LongPairRangeSet.java | 8 +------ .../ConcurrentOpenLongPairRangeSetTest.java | 7 ++++-- 5 files changed, 29 insertions(+), 16 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 7b962e46d8e6c..fdb8765d9c5a5 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -2843,15 +2843,27 @@ private List buildIndividualDeletedMessageRanges() { MLDataFormats.NestedPositionInfo.Builder nestedPositionBuilder = MLDataFormats.NestedPositionInfo .newBuilder(); - MLDataFormats.MessageRange.Builder messageRangeBuilder = MLDataFormats.MessageRange.newBuilder(); + MLDataFormats.MessageRange.Builder messageRangeBuilder = MLDataFormats.MessageRange + .newBuilder(); + AtomicInteger acksSerializedSize = new AtomicInteger(0); List rangeList = new ArrayList<>(); individualDeletedMessages.forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> { - MLDataFormats.NestedPositionInfo lowerPosition = nestedPositionBuilder.setLedgerId(lowerKey).setEntryId(lowerValue).build(); - MLDataFormats.NestedPositionInfo upperPosition = nestedPositionBuilder.setLedgerId(lowerKey).setEntryId(lowerValue).build(); - - MessageRange messageRange = messageRangeBuilder.setLowerEndpoint(lowerPosition).setUpperEndpoint(upperPosition).build(); + MLDataFormats.NestedPositionInfo lowerPosition = nestedPositionBuilder + .setLedgerId(lowerKey) + .setEntryId(lowerValue) + .build(); + + MLDataFormats.NestedPositionInfo upperPosition = nestedPositionBuilder + .setLedgerId(lowerKey) + .setEntryId(lowerValue) + .build(); + + MessageRange messageRange = messageRangeBuilder + .setLowerEndpoint(lowerPosition) + .setUpperEndpoint(upperPosition) + .build(); acksSerializedSize.addAndGet(messageRange.getSerializedSize()); rangeList.add(messageRange); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java index 4d5f1e79d18bf..0fb670e13a26a 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java @@ -48,7 +48,9 @@ public class RangeSetWrapper> implements LongPairRangeSe (LongPairConsumer) (key, value) -> key, (RangeBoundConsumer) key -> new LongPair(key, 0)); - public RangeSetWrapper(LongPairConsumer rangeConverter, RangeBoundConsumer rangeBoundConsumer, ManagedCursorImpl managedCursor) { + public RangeSetWrapper(LongPairConsumer rangeConverter, + RangeBoundConsumer rangeBoundConsumer, + ManagedCursorImpl managedCursor) { requireNonNull(managedCursor); this.config = managedCursor.getConfig(); this.rangeConverter = rangeConverter; diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java index 03925fdeb3ab2..89fbc26d41ae6 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapperTest.java @@ -70,7 +70,9 @@ public void clean() throws Exception { @Test public void testDirtyLedger() { - RangeSetWrapper rangeSetWrapper = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor); + RangeSetWrapper rangeSetWrapper = new RangeSetWrapper<>(consumer, + reverseConvert, + managedCursor); // Test add range rangeSetWrapper.addOpenClosed(10, 0, 20, 0); assertEquals(rangeSetWrapper.size(), 1); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java index 661decf946a4b..1bacb8877eece 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java @@ -109,13 +109,7 @@ public interface LongPairRangeSet> { * or action returns "false". Unless otherwise specified by the implementing class, * actions are performed in the order of entry set iteration (if an iteration order is specified.) * - * This method is optimized for reduce `Range` and `PositionImpl` object creation. - * Caller of this method can use either {@param rawRangeBoundConsumer} - * or {@param rangeBoundConvertFunction} to apply conversion directly - * on the raw LongPair like (long, long) - * or on the `T` (PositionImpl). - * - * Those convert function will apply on both bound of the range, then the result will pass to + * This method is optimized on reducing intermediate object creation. * {@param action} to do iteration jobs. * */ diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java index 05465e961f152..40bb337935742 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -486,8 +486,11 @@ public void testCardinality() { @Test public void testForEachResultTheSameAsForEachWithRangeBoundMapper() { - ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); - LongPairRangeSet.DefaultRangeSet defaultRangeSet = new LongPairRangeSet.DefaultRangeSet<>(consumer, reverseConsumer); + ConcurrentOpenLongPairRangeSet set = + new ConcurrentOpenLongPairRangeSet<>(consumer); + + LongPairRangeSet.DefaultRangeSet defaultRangeSet = + new LongPairRangeSet.DefaultRangeSet<>(consumer, reverseConsumer); set.addOpenClosed(1, 10, 1, 15); set.addOpenClosed(2, 25, 2, 28); From 2429b912ad8fdcaee7964f5f4c2dce9581fe0f69 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 19 Nov 2022 12:25:19 +0800 Subject: [PATCH 7/9] fix checkstyle --- .../bookkeeper/mledger/impl/RangeSetWrapper.java | 2 +- .../collections/ConcurrentOpenLongPairRangeSet.java | 2 +- .../common/util/collections/LongPairRangeSet.java | 10 ++++------ 3 files changed, 6 insertions(+), 8 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java index 0fb670e13a26a..02e43504482d8 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java @@ -122,7 +122,7 @@ public void forEach(RangeProcessor action, LongPairConsumer cons } @Override - public void forEachRawRange(RawRangeProcessor action) { + public void forEachRawRange(RawRangeProcessor action) { rangeSet.forEachRawRange(action); } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java index ad5fbb9442a14..72215d7296cc3 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java @@ -206,7 +206,7 @@ public void forEach(RangeProcessor action, LongPairConsumer cons } @Override - public void forEachRawRange(RawRangeProcessor processor) { + public void forEachRawRange(RawRangeProcessor processor) { AtomicBoolean completed = new AtomicBoolean(false); rangeBitSetMap.forEach((key, set) -> { if (completed.get()) { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java index 1bacb8877eece..1f519d0476aaa 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java @@ -113,7 +113,7 @@ public interface LongPairRangeSet> { * {@param action} to do iteration jobs. * */ - void forEachRawRange(RawRangeProcessor action); + void forEachRawRange(RawRangeProcessor action); /** * Returns total number of ranges into the set. @@ -176,10 +176,8 @@ interface RangeProcessor> { * The interface exposing a method for processing raw form of ranges. * This method will omit the process to convert (long, long) to `T` * create less object during the iteration. - * @param - The incoming type of data in the range object. - * @param - The output type of data after apply the conversion. */ - interface RawRangeProcessor { + interface RawRangeProcessor { boolean processRawRange(long lowerKey, long lowerValue, long upperKey, long upperValue); } @@ -304,7 +302,7 @@ public void forEach(RangeProcessor action) { } @Override - public void forEach(RangeProcessor action, LongPairConsumer __) { + public void forEach(RangeProcessor action, LongPairConsumer outerConsumer) { for (Range range : asRanges()) { if (!action.process(range)) { break; @@ -313,7 +311,7 @@ public void forEach(RangeProcessor action, LongPairConsumer __) } @Override - public void forEachRawRange(RawRangeProcessor action) { + public void forEachRawRange(RawRangeProcessor action) { for (Range range : asRanges()) { LongPair lowerEndpoint = this.rangeEndPointConsumer.apply(range.lowerEndpoint()); LongPair upperEndpoint = this.rangeEndPointConsumer.apply(range.upperEndpoint()); From bb216025fe754f54e1396371db0eb5977c1fb825 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 19 Nov 2022 14:00:39 +0800 Subject: [PATCH 8/9] fix wrong input --- .../org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index fdb8765d9c5a5..0d8005ce5984e 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -98,8 +98,8 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.common.util.collections.BitSetRecyclable; import org.apache.pulsar.common.util.collections.LongPairRangeSet; -import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeBoundConsumer; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeBoundConsumer; import org.apache.pulsar.metadata.api.Stat; import org.slf4j.Logger; import org.slf4j.LoggerFactory; From 653a7e9e1650fa564d6e0404ed5f5465e83d9394 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 19 Nov 2022 15:37:44 +0800 Subject: [PATCH 9/9] fix wrong input param --- .../org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java | 4 ++-- .../pulsar/common/util/collections/LongPairRangeSet.java | 3 +++ 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 0d8005ce5984e..d726bfe9d23d5 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -2856,8 +2856,8 @@ private List buildIndividualDeletedMessageRanges() { .build(); MLDataFormats.NestedPositionInfo upperPosition = nestedPositionBuilder - .setLedgerId(lowerKey) - .setEntryId(lowerValue) + .setLedgerId(upperKey) + .setEntryId(upperValue) .build(); MessageRange messageRange = messageRangeBuilder diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java index 1f519d0476aaa..8aad5587dfd38 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java @@ -176,6 +176,9 @@ interface RangeProcessor> { * The interface exposing a method for processing raw form of ranges. * This method will omit the process to convert (long, long) to `T` * create less object during the iteration. + * the parameter is the same as {@linkplain RangeProcessor} which + * means (lowerKey,lowerValue) in open bound + * (upperKey, upperValue) in close bound in Range */ interface RawRangeProcessor { boolean processRawRange(long lowerKey, long lowerValue,