Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@
import org.apache.pulsar.common.util.collections.BitSetRecyclable;
import org.apache.pulsar.common.util.collections.LongPairRangeSet;
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;
Expand Down Expand Up @@ -181,6 +182,10 @@ public class ManagedCursorImpl implements ManagedCursor {
private volatile ManagedCursorInfo managedCursorInfo;

private static final LongPairConsumer<PositionImpl> positionRangeConverter = PositionImpl::new;

private static final RangeBoundConsumer<PositionImpl> positionRangeReverseConverter =
(position) -> new LongPairRangeSet.LongPair(position.ledgerId, position.entryId);

private static final LongPairConsumer<PositionImplRecyclable> recyclePositionRangeConverter = (key, value) -> {
PositionImplRecyclable position = PositionImplRecyclable.create();
position.ledgerId = key;
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -2836,23 +2842,35 @@ private List<MLDataFormats.MessageRange> 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<MessageRange> 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();

individualDeletedMessages.forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> {
MLDataFormats.NestedPositionInfo lowerPosition = nestedPositionBuilder
.setLedgerId(lowerKey)
.setEntryId(lowerValue)
.build();

MLDataFormats.NestedPositionInfo upperPosition = nestedPositionBuilder
.setLedgerId(upperKey)
.setEntryId(upperValue)
.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();
return rangeList;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,15 +45,18 @@ public class RangeSetWrapper<T extends Comparable<T>> implements LongPairRangeSe
* Record which Ledger is dirty.
*/
private final DefaultRangeSet<Long> dirtyLedgers = new LongPairRangeSet.DefaultRangeSet<>(
(LongPairConsumer<Long>) (key, value) -> key);
(LongPairConsumer<Long>) (key, value) -> key,
(RangeBoundConsumer<Long>) key -> new LongPair(key, 0));

public RangeSetWrapper(LongPairConsumer<T> rangeConverter, ManagedCursorImpl managedCursor) {
public RangeSetWrapper(LongPairConsumer<T> rangeConverter,
RangeBoundConsumer<T> 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();
}

Expand Down Expand Up @@ -118,6 +121,11 @@ public void forEach(RangeProcessor<T> action, LongPairConsumer<? extends T> cons
rangeSet.forEach(action, consumer);
}

@Override
public void forEachRawRange(RawRangeProcessor action) {
rangeSet.forEachRawRange(action);
}

@Override
public int size() {
return rangeSet.size();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -40,6 +41,8 @@
public class RangeSetWrapperTest {

static final LongPairConsumer<LongPair> consumer = (key, value) -> new LongPair(key, value);
static final RangeBoundConsumer<LongPair> reverseConvert = (pair) -> pair;

ManagedLedgerImpl managedLedger;
RangeSetWrapper<LongPair> set;
ManagedLedgerConfig managedLedgerConfig;
Expand Down Expand Up @@ -67,7 +70,9 @@ public void clean() throws Exception {

@Test
public void testDirtyLedger() {
RangeSetWrapper<LongPair> rangeSetWrapper = new RangeSetWrapper<>(consumer, managedCursor);
RangeSetWrapper<LongPair> rangeSetWrapper = new RangeSetWrapper<>(consumer,
reverseConvert,
managedCursor);
// Test add range
rangeSetWrapper.addOpenClosed(10, 0, 20, 0);
assertEquals(rangeSetWrapper.size(), 1);
Expand All @@ -93,7 +98,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
Expand All @@ -113,7 +118,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);
Expand All @@ -131,7 +136,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);
Expand All @@ -148,7 +153,7 @@ public void testAddForDifferentKey2() {

@Test
public void testAddCompareCompareWithGuava() {
set = new RangeSetWrapper<>(consumer, managedCursor);
set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor);
com.google.common.collect.RangeSet<LongPair> gSet = TreeRangeSet.create();

// add 10K values for key 0
Expand Down Expand Up @@ -187,7 +192,7 @@ public void testAddCompareCompareWithGuava() {

@Test
public void testDeleteCompareWithGuava() throws Exception {
RangeSetWrapper<LongPair> set = new RangeSetWrapper<>(consumer, managedCursor);
RangeSetWrapper<LongPair> set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor);
com.google.common.collect.RangeSet<LongPair> gSet = TreeRangeSet.create();

// add 10K values for key 0
Expand Down Expand Up @@ -241,7 +246,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<LongPair> gSet = TreeRangeSet.create();
set.addOpenClosed(0, 97, 0, 99);
gSet.add(Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)));
Expand All @@ -266,7 +271,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)));
Expand All @@ -281,7 +286,7 @@ public void testFirstRange() {

@Test
public void testLastRange() {
set = new RangeSetWrapper<>(consumer, managedCursor);
set = new RangeSetWrapper<>(consumer, reverseConvert, managedCursor);
assertNull(set.lastRange());
Range<LongPair> range = Range.openClosed(new LongPair(0, 97), new LongPair(0, 99));
set.addOpenClosed(0, 97, 0, 99);
Expand All @@ -302,7 +307,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);
Expand All @@ -313,7 +318,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);
Expand Down Expand Up @@ -344,7 +349,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);
Expand All @@ -370,7 +375,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);
Expand All @@ -390,7 +395,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);
Expand All @@ -411,7 +416,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);
Expand Down Expand Up @@ -439,7 +444,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<LongPair> gSet = TreeRangeSet.create();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -195,7 +195,18 @@ public void forEach(RangeProcessor<T> action) {
}

@Override
public void forEach(RangeProcessor<T> action, LongPairConsumer<? extends T> consumer) {
public void forEach(RangeProcessor<T> action, LongPairConsumer<? extends T> consumerParam) {
forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> {
Range<T> range = Range.openClosed(
consumerParam.apply(lowerKey, lowerValue),
consumerParam.apply(upperKey, upperValue)
);
return action.process(range);
});
}

@Override
public void forEachRawRange(RawRangeProcessor processor) {
AtomicBoolean completed = new AtomicBoolean(false);
rangeBitSetMap.forEach((key, set) -> {
if (completed.get()) {
Expand All @@ -209,9 +220,8 @@ public void forEach(RangeProcessor<T> action, LongPairConsumer<? extends T> cons
int currentClosedMark = first;
while (currentClosedMark != -1 && currentClosedMark <= last) {
int nextOpenMark = set.nextClearBit(currentClosedMark);
Range<T> range = Range.openClosed(consumer.apply(key, currentClosedMark - 1),
consumer.apply(key, nextOpenMark - 1));
if (!action.process(range)) {
if (!processor.processRawRange(key, currentClosedMark - 1,
key, nextOpenMark - 1)) {
completed.set(true);
break;
}
Expand All @@ -220,6 +230,7 @@ public void forEach(RangeProcessor<T> action, LongPairConsumer<? extends T> cons
});
}


@Override
public Range<T> firstRange() {
if (rangeBitSetMap.isEmpty()) {
Expand Down Expand Up @@ -269,10 +280,13 @@ public int cardinality(long lowerKey, long lowerValue, long upperKey, long upper
public int size() {
if (updatedAfterCachedForSize) {
MutableInt size = new MutableInt(0);
forEach((range) -> {

// ignore result because we just want to count
forEachRawRange((lowerKey, lowerValue, upperKey, upperValue) -> {
size.increment();
return true;
});

cachedSize = size.intValue();
updatedAfterCachedForSize = false;
}
Expand Down
Loading