From e6e8338522ec04b43665d2a10d9e755757c700b8 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Tue, 12 Mar 2019 17:00:10 -0700 Subject: [PATCH 1/6] [pulsar-common] add open Concurrent LongPair RangeSet --- .../ConcurrentLongPairRangeSet.java | 339 ++++++++++++++++++ .../util/collections/LongPairRangeSet.java | 146 ++++++++ .../ConcurrentLongPairRangeSetTest.java | 328 +++++++++++++++++ 3 files changed, 813 insertions(+) create mode 100644 pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java create mode 100644 pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java create mode 100644 pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java new file mode 100644 index 0000000000000..9147b9597da4f --- /dev/null +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java @@ -0,0 +1,339 @@ +package org.apache.pulsar.common.util.collections; + +import java.util.ArrayList; +import java.util.BitSet; +import java.util.List; +import java.util.Map.Entry; +import java.util.NavigableMap; +import java.util.concurrent.ConcurrentSkipListMap; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.locks.StampedLock; + +import static com.google.common.base.Preconditions.checkNotNull; +import com.google.common.collect.BoundType; +import com.google.common.collect.ComparisonChain; +import com.google.common.collect.Range; + +/** + * A Concurrent set comprising zero or more ranges of type {@link LongPair}. This can be alternative of + * {@link com.google.common.collect.RangeSet} and can be used if {@code range} type is {@link LongPair}
+ * + *
+ *  
+ * Usage:
+ * a. This can be used if one doesn't want to create object for every new inserted {@code range}
+ * b. It creates {@link BitSet} for every unique first-key of the range. 
+ * So, this rangeSet is not appropriate for large number of unique first-keys.
+ * 
+ * + * + */ +public class ConcurrentLongPairRangeSet implements LongPairRangeSet { + + protected final NavigableMap rangeBitSetMap = new ConcurrentSkipListMap<>(); + private boolean threadSafe = true; + private final int size; + + public ConcurrentLongPairRangeSet() { + this(1024, true); + } + + public ConcurrentLongPairRangeSet(int size) { + this(size, true); + } + + public ConcurrentLongPairRangeSet(int size, boolean threadSafe) { + this.threadSafe = threadSafe; + this.size = size; + } + + public void clear() { + rangeBitSetMap.clear(); + } + + class ConcurrentBitSet extends BitSet { + private static final long serialVersionUID = 1L; + private final StampedLock rwLock = new StampedLock(); + + /** + * Creates a bit set whose initial size is large enough to explicitly represent bits with indices in the range + * {@code 0} through {@code nbits-1}. All bits are initially {@code false}. + * + * @param nbits + * the initial size of the bit set + * @throws NegativeArraySizeException + * if the specified initial size is negative + */ + public ConcurrentBitSet(int nbits) { + super(nbits); + } + + @Override + public boolean get(int bitIndex) { + return super.get(bitIndex); + } + + @Override + public void set(int bitIndex) { + long stamp = rwLock.writeLock(); + try { + super.set(bitIndex); + } finally { + rwLock.unlockWrite(stamp); + } + } + + @Override + public void set(int fromIndex, int toIndex) { + long stamp = rwLock.writeLock(); + try { + super.set(fromIndex, toIndex); + } finally { + rwLock.unlockWrite(stamp); + } + } + + @Override + public int nextSetBit(int fromIndex) { + long stamp = rwLock.readLock(); + try { + return super.nextSetBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + + @Override + public int nextClearBit(int fromIndex) { + long stamp = rwLock.readLock(); + try { + return super.nextClearBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + + @Override + public int previousSetBit(int fromIndex) { + long stamp = rwLock.readLock(); + try { + return super.previousSetBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + + @Override + public int previousClearBit(int fromIndex) { + long stamp = rwLock.readLock(); + try { + return super.previousClearBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + + @Override + public boolean isEmpty() { + long stamp = rwLock.tryOptimisticRead(); + boolean isEmpty = super.isEmpty(); + if (!rwLock.validate(stamp)) { + // Fallback to read lock + stamp = rwLock.readLock(); + try { + isEmpty = super.isEmpty(); + } finally { + rwLock.unlockRead(stamp); + } + } + return isEmpty; + } + + } + + /** + * Adds the specified range to this {@code RangeSet} (optional operation). That is, for equal range sets a and b, + * the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which both + * {@code a.enclosesAll(b)} and {@code a.encloses(range)}. + * + *

+ * Note that {@code range} will merge given {@code range} with any ranges in the range set that are + * {@linkplain Range#isConnected(Range) connected} with it. Moreover, if {@code range} is empty, this is a no-op. + */ + public void add(Range range) { + LongPair lowerEndpoint = range.hasLowerBound() ? range.lowerEndpoint() : LongPair.earliest; + LongPair upperEndpoint = range.hasUpperBound() ? range.upperEndpoint() : LongPair.latest; + + int lower = (range.hasLowerBound() && range.lowerBoundType().equals(BoundType.CLOSED)) + ? getSafeEntry(lowerEndpoint) + : getSafeEntry(lowerEndpoint) + 1; + int upper = (range.hasUpperBound() && range.lowerBoundType().equals(BoundType.CLOSED)) + ? getSafeEntry(upperEndpoint) + : getSafeEntry(upperEndpoint) + 1; + if (lowerEndpoint.getKey() != upperEndpoint.getKey()) { + // (1) set lower to last in lowerRange.getKey() + if (!lowerEndpoint.equals(LongPair.earliest)) { + BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(lowerEndpoint.getKey(), (key) -> createNewBitSet()); + if (rangeBitSet != null) { + int lastValue = rangeBitSet.previousSetBit(rangeBitSet.size()); + rangeBitSet.set(lower, Math.max(lastValue, lower) + 1); + } + } + // (2) set 0 to upper in upperRange.getKey() + if (!upperEndpoint.equals(LongPair.latest)) { + BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(upperEndpoint.getKey(), (key) -> createNewBitSet()); + if (rangeBitSet != null) { + rangeBitSet.set(0, upper + 1); + } + } + } else { + long key = lowerEndpoint.getKey(); + BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(key, (k) -> createNewBitSet()); + rangeBitSet.set(lower, upper + 1); + } + } + + public boolean contains(LongPair position) { + checkNotNull(position, "argument can't be null"); + BitSet rangeBitSet = rangeBitSetMap.get(position.getKey()); + if (rangeBitSet != null) { + return rangeBitSet.get(getSafeEntry(position)); + } + return false; + } + + public Range rangeContaining(LongPair position) { + checkNotNull(position, "argument can't be null"); + BitSet rangeBitSet = rangeBitSetMap.get(position.getKey()); + if (rangeBitSet != null) { + if (!rangeBitSet.get(getSafeEntry(position))) { + // TODO: document it + // if position is not part of any range then return null + return null; + } + final LongPair lower = new LongPair(position.getKey(), + rangeBitSet.previousClearBit(getSafeEntry(position)) + 1); + final LongPair upper = new LongPair(position.getKey(), + Math.max(rangeBitSet.nextClearBit(getSafeEntry(position)) - 1, lower.getValue())); + return Range.closed(lower, upper); + } else { + // position's key doesn't exist so, range should be last entry Of previous key and first entry of next + // key + Entry previousRangeBitSet = rangeBitSetMap.ceilingEntry(position.getKey()); + final LongPair lower = (previousRangeBitSet != null) + ? new LongPair(previousRangeBitSet.getKey(), + previousRangeBitSet.getValue().previousSetBit(previousRangeBitSet.getValue().size())) + : LongPair.earliest; + Entry nextRangeBitSet = rangeBitSetMap.floorEntry(position.getKey()); + final LongPair upper = (nextRangeBitSet != null) + ? new LongPair(nextRangeBitSet.getKey(), nextRangeBitSet.getValue().nextSetBit(0)) + : LongPair.latest; + return Range.closed(lower, upper); + } + } + + public void remove(Range range) { + LongPair lowerEndpoint = range.hasLowerBound() ? range.lowerEndpoint() : LongPair.earliest; + LongPair upperEndpoint = range.hasUpperBound() ? range.upperEndpoint() : LongPair.latest; + + int lower = (range.hasLowerBound() && range.lowerBoundType().equals(BoundType.CLOSED)) + ? getSafeEntry(lowerEndpoint) + : getSafeEntry(lowerEndpoint) + 1; + int upper = (range.hasUpperBound() && range.upperBoundType().equals(BoundType.CLOSED)) + ? getSafeEntry(upperEndpoint) + : getSafeEntry(upperEndpoint) - 1; + + // if lower-bound is not set then remove all the keys less than given upper-bound range + if (lowerEndpoint.equals(LongPair.earliest)) { + // remove all keys with + rangeBitSetMap.forEach((key, set) -> { + if (key < upperEndpoint.getKey()) { + rangeBitSetMap.remove(key); + } + }); + } + + // if upper-bound is not set then remove all the keys greater than given lower-bound range + if (upperEndpoint.equals(LongPair.latest)) { + // remove all keys with + rangeBitSetMap.forEach((key, set) -> { + if (key > lowerEndpoint.getKey()) { + rangeBitSetMap.remove(key); + } + }); + } + + // remove all the keys between two endpoint keys + rangeBitSetMap.forEach((key, set) -> { + if (lowerEndpoint.getKey() == upperEndpoint.getKey()) { + set.clear(lower, upper + 1); + } else { + // eg: remove-range: [(3,5) - (5,5)] -> Delete all items from 3,6->3,N,4.*,5,0->5,5 + if (key == lowerEndpoint.getKey()) { + // remove all entries from given position to last position + set.clear(lower, set.previousSetBit(set.size())); + } else if (key == upperEndpoint.getKey()) { + // remove all entries from 0 to given position + set.clear(0, upper + 1); + } else if (key > lowerEndpoint.getKey() && key < upperEndpoint.getKey()) { + rangeBitSetMap.remove(key); + } + } + // remove bit-set if set is empty + if (set.isEmpty()) { + rangeBitSetMap.remove(key); + } + }); + } + + public boolean isEmpty() { + if (rangeBitSetMap.isEmpty()) { + return true; + } + AtomicBoolean isEmpty = new AtomicBoolean(false); + rangeBitSetMap.forEach((key, val) -> { + if (!isEmpty.get()) { + return; + } + isEmpty.set(val.isEmpty()); + }); + return isEmpty.get(); + } + + public Range span() { + Entry firstSet = rangeBitSetMap.firstEntry(); + Entry lastSet = rangeBitSetMap.lastEntry(); + int first = firstSet.getValue().nextSetBit(0); + int last = lastSet.getValue().previousSetBit(lastSet.getValue().size()); + return Range.closed(new LongPair(firstSet.getKey(), first), new LongPair(lastSet.getKey(), last)); + } + + public List> asRanges() { + List> ranges = new ArrayList<>(); + rangeBitSetMap.forEach((key, set) -> { + 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.closed(new LongPair(key, currentClosedMark), + new LongPair(key, nextOpenMark - 1)); + ranges.add(range); + currentClosedMark = set.nextSetBit(nextOpenMark); + } + }); + return ranges; + } + + private int getSafeEntry(LongPair position) { + return (int) (position.getValue() > 0 ? position.getValue() : 0); + } + + private BitSet createNewBitSet() { + return this.threadSafe ? new ConcurrentBitSet(size) : new BitSet(size); + } + +} 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 new file mode 100644 index 0000000000000..986d5dbbf0071 --- /dev/null +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongPairRangeSet.java @@ -0,0 +1,146 @@ +package org.apache.pulsar.common.util.collections; + +import java.util.Collection; +import java.util.Set; + +import com.google.common.collect.ComparisonChain; +import com.google.common.collect.Range; +import com.google.common.collect.RangeSet; +import com.google.common.collect.TreeRangeSet; + +/** + * A set comprising zero or more ranges type of {@link LongPair} + */ +public interface LongPairRangeSet { + + /** + * Default RangeSet implementation based on {@link TreeRangeSet} + * + * @return + */ + static LongPairRangeSet create() { + return new DefaultRangeSet(); + } + + void clear(); + + /** + * Adds the specified range to this {@code RangeSet} (optional operation). That is, for equal range sets a and b, + * the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which both + * {@code a.enclosesAll(b)} and {@code a.encloses(range)}. + * + *

+ * Note that {@code range} will merge given {@code range} with any ranges in the range set that are + * {@linkplain Range#isConnected(Range) connected} with it. Moreover, if {@code range} is empty, this is a no-op. + */ + void add(Range range); + + /** Determines whether any of this range set's member ranges contains {@code value}. */ + boolean contains(LongPair position); + + /** + * Returns the unique range from this range set that {@linkplain Range#contains contains} {@code value}, or + * {@code null} if this range set does not contain {@code value}. + */ + Range rangeContaining(LongPair position); + + /** + * Removes the specified range from this {@code RangeSet} + * + * @param range + */ + void remove(Range range); + + boolean isEmpty(); + + /** + * Returns the minimal range which {@linkplain Range#encloses(Range) encloses} all ranges in this range set. + * + * @return + */ + Range span(); + + /** + * Returns a view of the {@linkplain Range#isConnected disconnected} ranges that make up this range set. + * + * @return + */ + Collection> asRanges(); + + public static class LongPair implements Comparable { + + public static final LongPair earliest = new LongPair(-1, -1); + public static final LongPair latest = new LongPair(Integer.MAX_VALUE, Integer.MAX_VALUE); + + private long key; + private int value; + + public LongPair(long key, int value) { + this.key = key; + this.value = value; + } + + public long getKey() { + return this.key; + } + + public int getValue() { + return this.value; + } + + @Override + public int compareTo(LongPair o) { + return ComparisonChain.start().compare(key, o.getKey()).compare(value, o.getValue()).result(); + } + + @Override + public String toString() { + return String.format("%d:%d", key, value); + } + } + + public static class DefaultRangeSet implements LongPairRangeSet { + + RangeSet set = TreeRangeSet.create(); + + @Override + public void clear() { + set.clear(); + } + + @Override + public void add(Range range) { + set.add(range); + } + + @Override + public boolean contains(LongPair position) { + return set.contains(position); + } + + @Override + public Range rangeContaining(LongPair position) { + return set.rangeContaining(position); + } + + @Override + public void remove(Range range) { + set.remove(range); + } + + @Override + public boolean isEmpty() { + return set.isEmpty(); + } + + @Override + public Range span() { + return set.span(); + } + + @Override + public Set> asRanges() { + return set.asRanges(); + } + } +} diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java new file mode 100644 index 0000000000000..1caee023a2033 --- /dev/null +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java @@ -0,0 +1,328 @@ +/** + * 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.common.util.collections; + +import static org.testng.Assert.assertEquals; + +import java.util.List; +import java.util.Set; + +import org.apache.pulsar.common.util.collections.ConcurrentLongPairRangeSet.LongPair; +import org.testng.annotations.Test; + +import com.google.common.collect.BoundType; +import com.google.common.collect.Lists; +import com.google.common.collect.Range; +import com.google.common.collect.TreeRangeSet; + +public class ConcurrentLongPairRangeSetTest { + + @Test + public void testAddForSameKey() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + // add 0 to 5 + set.add(Range.closed(new LongPair(0, 0), new LongPair(0, 5))); + // add 8,9,10 + set.add(Range.closed(new LongPair(0, 8), new LongPair(0, 8))); + set.add(Range.closed(new LongPair(0, 9), new LongPair(0, 9))); + set.add(Range.closed(new LongPair(0, 10), new LongPair(0, 10))); + // add 98 to 99 and 102,105 + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 102), new LongPair(0, 106))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 0), new LongPair(0, 5)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 8), new LongPair(0, 10)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 102), new LongPair(0, 106)))); + } + + @Test + public void testAddForDifferentKey() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + // [98,100],[(1,5),(1,5)],[(1,10,1,15)],[(1,20),(1,20)],[(2,0),(2,10)] + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + set.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); + set.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 100)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 0), new LongPair(1, 5)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 10), new LongPair(1, 15)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 20), new LongPair(1, 20)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 0), new LongPair(2, 10)))); + } + + @Test + public void testAddCompareCompareWithGuava() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); + + // add 10K values for key 0 + int totalInsert = 10_000; + // add single values + for (int i = 0; i < totalInsert; i++) { + if (i % 3 == 0 || i % 6 == 0 || i % 8 == 0) { + set.add(Range.closed(new LongPair(0, i), new LongPair(0, i))); + gSet.add(Range.closed(new LongPair(0, i), new LongPair(0, i))); + } + } + // add batches + for (int i = totalInsert; i < (totalInsert * 2); i++) { + if (i % 5 == 0) { + set.add(Range.closed(new LongPair(0, i - 3), new LongPair(0, i + 3))); + gSet.add(Range.closed(new LongPair(0, i - 3), new LongPair(0, i + 3))); + } + } + List> ranges = set.asRanges(); + Set> gRanges = gSet.asRanges(); + + List> gRangeConnected = getConnectedRange(gRanges); + assertEquals(gRangeConnected.size(), ranges.size()); + int i = 0; + for (Range range : gRangeConnected) { + assertEquals(range, ranges.get(i)); + i++; + } + } + + @Test + public void testDeleteCompareWithGuava() { + + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); + + // add 10K values for key 0 + int totalInsert = 10_000; + // add single values + List> removedRanges = Lists.newArrayList(); + for (int i = 0; i < totalInsert; i++) { + if (i % 3 == 0 || i % 7 == 0 || i % 11 == 0) { + continue; + } + Range range = Range.closed(new LongPair(0, i), new LongPair(0, i)); + set.add(range); + gSet.add(range); + if (i % 4 == 0) { + removedRanges.add(range); + } + } + // add batches + for (int i = totalInsert; i < (totalInsert * 2); i++) { + Range range = Range.closed(new LongPair(0, i - 3), new LongPair(0, i + 3)); + if (i % 5 != 0) { + set.add(range); + gSet.add(range); + } + if (i % 4 == 0) { + removedRanges.add(range); + } + } + // remove records + for (Range range : removedRanges) { + set.remove(range); + gSet.remove(range); + } + + List> ranges = set.asRanges(); + Set> gRanges = gSet.asRanges(); + List> gRangeConnected = getConnectedRange(gRanges); + assertEquals(gRangeConnected.size(), ranges.size()); + int i = 0; + for (Range range : gRangeConnected) { + assertEquals(range, ranges.get(i)); + i++; + } + } + + @Test + public void testSpanWithGuava() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + 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(); + gSet.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + gSet.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + assertEquals(set.span(), gSet.span()); + assertEquals(set.span(), Range.closed(new LongPair(0, 98), new LongPair(1, 5))); + + 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))); + gSet.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); + gSet.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); + gSet.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); + gSet.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); + gSet.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); + assertEquals(set.span(), gSet.span()); + assertEquals(set.span(), Range.closed(new LongPair(0, 98), new LongPair(4, 20))); + + } + + @Test + public void testDeleteForDifferentKey() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + 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))); + + // delete only (0,100) + set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); + + /** + * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] + */ + set.remove(Range.closed(new LongPair(2, 27), new LongPair(4, 15))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 0), new LongPair(1, 5)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 10), new LongPair(1, 15)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 20), new LongPair(1, 20)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 0), new LongPair(2, 10)))); + + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 25), new LongPair(2, 26)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 28), new LongPair(2, 28)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(4, 16), new LongPair(4, 20)))); + } + + @Test + public void testDeleteWithAtMost() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + 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))); + + // delete only (0,100) + set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); + + /** + * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] + */ + set.remove(Range.atMost(new LongPair(2, 27))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 28), new LongPair(2, 28)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(3, 12), new LongPair(3, 20)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(4, 12), new LongPair(4, 20)))); + } + + @Test + public void testDeleteWithLeastMost() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + 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))); + + // delete only (0,100) + set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); + + /** + * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] + */ + set.remove(Range.atLeast(new LongPair(2, 27))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 0), new LongPair(1, 5)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 10), new LongPair(1, 15)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 20), new LongPair(1, 20)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 0), new LongPair(2, 10)))); + assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 25), new LongPair(2, 26)))); + } + + @Test + public void testRangeContaining() { + ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); + 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(); + gSet.add(Range.closed(new LongPair(0, 98), new LongPair(0, 100))); + gSet.add(Range.closed(new LongPair(0, 101), new LongPair(1, 5))); + 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))); + gSet.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); + gSet.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); + gSet.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); + gSet.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); + gSet.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); + + LongPair position = new LongPair(0, 99); + assertEquals(set.rangeContaining(position), Range.closed(new LongPair(0, 98), new LongPair(0, 100))); + assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); + + position = new LongPair(2, 30); + assertEquals(set.rangeContaining(position), null); + assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); + + position = new LongPair(3, 13); + assertEquals(set.rangeContaining(position), Range.closed(new LongPair(3, 12), new LongPair(3, 20))); + assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); + + position = new LongPair(3, 22); + assertEquals(set.rangeContaining(position), null); + assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); + } + + private List> getConnectedRange(Set> gRanges) { + List> gRangeConnected = Lists.newArrayList(); + Range lastRange = null; + for (Range range : gRanges) { + if (lastRange == null) { + lastRange = range; + continue; + } + if ((lastRange.upperEndpoint().getValue() + 1) == (range.lowerEndpoint().getValue())) { + lastRange = Range.closed(lastRange.lowerEndpoint(), range.upperEndpoint()); + } else { + gRangeConnected.add(lastRange); + lastRange = range; + } + } + lastRange = lastRange.lowerBoundType().equals(BoundType.CLOSED) ? lastRange + : Range.closed( + new LongPair(lastRange.lowerEndpoint().getKey(), lastRange.lowerEndpoint().getValue() + 1), + lastRange.upperEndpoint()); + gRangeConnected.add(lastRange); + return gRangeConnected; + } +} From ca4c38beb7d6c4cde1e49a478abd03c56cfd2402 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Wed, 13 Mar 2019 14:09:09 -0700 Subject: [PATCH 2/6] add open-range set methods --- ...va => ConcurrentOpenLongPairRangeSet.java} | 345 +++++++++++----- .../util/collections/LongPairRangeSet.java | 148 +++++-- .../ConcurrentLongPairRangeSetTest.java | 328 --------------- .../ConcurrentOpenLongPairRangeSetTest.java | 384 ++++++++++++++++++ 4 files changed, 726 insertions(+), 479 deletions(-) rename pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/{ConcurrentLongPairRangeSet.java => ConcurrentOpenLongPairRangeSet.java} (51%) delete mode 100644 pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java create mode 100644 pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java similarity index 51% rename from pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java rename to pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java index 9147b9597da4f..51ea3cec55bbc 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSet.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSet.java @@ -1,5 +1,25 @@ +/** + * 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.common.util.collections; +import static com.google.common.base.Preconditions.checkNotNull; + import java.util.ArrayList; import java.util.BitSet; import java.util.List; @@ -7,11 +27,10 @@ import java.util.NavigableMap; import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.StampedLock; -import static com.google.common.base.Preconditions.checkNotNull; import com.google.common.collect.BoundType; -import com.google.common.collect.ComparisonChain; import com.google.common.collect.Range; /** @@ -23,32 +42,35 @@ * Usage: * a. This can be used if one doesn't want to create object for every new inserted {@code range} * b. It creates {@link BitSet} for every unique first-key of the range. - * So, this rangeSet is not appropriate for large number of unique first-keys. + * So, this rangeSet is not suitable for large number of unique keys. * * * */ -public class ConcurrentLongPairRangeSet implements LongPairRangeSet { +public class ConcurrentOpenLongPairRangeSet> implements LongPairRangeSet { protected final NavigableMap rangeBitSetMap = new ConcurrentSkipListMap<>(); private boolean threadSafe = true; - private final int size; + private final int bitSetSize; + private final LongPairConsumer consumer; - public ConcurrentLongPairRangeSet() { - this(1024, true); - } + // caching place-holder for cpu-optimization to avoid calculating ranges again + private volatile int cachedSize = 0; + private volatile String cachedToString = "[]"; + private volatile boolean updatedAfterCached = true; - public ConcurrentLongPairRangeSet(int size) { - this(size, true); + public ConcurrentOpenLongPairRangeSet(LongPairConsumer consumer) { + this(1024, true, consumer); } - public ConcurrentLongPairRangeSet(int size, boolean threadSafe) { - this.threadSafe = threadSafe; - this.size = size; + public ConcurrentOpenLongPairRangeSet(int size, LongPairConsumer consumer) { + this(size, true, consumer); } - public void clear() { - rangeBitSetMap.clear(); + public ConcurrentOpenLongPairRangeSet(int size, boolean threadSafe, LongPairConsumer consumer) { + this.threadSafe = threadSafe; + this.bitSetSize = size; + this.consumer = consumer; } class ConcurrentBitSet extends BitSet { @@ -155,91 +177,230 @@ public boolean isEmpty() { * Adds the specified range to this {@code RangeSet} (optional operation). That is, for equal range sets a and b, * the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which both * {@code a.enclosesAll(b)} and {@code a.encloses(range)}. - * *

* Note that {@code range} will merge given {@code range} with any ranges in the range set that are * {@linkplain Range#isConnected(Range) connected} with it. Moreover, if {@code range} is empty, this is a no-op. */ - public void add(Range range) { - LongPair lowerEndpoint = range.hasLowerBound() ? range.lowerEndpoint() : LongPair.earliest; - LongPair upperEndpoint = range.hasUpperBound() ? range.upperEndpoint() : LongPair.latest; - - int lower = (range.hasLowerBound() && range.lowerBoundType().equals(BoundType.CLOSED)) - ? getSafeEntry(lowerEndpoint) - : getSafeEntry(lowerEndpoint) + 1; - int upper = (range.hasUpperBound() && range.lowerBoundType().equals(BoundType.CLOSED)) - ? getSafeEntry(upperEndpoint) - : getSafeEntry(upperEndpoint) + 1; - if (lowerEndpoint.getKey() != upperEndpoint.getKey()) { + @Override + public void addOpenClosed(long lowerKey, long lowerValueOpen, long upperKey, long upperValue) { + long lowerValue = lowerValueOpen + 1; + if (lowerKey != upperKey) { // (1) set lower to last in lowerRange.getKey() - if (!lowerEndpoint.equals(LongPair.earliest)) { - BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(lowerEndpoint.getKey(), (key) -> createNewBitSet()); - if (rangeBitSet != null) { + if (isValid(lowerKey, lowerValue)) { + BitSet rangeBitSet = rangeBitSetMap.get(lowerKey); + // if lower and upper has different key/ledger then set ranges for lower-key only if + // a. bitSet already exist and given value is not the last value in the bitset. + // it will prevent setting up values which are not actually expected to set + // eg: (2:10..4:10] in this case , don't set any value for 2:10 and set [4:0..4:10] + if (rangeBitSet != null && (rangeBitSet.previousSetBit(rangeBitSet.size()) > lowerValueOpen)) { int lastValue = rangeBitSet.previousSetBit(rangeBitSet.size()); - rangeBitSet.set(lower, Math.max(lastValue, lower) + 1); + rangeBitSet.set((int) lowerValue, (int) Math.max(lastValue, lowerValue) + 1); } } - // (2) set 0 to upper in upperRange.getKey() - if (!upperEndpoint.equals(LongPair.latest)) { - BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(upperEndpoint.getKey(), (key) -> createNewBitSet()); + // (2) set 0th-index to upper-index in upperRange.getKey() + if (isValid(upperKey, upperValue)) { + BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(upperKey, (key) -> createNewBitSet()); if (rangeBitSet != null) { - rangeBitSet.set(0, upper + 1); + rangeBitSet.set(0, (int) upperValue + 1); } } + // No-op if values are not valid eg: if lower == LongPair.earliest or upper == LongPair.latest then nothing + // to set } else { - long key = lowerEndpoint.getKey(); + long key = lowerKey; BitSet rangeBitSet = rangeBitSetMap.computeIfAbsent(key, (k) -> createNewBitSet()); - rangeBitSet.set(lower, upper + 1); + rangeBitSet.set((int) lowerValue, (int) upperValue + 1); } + updatedAfterCached = true; } - public boolean contains(LongPair position) { - checkNotNull(position, "argument can't be null"); - BitSet rangeBitSet = rangeBitSetMap.get(position.getKey()); + private boolean isValid(long key, long value) { + return key != LongPair.earliest.getKey() && value != LongPair.earliest.getValue() + && key != LongPair.latest.getKey() && value != LongPair.latest.getValue(); + } + + @Override + public boolean contains(long key, long value) { + + BitSet rangeBitSet = rangeBitSetMap.get(key); if (rangeBitSet != null) { - return rangeBitSet.get(getSafeEntry(position)); + return rangeBitSet.get(getSafeEntry(value)); } return false; } - public Range rangeContaining(LongPair position) { - checkNotNull(position, "argument can't be null"); - BitSet rangeBitSet = rangeBitSetMap.get(position.getKey()); + @Override + public Range rangeContaining(long key, long value) { + BitSet rangeBitSet = rangeBitSetMap.get(key); if (rangeBitSet != null) { - if (!rangeBitSet.get(getSafeEntry(position))) { - // TODO: document it + if (!rangeBitSet.get(getSafeEntry(value))) { // if position is not part of any range then return null return null; } - final LongPair lower = new LongPair(position.getKey(), - rangeBitSet.previousClearBit(getSafeEntry(position)) + 1); - final LongPair upper = new LongPair(position.getKey(), - Math.max(rangeBitSet.nextClearBit(getSafeEntry(position)) - 1, lower.getValue())); - return Range.closed(lower, upper); - } else { - // position's key doesn't exist so, range should be last entry Of previous key and first entry of next - // key - Entry previousRangeBitSet = rangeBitSetMap.ceilingEntry(position.getKey()); - final LongPair lower = (previousRangeBitSet != null) - ? new LongPair(previousRangeBitSet.getKey(), - previousRangeBitSet.getValue().previousSetBit(previousRangeBitSet.getValue().size())) - : LongPair.earliest; - Entry nextRangeBitSet = rangeBitSetMap.floorEntry(position.getKey()); - final LongPair upper = (nextRangeBitSet != null) - ? new LongPair(nextRangeBitSet.getKey(), nextRangeBitSet.getValue().nextSetBit(0)) - : LongPair.latest; + int lowerValue = rangeBitSet.previousClearBit(getSafeEntry(value)) + 1; + final T lower = consumer.apply(key, lowerValue); + final T upper = consumer.apply(key, + Math.max(rangeBitSet.nextClearBit(getSafeEntry(value)) - 1, lowerValue)); return Range.closed(lower, upper); } + return null; + } + + @Override + public void removeAtMost(long key, long value) { + this.remove(Range.atMost(new LongPair(key, value))); + } + + @Override + public boolean isEmpty() { + if (rangeBitSetMap.isEmpty()) { + return true; + } + AtomicBoolean isEmpty = new AtomicBoolean(false); + rangeBitSetMap.forEach((key, val) -> { + if (!isEmpty.get()) { + return; + } + isEmpty.set(val.isEmpty()); + }); + return isEmpty.get(); + } + + @Override + public void clear() { + rangeBitSetMap.clear(); + updatedAfterCached = true; + } + + @Override + public Range span() { + Entry firstSet = rangeBitSetMap.firstEntry(); + Entry lastSet = rangeBitSetMap.lastEntry(); + int first = firstSet.getValue().nextSetBit(0); + int last = lastSet.getValue().previousSetBit(lastSet.getValue().size()); + return Range.openClosed(consumer.apply(firstSet.getKey(), first - 1), consumer.apply(lastSet.getKey(), last)); + } + + @Override + public List> asRanges() { + List> ranges = new ArrayList<>(); + createRanges(ranges, null); + return ranges; + } + + @Override + public Range firstRange() { + Entry firstSet = rangeBitSetMap.firstEntry(); + int lower = firstSet.getValue().nextSetBit(0); + int upper = Math.max(lower, firstSet.getValue().nextClearBit(lower) - 1); + return Range.openClosed(consumer.apply(firstSet.getKey(), lower - 1), consumer.apply(firstSet.getKey(), upper)); + } + + @Override + public int size() { + if (updatedAfterCached) { + cachedSize = createRanges(null, null); + updatedAfterCached = false; + } + return cachedSize; + } + + @Override + public String toString() { + if (updatedAfterCached) { + StringBuilder toString = new StringBuilder(); + createRanges(null, toString); + cachedToString = toString.toString(); + updatedAfterCached = false; + } + return cachedToString; + } + + /** + * It creates ranges and add into given list {@code ranges} and also returns total number of ranges into the set. If + * given List {@code ranges} is null then it just returns number of ranges. + * + * @param ranges + * @return + */ + private int createRanges(List> ranges, StringBuilder toString) { + AtomicInteger size = new AtomicInteger(0); + if (toString != null) { + toString.append("["); + } + rangeBitSetMap.forEach((key, set) -> { + if (set.isEmpty()) { + return; + } + int first = set.nextSetBit(0); + int last = set.previousSetBit(set.size()); + int currentClosedMark = first; + // TODO: remove: sometime previous-keyset (previous ledger) is connected to the current one, in that case + // merge the range + 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 (ranges != null) { + ranges.add(range); + } + if (toString != null) { + if (size.get() != 0) { + toString.append(","); + } + toString.append(range); + } + size.getAndIncrement(); + currentClosedMark = set.nextSetBit(nextOpenMark); + } + }); + if (toString != null) { + toString.append("]"); + } + return size.get(); + } + + /** + * Adds the specified range to this {@code RangeSet} (optional operation). That is, for equal range sets a and b, + * the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which both + * {@code a.enclosesAll(b)} and {@code a.encloses(range)}. + * + *

+ * Note that {@code range} will merge given {@code range} with any ranges in the range set that are + * {@linkplain Range#isConnected(Range) connected} with it. Moreover, if {@code range} is empty/invalid, this is a + * no-op. + */ + public void add(Range range) { + LongPair lowerEndpoint = range.hasLowerBound() ? range.lowerEndpoint() : LongPair.earliest; + LongPair upperEndpoint = range.hasUpperBound() ? range.upperEndpoint() : LongPair.latest; + + long lowerValueOpen = (range.hasLowerBound() && range.lowerBoundType().equals(BoundType.CLOSED)) + ? getSafeEntry(lowerEndpoint) - 1 + : getSafeEntry(lowerEndpoint); + long upperValueClosed = (range.hasUpperBound() && range.upperBoundType().equals(BoundType.CLOSED)) + ? getSafeEntry(upperEndpoint) + : getSafeEntry(upperEndpoint) + 1; + + // #addOpenClosed doesn't create bitSet for lower-key because it avoids setting up values for non-exist items + // into the key-ledger. so, create bitSet and initialize so, it can't be ignored at #addOpenClosed + rangeBitSetMap.computeIfAbsent(lowerEndpoint.getKey(), (key) -> createNewBitSet()) + .set((int) lowerValueOpen + 1); + this.addOpenClosed(lowerEndpoint.getKey(), lowerValueOpen, upperEndpoint.getKey(), upperValueClosed); + } + + public boolean contains(LongPair position) { + checkNotNull(position, "argument can't be null"); + return contains(position.getKey(), position.getValue()); } public void remove(Range range) { LongPair lowerEndpoint = range.hasLowerBound() ? range.lowerEndpoint() : LongPair.earliest; LongPair upperEndpoint = range.hasUpperBound() ? range.upperEndpoint() : LongPair.latest; - int lower = (range.hasLowerBound() && range.lowerBoundType().equals(BoundType.CLOSED)) + long lower = (range.hasLowerBound() && range.lowerBoundType().equals(BoundType.CLOSED)) ? getSafeEntry(lowerEndpoint) : getSafeEntry(lowerEndpoint) + 1; - int upper = (range.hasUpperBound() && range.upperBoundType().equals(BoundType.CLOSED)) + long upper = (range.hasUpperBound() && range.upperBoundType().equals(BoundType.CLOSED)) ? getSafeEntry(upperEndpoint) : getSafeEntry(upperEndpoint) - 1; @@ -266,15 +427,15 @@ public void remove(Range range) { // remove all the keys between two endpoint keys rangeBitSetMap.forEach((key, set) -> { if (lowerEndpoint.getKey() == upperEndpoint.getKey()) { - set.clear(lower, upper + 1); + set.clear((int) lower, (int) upper + 1); } else { // eg: remove-range: [(3,5) - (5,5)] -> Delete all items from 3,6->3,N,4.*,5,0->5,5 if (key == lowerEndpoint.getKey()) { // remove all entries from given position to last position - set.clear(lower, set.previousSetBit(set.size())); + set.clear((int) lower, set.previousSetBit(set.size())); } else if (key == upperEndpoint.getKey()) { // remove all entries from 0 to given position - set.clear(0, upper + 1); + set.clear(0, (int) upper + 1); } else if (key > lowerEndpoint.getKey() && key < upperEndpoint.getKey()) { rangeBitSetMap.remove(key); } @@ -284,56 +445,20 @@ public void remove(Range range) { rangeBitSetMap.remove(key); } }); - } - public boolean isEmpty() { - if (rangeBitSetMap.isEmpty()) { - return true; - } - AtomicBoolean isEmpty = new AtomicBoolean(false); - rangeBitSetMap.forEach((key, val) -> { - if (!isEmpty.get()) { - return; - } - isEmpty.set(val.isEmpty()); - }); - return isEmpty.get(); + updatedAfterCached = true; } - public Range span() { - Entry firstSet = rangeBitSetMap.firstEntry(); - Entry lastSet = rangeBitSetMap.lastEntry(); - int first = firstSet.getValue().nextSetBit(0); - int last = lastSet.getValue().previousSetBit(lastSet.getValue().size()); - return Range.closed(new LongPair(firstSet.getKey(), first), new LongPair(lastSet.getKey(), last)); - } - - public List> asRanges() { - List> ranges = new ArrayList<>(); - rangeBitSetMap.forEach((key, set) -> { - 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.closed(new LongPair(key, currentClosedMark), - new LongPair(key, nextOpenMark - 1)); - ranges.add(range); - currentClosedMark = set.nextSetBit(nextOpenMark); - } - }); - return ranges; + private int getSafeEntry(LongPair position) { + return (int) Math.max(position.getValue(), -1); } - private int getSafeEntry(LongPair position) { - return (int) (position.getValue() > 0 ? position.getValue() : 0); + private int getSafeEntry(long value) { + return (int) Math.max(value, -1); } private BitSet createNewBitSet() { - return this.threadSafe ? new ConcurrentBitSet(size) : new BitSet(size); + return this.threadSafe ? new ConcurrentBitSet(bitSetSize) : new BitSet(bitSetSize); } } 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 986d5dbbf0071..02736a503d9ef 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 @@ -1,3 +1,21 @@ +/** + * 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.common.util.collections; import java.util.Collection; @@ -9,63 +27,78 @@ import com.google.common.collect.TreeRangeSet; /** - * A set comprising zero or more ranges type of {@link LongPair} + * A set comprising zero or more ranges type of key-value pair. */ -public interface LongPairRangeSet { +public interface LongPairRangeSet> { /** - * Default RangeSet implementation based on {@link TreeRangeSet} + * Adds the specified range (range that contains all values strictly greater than {@code + * lower} and less than or equal to {@code upper}.) to this {@code RangeSet} (optional operation). That is, for equal + * range sets a and b, the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which + * both {@code a.enclosesAll(b)} and {@code a.encloses(range)}. * - * @return - */ - static LongPairRangeSet create() { - return new DefaultRangeSet(); - } - - void clear(); - - /** - * Adds the specified range to this {@code RangeSet} (optional operation). That is, for equal range sets a and b, - * the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which both - * {@code a.enclosesAll(b)} and {@code a.encloses(range)}. - * - *

- * Note that {@code range} will merge given {@code range} with any ranges in the range set that are - * {@linkplain Range#isConnected(Range) connected} with it. Moreover, if {@code range} is empty, this is a no-op. + *

+     *  
+     * @param lowerKey :  value for key of lowerEndpoint of Range
+     * @param lowerValue: value for value of lowerEndpoint of Range
+     * @param upperKey  : value for key of upperEndpoint of Range
+     * @param upperValue: value for value of upperEndpoint of Range
+     * 
*/ - void add(Range range); + void addOpenClosed(long lowerKey, long lowerValue, long upperKey, long upperValue); /** Determines whether any of this range set's member ranges contains {@code value}. */ - boolean contains(LongPair position); + boolean contains(long key, long value); /** * Returns the unique range from this range set that {@linkplain Range#contains contains} {@code value}, or * {@code null} if this range set does not contain {@code value}. */ - Range rangeContaining(LongPair position); + Range rangeContaining(long key, long value); /** - * Removes the specified range from this {@code RangeSet} + * Remove range that contains all values less than or equal to given key-value. * - * @param range + * @param key + * @param value */ - void remove(Range range); + void removeAtMost(long key, long value); boolean isEmpty(); + void clear(); + /** * Returns the minimal range which {@linkplain Range#encloses(Range) encloses} all ranges in this range set. * * @return */ - Range span(); + Range span(); /** * Returns a view of the {@linkplain Range#isConnected disconnected} ranges that make up this range set. * * @return */ - Collection> asRanges(); + Collection> asRanges(); + + /** + * Returns total number of ranges into the set. + * + * @return + */ + int size(); + + /** + * It returns very first smallest range in the rangeSet. + * + * @return Range first smallest range into the set + */ + Range firstRange(); + + public static interface LongPairConsumer { + T apply(long key, long value); + } public static class LongPair implements Comparable { @@ -73,9 +106,9 @@ public static class LongPair implements Comparable { public static final LongPair latest = new LongPair(Integer.MAX_VALUE, Integer.MAX_VALUE); private long key; - private int value; + private long value; - public LongPair(long key, int value) { + public LongPair(long key, long value) { this.key = key; this.value = value; } @@ -84,7 +117,7 @@ public long getKey() { return this.key; } - public int getValue() { + public long getValue() { return this.value; } @@ -99,9 +132,15 @@ public String toString() { } } - public static class DefaultRangeSet implements LongPairRangeSet { + public static class DefaultRangeSet> implements LongPairRangeSet { + + RangeSet set = TreeRangeSet.create(); - RangeSet set = TreeRangeSet.create(); + private final LongPairConsumer consumer; + + public DefaultRangeSet(LongPairConsumer consumer) { + this.consumer = consumer; + } @Override public void clear() { @@ -109,38 +148,65 @@ public void clear() { } @Override - public void add(Range range) { - set.add(range); + public void addOpenClosed(long key1, long value1, long key2, long value2) { + set.add(Range.openClosed(consumer.apply(key1, value1), consumer.apply(key2, value2))); } - @Override - public boolean contains(LongPair position) { + public boolean contains(T position) { return set.contains(position); } - @Override - public Range rangeContaining(LongPair position) { + public Range rangeContaining(T position) { return set.rangeContaining(position); } @Override - public void remove(Range range) { + public Range rangeContaining(long key, long value) { + return this.rangeContaining(consumer.apply(key, value)); + } + + public void remove(Range range) { set.remove(range); } + @Override + public void removeAtMost(long key, long value) { + set.remove(Range.atMost(consumer.apply(key, value))); + } + @Override public boolean isEmpty() { return set.isEmpty(); } @Override - public Range span() { + public Range span() { return set.span(); } @Override - public Set> asRanges() { + public Set> asRanges() { return set.asRanges(); } + + @Override + public boolean contains(long key, long value) { + return this.contains(consumer.apply(key, value)); + } + + @Override + public Range firstRange() { + return set.asRanges().iterator().next(); + } + + @Override + public int size() { + return set.asRanges().size(); + } + + @Override + public String toString() { + return set.toString(); + } } } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java deleted file mode 100644 index 1caee023a2033..0000000000000 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentLongPairRangeSetTest.java +++ /dev/null @@ -1,328 +0,0 @@ -/** - * 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.common.util.collections; - -import static org.testng.Assert.assertEquals; - -import java.util.List; -import java.util.Set; - -import org.apache.pulsar.common.util.collections.ConcurrentLongPairRangeSet.LongPair; -import org.testng.annotations.Test; - -import com.google.common.collect.BoundType; -import com.google.common.collect.Lists; -import com.google.common.collect.Range; -import com.google.common.collect.TreeRangeSet; - -public class ConcurrentLongPairRangeSetTest { - - @Test - public void testAddForSameKey() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - // add 0 to 5 - set.add(Range.closed(new LongPair(0, 0), new LongPair(0, 5))); - // add 8,9,10 - set.add(Range.closed(new LongPair(0, 8), new LongPair(0, 8))); - set.add(Range.closed(new LongPair(0, 9), new LongPair(0, 9))); - set.add(Range.closed(new LongPair(0, 10), new LongPair(0, 10))); - // add 98 to 99 and 102,105 - set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); - set.add(Range.closed(new LongPair(0, 102), new LongPair(0, 106))); - - List> ranges = set.asRanges(); - int count = 0; - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 0), new LongPair(0, 5)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 8), new LongPair(0, 10)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 99)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 102), new LongPair(0, 106)))); - } - - @Test - public void testAddForDifferentKey() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - // [98,100],[(1,5),(1,5)],[(1,10,1,15)],[(1,20),(1,20)],[(2,0),(2,10)] - set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); - set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); - set.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); - set.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); - - List> ranges = set.asRanges(); - int count = 0; - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 100)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 0), new LongPair(1, 5)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 10), new LongPair(1, 15)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 20), new LongPair(1, 20)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 0), new LongPair(2, 10)))); - } - - @Test - public void testAddCompareCompareWithGuava() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); - - // add 10K values for key 0 - int totalInsert = 10_000; - // add single values - for (int i = 0; i < totalInsert; i++) { - if (i % 3 == 0 || i % 6 == 0 || i % 8 == 0) { - set.add(Range.closed(new LongPair(0, i), new LongPair(0, i))); - gSet.add(Range.closed(new LongPair(0, i), new LongPair(0, i))); - } - } - // add batches - for (int i = totalInsert; i < (totalInsert * 2); i++) { - if (i % 5 == 0) { - set.add(Range.closed(new LongPair(0, i - 3), new LongPair(0, i + 3))); - gSet.add(Range.closed(new LongPair(0, i - 3), new LongPair(0, i + 3))); - } - } - List> ranges = set.asRanges(); - Set> gRanges = gSet.asRanges(); - - List> gRangeConnected = getConnectedRange(gRanges); - assertEquals(gRangeConnected.size(), ranges.size()); - int i = 0; - for (Range range : gRangeConnected) { - assertEquals(range, ranges.get(i)); - i++; - } - } - - @Test - public void testDeleteCompareWithGuava() { - - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); - - // add 10K values for key 0 - int totalInsert = 10_000; - // add single values - List> removedRanges = Lists.newArrayList(); - for (int i = 0; i < totalInsert; i++) { - if (i % 3 == 0 || i % 7 == 0 || i % 11 == 0) { - continue; - } - Range range = Range.closed(new LongPair(0, i), new LongPair(0, i)); - set.add(range); - gSet.add(range); - if (i % 4 == 0) { - removedRanges.add(range); - } - } - // add batches - for (int i = totalInsert; i < (totalInsert * 2); i++) { - Range range = Range.closed(new LongPair(0, i - 3), new LongPair(0, i + 3)); - if (i % 5 != 0) { - set.add(range); - gSet.add(range); - } - if (i % 4 == 0) { - removedRanges.add(range); - } - } - // remove records - for (Range range : removedRanges) { - set.remove(range); - gSet.remove(range); - } - - List> ranges = set.asRanges(); - Set> gRanges = gSet.asRanges(); - List> gRangeConnected = getConnectedRange(gRanges); - assertEquals(gRangeConnected.size(), ranges.size()); - int i = 0; - for (Range range : gRangeConnected) { - assertEquals(range, ranges.get(i)); - i++; - } - } - - @Test - public void testSpanWithGuava() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - 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(); - gSet.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); - gSet.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); - assertEquals(set.span(), gSet.span()); - assertEquals(set.span(), Range.closed(new LongPair(0, 98), new LongPair(1, 5))); - - 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))); - gSet.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); - gSet.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); - gSet.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); - gSet.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); - gSet.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); - assertEquals(set.span(), gSet.span()); - assertEquals(set.span(), Range.closed(new LongPair(0, 98), new LongPair(4, 20))); - - } - - @Test - public void testDeleteForDifferentKey() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); - set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); - 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))); - - // delete only (0,100) - set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); - - /** - * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] - */ - set.remove(Range.closed(new LongPair(2, 27), new LongPair(4, 15))); - - List> ranges = set.asRanges(); - int count = 0; - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 99)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 0), new LongPair(1, 5)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 10), new LongPair(1, 15)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 20), new LongPair(1, 20)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 0), new LongPair(2, 10)))); - - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 25), new LongPair(2, 26)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 28), new LongPair(2, 28)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(4, 16), new LongPair(4, 20)))); - } - - @Test - public void testDeleteWithAtMost() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); - set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); - 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))); - - // delete only (0,100) - set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); - - /** - * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] - */ - set.remove(Range.atMost(new LongPair(2, 27))); - - List> ranges = set.asRanges(); - int count = 0; - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 28), new LongPair(2, 28)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(3, 12), new LongPair(3, 20)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(4, 12), new LongPair(4, 20)))); - } - - @Test - public void testDeleteWithLeastMost() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); - set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); - 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))); - - // delete only (0,100) - set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); - - /** - * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] - */ - set.remove(Range.atLeast(new LongPair(2, 27))); - - List> ranges = set.asRanges(); - int count = 0; - assertEquals(ranges.get(count++), (Range.closed(new LongPair(0, 98), new LongPair(0, 99)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 0), new LongPair(1, 5)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 10), new LongPair(1, 15)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(1, 20), new LongPair(1, 20)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 0), new LongPair(2, 10)))); - assertEquals(ranges.get(count++), (Range.closed(new LongPair(2, 25), new LongPair(2, 26)))); - } - - @Test - public void testRangeContaining() { - ConcurrentLongPairRangeSet set = new ConcurrentLongPairRangeSet(); - 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(); - gSet.add(Range.closed(new LongPair(0, 98), new LongPair(0, 100))); - gSet.add(Range.closed(new LongPair(0, 101), new LongPair(1, 5))); - 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))); - gSet.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); - gSet.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); - gSet.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); - gSet.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); - gSet.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); - - LongPair position = new LongPair(0, 99); - assertEquals(set.rangeContaining(position), Range.closed(new LongPair(0, 98), new LongPair(0, 100))); - assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); - - position = new LongPair(2, 30); - assertEquals(set.rangeContaining(position), null); - assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); - - position = new LongPair(3, 13); - assertEquals(set.rangeContaining(position), Range.closed(new LongPair(3, 12), new LongPair(3, 20))); - assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); - - position = new LongPair(3, 22); - assertEquals(set.rangeContaining(position), null); - assertEquals(set.rangeContaining(position), gSet.rangeContaining(position)); - } - - private List> getConnectedRange(Set> gRanges) { - List> gRangeConnected = Lists.newArrayList(); - Range lastRange = null; - for (Range range : gRanges) { - if (lastRange == null) { - lastRange = range; - continue; - } - if ((lastRange.upperEndpoint().getValue() + 1) == (range.lowerEndpoint().getValue())) { - lastRange = Range.closed(lastRange.lowerEndpoint(), range.upperEndpoint()); - } else { - gRangeConnected.add(lastRange); - lastRange = range; - } - } - lastRange = lastRange.lowerBoundType().equals(BoundType.CLOSED) ? lastRange - : Range.closed( - new LongPair(lastRange.lowerEndpoint().getKey(), lastRange.lowerEndpoint().getValue() + 1), - lastRange.upperEndpoint()); - gRangeConnected.add(lastRange); - return gRangeConnected; - } -} 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 new file mode 100644 index 0000000000000..54bcecc6e6fb6 --- /dev/null +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -0,0 +1,384 @@ +/** + * 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.common.util.collections; + +import static org.testng.Assert.assertEquals; + +import java.util.List; +import java.util.Set; + +import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPair; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; +import org.testng.annotations.Test; + +import com.google.common.collect.BoundType; +import com.google.common.collect.Lists; +import com.google.common.collect.Range; +import com.google.common.collect.TreeRangeSet; + +public class ConcurrentOpenLongPairRangeSetTest { + + static final LongPairConsumer consumer = (key, value) -> new LongPair(key, value); + + @Test + public void testAddForSameKey() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + // add 0 to 5 + set.add(Range.closed(new LongPair(0, 0), new LongPair(0, 5))); + // add 8,9,10 + set.add(Range.closed(new LongPair(0, 8), new LongPair(0, 8))); + set.add(Range.closed(new LongPair(0, 9), new LongPair(0, 9))); + set.add(Range.closed(new LongPair(0, 10), new LongPair(0, 10))); + // add 98 to 99 and 102,105 + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 102), new LongPair(0, 106))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, -1), new LongPair(0, 5)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, 7), new LongPair(0, 10)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, 101), new LongPair(0, 106)))); + } + + @Test + public void testAddForDifferentKey() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + // [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); + set.addOpenClosed(1, 10, 1, 15); + set.addOpenClosed(1, 20, 2, 10); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, 98), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, -1), new LongPair(1, 5)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, 10), new LongPair(1, 15)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, -1), new LongPair(2, 10)))); + } + + @Test + public void testAddCompareCompareWithGuava() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); + + // add 10K values for key 0 + int totalInsert = 10_000; + // add single values + for (int i = 0; i < totalInsert; i++) { + if (i % 3 == 0 || i % 6 == 0 || i % 8 == 0) { + LongPair lower = new LongPair(0, i - 1); + LongPair upper = new LongPair(0, i); + // set.add(Range.openClosed(lower, upper)); + set.addOpenClosed(lower.getKey(), lower.getValue(), upper.getKey(), upper.getValue()); + gSet.add(Range.openClosed(lower, upper)); + } + } + // add batches + for (int i = totalInsert; i < (totalInsert * 2); i++) { + if (i % 5 == 0) { + LongPair lower = new LongPair(0, i - 3 - 1); + LongPair upper = new LongPair(0, i + 3); + // set.add(Range.openClosed(lower, upper)); + set.addOpenClosed(lower.getKey(), lower.getValue(), upper.getKey(), upper.getValue()); + gSet.add(Range.openClosed(lower, upper)); + } + } + List> ranges = set.asRanges(); + Set> gRanges = gSet.asRanges(); + + List> gRangeConnected = getConnectedRange(gRanges); + assertEquals(gRangeConnected.size(), ranges.size()); + int i = 0; + for (Range range : gRangeConnected) { + assertEquals(range, ranges.get(i)); + i++; + } + } + + @Test + public void testDeleteCompareWithGuava() { + + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); + + // add 10K values for key 0 + int totalInsert = 10_000; + // add single values + List> removedRanges = Lists.newArrayList(); + for (int i = 0; i < totalInsert; i++) { + if (i % 3 == 0 || i % 7 == 0 || i % 11 == 0) { + continue; + } + LongPair lower = new LongPair(0, i - 1); + LongPair upper = new LongPair(0, i); + Range range = Range.openClosed(lower, upper); + // set.add(range); + set.addOpenClosed(lower.getKey(), lower.getValue(), upper.getKey(), upper.getValue()); + gSet.add(range); + if (i % 4 == 0) { + removedRanges.add(range); + } + } + // add batches + for (int i = totalInsert; i < (totalInsert * 2); i++) { + LongPair lower = new LongPair(0, i - 3 - 1); + LongPair upper = new LongPair(0, i + 3); + Range range = Range.openClosed(lower, upper); + if (i % 5 != 0) { + // set.add(range); + set.addOpenClosed(lower.getKey(), lower.getValue(), upper.getKey(), upper.getValue()); + gSet.add(range); + } + if (i % 4 == 0) { + removedRanges.add(range); + } + } + // remove records + for (Range range : removedRanges) { + set.remove(range); + gSet.remove(range); + } + + List> ranges = set.asRanges(); + Set> gRanges = gSet.asRanges(); + List> gRangeConnected = getConnectedRange(gRanges); + assertEquals(gRangeConnected.size(), ranges.size()); + int i = 0; + for (Range range : gRangeConnected) { + assertEquals(range, ranges.get(i)); + i++; + } + } + + @Test + public void testSpanWithGuava() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + com.google.common.collect.RangeSet gSet = TreeRangeSet.create(); + set.add(Range.openClosed(new LongPair(0, 97), new LongPair(0, 99))); + gSet.add(Range.openClosed(new LongPair(0, 97), new LongPair(0, 99))); + set.add(Range.openClosed(new LongPair(0, 99), new LongPair(1, 5))); + gSet.add(Range.openClosed(new LongPair(0, 99), new LongPair(1, 5))); + assertEquals(set.span(), gSet.span()); + assertEquals(set.span(), Range.openClosed(new LongPair(0, 97), new LongPair(1, 5))); + + set.add(Range.openClosed(new LongPair(1, 9), new LongPair(1, 15))); + set.add(Range.openClosed(new LongPair(1, 19), new LongPair(2, 10))); + set.add(Range.openClosed(new LongPair(2, 24), new LongPair(2, 28))); + set.add(Range.openClosed(new LongPair(3, 11), new LongPair(3, 20))); + set.add(Range.openClosed(new LongPair(4, 11), new LongPair(4, 20))); + gSet.add(Range.openClosed(new LongPair(1, 9), new LongPair(1, 15))); + gSet.add(Range.openClosed(new LongPair(1, 19), new LongPair(2, 10))); + gSet.add(Range.openClosed(new LongPair(2, 24), new LongPair(2, 28))); + gSet.add(Range.openClosed(new LongPair(3, 11), new LongPair(3, 20))); + gSet.add(Range.openClosed(new LongPair(4, 11), new LongPair(4, 20))); + assertEquals(set.span(), gSet.span()); + assertEquals(set.span(), Range.openClosed(new LongPair(0, 97), new LongPair(4, 20))); + } + + @Test + public void testFirstRange() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + Range range = Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)); + set.add(range); + assertEquals(set.firstRange(), range); + assertEquals(set.size(), 1); + range = Range.openClosed(new LongPair(0, 98), new LongPair(0, 105)); + set.add(range); + assertEquals(set.firstRange(), Range.openClosed(new LongPair(0, 97), new LongPair(0, 105))); + assertEquals(set.size(), 1); + range = Range.openClosed(new LongPair(0, 5), new LongPair(0, 75)); + set.add(range); + assertEquals(set.firstRange(), range); + assertEquals(set.size(), 2); + } + + @Test + public void testToString() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + Range range = Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)); + set.add(range); + assertEquals(set.toString(), "[(0:97..0:99]]"); + range = Range.openClosed(new LongPair(0, 98), new LongPair(0, 105)); + set.add(range); + assertEquals(set.toString(), "[(0:97..0:105]]"); + range = Range.openClosed(new LongPair(0, 5), new LongPair(0, 75)); + set.add(range); + assertEquals(set.toString(), "[(0:5..0:75],(0:97..0:105]]"); + } + + @Test + public void testDeleteForDifferentKey() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + set.addOpenClosed(0, 97, 0, 99); + set.addOpenClosed(0, 99, 1, 5); + set.addOpenClosed(1, 9, 1, 15); + set.addOpenClosed(1, 19, 2, 10); + set.addOpenClosed(2, 24, 2, 28); + set.addOpenClosed(3, 11, 3, 20); + set.addOpenClosed(4, 11, 4, 20); + + // delete only (0,100) + set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); + + /** + * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] + */ + set.remove(Range.closed(new LongPair(2, 27), new LongPair(4, 15))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, -1), new LongPair(1, 5)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, 9), new LongPair(1, 15)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, -1), new LongPair(2, 10)))); + + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, 24), new LongPair(2, 26)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, 27), new LongPair(2, 28)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(4, 15), new LongPair(4, 20)))); + } + + @Test + public void testDeleteWithAtMost() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + 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))); + + // delete only (0,100) + set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); + + /** + * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] + */ + set.remove(Range.atMost(new LongPair(2, 27))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, 27), new LongPair(2, 28)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(3, 11), new LongPair(3, 20)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(4, 11), new LongPair(4, 20)))); + } + + @Test + public void testDeleteWithLeastMost() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + set.add(Range.closed(new LongPair(0, 98), new LongPair(0, 99))); + set.add(Range.closed(new LongPair(0, 100), new LongPair(1, 5))); + 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))); + + // delete only (0,100) + set.remove(Range.open(new LongPair(0, 99), new LongPair(0, 105))); + + /** + * delete all keys from [2,27]->[4,15] : remaining [2,25..26,28], [4,16..20] + */ + set.remove(Range.atLeast(new LongPair(2, 27))); + + List> ranges = set.asRanges(); + int count = 0; + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(0, 97), new LongPair(0, 99)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, -1), new LongPair(1, 5)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, 9), new LongPair(1, 15)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(1, 19), new LongPair(1, 20)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, -1), new LongPair(2, 10)))); + assertEquals(ranges.get(count++), (Range.openClosed(new LongPair(2, 24), new LongPair(2, 26)))); + } + + @Test + public void testRangeContaining() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + 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(); + gSet.add(Range.closed(new LongPair(0, 98), new LongPair(0, 100))); + gSet.add(Range.closed(new LongPair(0, 101), new LongPair(1, 5))); + 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))); + gSet.add(Range.closed(new LongPair(1, 10), new LongPair(1, 15))); + gSet.add(Range.closed(new LongPair(1, 20), new LongPair(2, 10))); + gSet.add(Range.closed(new LongPair(2, 25), new LongPair(2, 28))); + gSet.add(Range.closed(new LongPair(3, 12), new LongPair(3, 20))); + gSet.add(Range.closed(new LongPair(4, 12), new LongPair(4, 20))); + + LongPair position = new LongPair(0, 99); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), + Range.closed(new LongPair(0, 98), new LongPair(0, 100))); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), gSet.rangeContaining(position)); + + position = new LongPair(2, 30); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), null); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), gSet.rangeContaining(position)); + + position = new LongPair(3, 13); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), + Range.closed(new LongPair(3, 12), new LongPair(3, 20))); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), gSet.rangeContaining(position)); + + position = new LongPair(3, 22); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), null); + assertEquals(set.rangeContaining(position.getKey(), position.getValue()), gSet.rangeContaining(position)); + } + + private List> getConnectedRange(Set> gRanges) { + List> gRangeConnected = Lists.newArrayList(); + Range lastRange = null; + for (Range range : gRanges) { + if (lastRange == null) { + lastRange = range; + continue; + } + LongPair previousUpper = lastRange.upperEndpoint(); + LongPair currentLower = range.lowerEndpoint(); + int previousUpperValue = (int) (lastRange.upperBoundType().equals(BoundType.CLOSED) + ? previousUpper.getValue() + : previousUpper.getValue() - 1); + int currentLowerValue = (int) (range.lowerBoundType().equals(BoundType.CLOSED) ? currentLower.getValue() + : currentLower.getValue() + 1); + boolean connected = (previousUpper.getKey() == currentLower.getKey()) + ? (previousUpperValue >= currentLowerValue) + : false; + if (connected) { + lastRange = Range.closed(lastRange.lowerEndpoint(), range.upperEndpoint()); + } else { + gRangeConnected.add(lastRange); + lastRange = range; + } + } + int lowerOpenValue = (int) (lastRange.lowerBoundType().equals(BoundType.CLOSED) + ? (lastRange.lowerEndpoint().getValue() - 1) + : lastRange.lowerEndpoint().getValue()); + lastRange = Range.openClosed(new LongPair(lastRange.lowerEndpoint().getKey(), lowerOpenValue), + lastRange.upperEndpoint()); + gRangeConnected.add(lastRange); + return gRangeConnected; + } +} From ec7e33ab5b0341c8db041836591356625875e615 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Thu, 14 Mar 2019 12:43:40 -0700 Subject: [PATCH 3/6] add forEach --- .../ConcurrentOpenLongPairRangeSet.java | 104 ++++++++++-------- .../util/collections/LongPairRangeSet.java | 44 +++++++- 2 files changed, 100 insertions(+), 48 deletions(-) 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 51ea3cec55bbc..7ed1af71abb7a 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 @@ -30,6 +30,9 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.StampedLock; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; +import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeProcessor; + import com.google.common.collect.BoundType; import com.google.common.collect.Range; @@ -284,10 +287,44 @@ public Range span() { @Override public List> asRanges() { List> ranges = new ArrayList<>(); - createRanges(ranges, null); + forEach((range) -> { + ranges.add(range); + return true; + }); return ranges; } + @Override + public void forEach(RangeProcessor action) { + forEach(action, consumer); + } + + @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); + } + }); + } + @Override public Range firstRange() { Entry firstSet = rangeBitSetMap.firstEntry(); @@ -299,7 +336,12 @@ public Range firstRange() { @Override public int size() { if (updatedAfterCached) { - cachedSize = createRanges(null, null); + AtomicInteger size = new AtomicInteger(0); + forEach((range) -> { + size.getAndIncrement(); + return true; + }); + cachedSize = size.get(); updatedAfterCached = false; } return cachedSize; @@ -309,55 +351,23 @@ public int size() { public String toString() { if (updatedAfterCached) { StringBuilder toString = new StringBuilder(); - createRanges(null, toString); - cachedToString = toString.toString(); - updatedAfterCached = false; - } - return cachedToString; - } - - /** - * It creates ranges and add into given list {@code ranges} and also returns total number of ranges into the set. If - * given List {@code ranges} is null then it just returns number of ranges. - * - * @param ranges - * @return - */ - private int createRanges(List> ranges, StringBuilder toString) { - AtomicInteger size = new AtomicInteger(0); - if (toString != null) { - toString.append("["); - } - rangeBitSetMap.forEach((key, set) -> { - if (set.isEmpty()) { - return; + AtomicBoolean first = new AtomicBoolean(true); + if (toString != null) { + toString.append("["); } - int first = set.nextSetBit(0); - int last = set.previousSetBit(set.size()); - int currentClosedMark = first; - // TODO: remove: sometime previous-keyset (previous ledger) is connected to the current one, in that case - // merge the range - 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 (ranges != null) { - ranges.add(range); - } - if (toString != null) { - if (size.get() != 0) { - toString.append(","); - } - toString.append(range); + forEach((range) -> { + if (!first.get()) { + toString.append(","); } - size.getAndIncrement(); - currentClosedMark = set.nextSetBit(nextOpenMark); - } - }); - if (toString != null) { + toString.append(range); + first.set(false); + return true; + }); toString.append("]"); + cachedToString = toString.toString(); + updatedAfterCached = false; } - return size.get(); + return cachedToString; } /** 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 02736a503d9ef..0f4a48d7ac381 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 @@ -81,7 +81,26 @@ public interface LongPairRangeSet> { * @return */ Collection> asRanges(); - + + /** + * 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.) + * + * @param action + */ + void forEach(RangeProcessor action); + + /** + * 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.) + * + * @param action + * @param consumer + */ + void forEach(RangeProcessor action, LongPairConsumer consumer); + /** * Returns total number of ranges into the set. * @@ -100,6 +119,15 @@ public static interface LongPairConsumer { T apply(long key, long value); } + public static interface RangeProcessor> { + /** + * + * @param range + * @return false if there is no further processing required + */ + boolean process(Range range); + } + public static class LongPair implements Comparable { public static final LongPair earliest = new LongPair(-1, -1); @@ -189,6 +217,20 @@ public Set> asRanges() { return set.asRanges(); } + @Override + public void forEach(RangeProcessor action) { + forEach(action, consumer); + } + + @Override + public void forEach(RangeProcessor action, LongPairConsumer consumer) { + for (Range range : asRanges()) { + if (!action.process(range)) { + break; + } + } + } + @Override public boolean contains(long key, long value) { return this.contains(consumer.apply(key, value)); From 01e873fe14936c321c587002f7d8f40b648d69b4 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Thu, 14 Mar 2019 21:55:12 -0700 Subject: [PATCH 4/6] add forEach with consumer --- .../util/collections/ConcurrentOpenLongPairRangeSet.java | 2 +- .../pulsar/common/util/collections/LongPairRangeSet.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) 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 7ed1af71abb7a..36be51f3d2ed7 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,7 @@ public void forEach(RangeProcessor action) { } @Override - public void forEach(RangeProcessor action, LongPairConsumer consumer) { + public void forEach(RangeProcessor action, LongPairConsumer consumer) { 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 0f4a48d7ac381..daad0ad2a54f2 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 @@ -99,7 +99,7 @@ public interface LongPairRangeSet> { * @param action * @param consumer */ - void forEach(RangeProcessor action, LongPairConsumer consumer); + void forEach(RangeProcessor action, LongPairConsumer consumer); /** * Returns total number of ranges into the set. @@ -223,7 +223,7 @@ public void forEach(RangeProcessor action) { } @Override - public void forEach(RangeProcessor action, LongPairConsumer consumer) { + public void forEach(RangeProcessor action, LongPairConsumer consumer) { for (Range range : asRanges()) { if (!action.process(range)) { break; From 3cb7a7017b2d19e808806cd4f843ec6a5e740e17 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Fri, 17 May 2019 13:56:06 -0700 Subject: [PATCH 5/6] Fix stamp-lock usage --- .../ConcurrentOpenLongPairRangeSet.java | 103 ++++++++++++------ 1 file changed, 69 insertions(+), 34 deletions(-) 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 36be51f3d2ed7..392faf4b383de 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 @@ -30,9 +30,6 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.StampedLock; -import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; -import org.apache.pulsar.common.util.collections.LongPairRangeSet.RangeProcessor; - import com.google.common.collect.BoundType; import com.google.common.collect.Range; @@ -95,67 +92,105 @@ public ConcurrentBitSet(int nbits) { @Override public boolean get(int bitIndex) { - return super.get(bitIndex); + long stamp = rwLock.tryOptimisticRead(); + boolean isSet = super.get(bitIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + isSet = super.get(bitIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + return isSet; } @Override public void set(int bitIndex) { - long stamp = rwLock.writeLock(); - try { - super.set(bitIndex); - } finally { - rwLock.unlockWrite(stamp); + long stamp = rwLock.tryOptimisticRead(); + super.set(bitIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + super.set(bitIndex); + } finally { + rwLock.unlockRead(stamp); + } } } @Override public void set(int fromIndex, int toIndex) { - long stamp = rwLock.writeLock(); - try { - super.set(fromIndex, toIndex); - } finally { - rwLock.unlockWrite(stamp); + long stamp = rwLock.tryOptimisticRead(); + super.set(fromIndex, toIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + super.set(fromIndex, toIndex); + } finally { + rwLock.unlockRead(stamp); + } } } @Override public int nextSetBit(int fromIndex) { - long stamp = rwLock.readLock(); - try { - return super.nextSetBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); + long stamp = rwLock.tryOptimisticRead(); + int bit = super.nextSetBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.nextSetBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } } + return bit; } @Override public int nextClearBit(int fromIndex) { - long stamp = rwLock.readLock(); - try { - return super.nextClearBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); + long stamp = rwLock.tryOptimisticRead(); + int bit = super.nextClearBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.nextClearBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } } + return bit; } @Override public int previousSetBit(int fromIndex) { - long stamp = rwLock.readLock(); - try { - return super.previousSetBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); + long stamp = rwLock.tryOptimisticRead(); + int bit = super.previousSetBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.previousSetBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } } + return bit; } @Override public int previousClearBit(int fromIndex) { - long stamp = rwLock.readLock(); - try { - return super.previousClearBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); + long stamp = rwLock.tryOptimisticRead(); + int bit = super.previousClearBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.previousClearBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } } + return bit; } @Override From e462debd8b1ec14942287ac37a587859b31fe2e7 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Fri, 17 May 2019 13:59:28 -0700 Subject: [PATCH 6/6] Move ConcurrentBitSet to separate class --- .../util/collections/ConcurrentBitSet.java | 160 ++++++++++++++++++ .../ConcurrentOpenLongPairRangeSet.java | 139 --------------- 2 files changed, 160 insertions(+), 139 deletions(-) create mode 100644 pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentBitSet.java diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentBitSet.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentBitSet.java new file mode 100644 index 0000000000000..739ad1e393b7e --- /dev/null +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentBitSet.java @@ -0,0 +1,160 @@ +/** + * 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.common.util.collections; + +import java.util.BitSet; +import java.util.concurrent.locks.StampedLock; + +public class ConcurrentBitSet extends BitSet { + + private static final long serialVersionUID = 1L; + private final StampedLock rwLock = new StampedLock(); + + /** + * Creates a bit set whose initial size is large enough to explicitly represent bits with indices in the range + * {@code 0} through {@code nbits-1}. All bits are initially {@code false}. + * + * @param nbits + * the initial size of the bit set + * @throws NegativeArraySizeException + * if the specified initial size is negative + */ + public ConcurrentBitSet(int nbits) { + super(nbits); + } + + @Override + public boolean get(int bitIndex) { + long stamp = rwLock.tryOptimisticRead(); + boolean isSet = super.get(bitIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + isSet = super.get(bitIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + return isSet; + } + + @Override + public void set(int bitIndex) { + long stamp = rwLock.tryOptimisticRead(); + super.set(bitIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + super.set(bitIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + } + + @Override + public void set(int fromIndex, int toIndex) { + long stamp = rwLock.tryOptimisticRead(); + super.set(fromIndex, toIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + super.set(fromIndex, toIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + } + + @Override + public int nextSetBit(int fromIndex) { + long stamp = rwLock.tryOptimisticRead(); + int bit = super.nextSetBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.nextSetBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + return bit; + } + + @Override + public int nextClearBit(int fromIndex) { + long stamp = rwLock.tryOptimisticRead(); + int bit = super.nextClearBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.nextClearBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + return bit; + } + + @Override + public int previousSetBit(int fromIndex) { + long stamp = rwLock.tryOptimisticRead(); + int bit = super.previousSetBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.previousSetBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + return bit; + } + + @Override + public int previousClearBit(int fromIndex) { + long stamp = rwLock.tryOptimisticRead(); + int bit = super.previousClearBit(fromIndex); + if (!rwLock.validate(stamp)) { + stamp = rwLock.readLock(); + try { + bit = super.previousClearBit(fromIndex); + } finally { + rwLock.unlockRead(stamp); + } + } + return bit; + } + + @Override + public boolean isEmpty() { + long stamp = rwLock.tryOptimisticRead(); + boolean isEmpty = super.isEmpty(); + if (!rwLock.validate(stamp)) { + // Fallback to read lock + stamp = rwLock.readLock(); + try { + isEmpty = super.isEmpty(); + } finally { + rwLock.unlockRead(stamp); + } + } + return isEmpty; + } +} \ No newline at end of file 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 392faf4b383de..a88b8bf3c3cf1 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 @@ -28,7 +28,6 @@ import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.locks.StampedLock; import com.google.common.collect.BoundType; import com.google.common.collect.Range; @@ -73,144 +72,6 @@ public ConcurrentOpenLongPairRangeSet(int size, boolean threadSafe, LongPairCons this.consumer = consumer; } - class ConcurrentBitSet extends BitSet { - private static final long serialVersionUID = 1L; - private final StampedLock rwLock = new StampedLock(); - - /** - * Creates a bit set whose initial size is large enough to explicitly represent bits with indices in the range - * {@code 0} through {@code nbits-1}. All bits are initially {@code false}. - * - * @param nbits - * the initial size of the bit set - * @throws NegativeArraySizeException - * if the specified initial size is negative - */ - public ConcurrentBitSet(int nbits) { - super(nbits); - } - - @Override - public boolean get(int bitIndex) { - long stamp = rwLock.tryOptimisticRead(); - boolean isSet = super.get(bitIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - isSet = super.get(bitIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - return isSet; - } - - @Override - public void set(int bitIndex) { - long stamp = rwLock.tryOptimisticRead(); - super.set(bitIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - super.set(bitIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - } - - @Override - public void set(int fromIndex, int toIndex) { - long stamp = rwLock.tryOptimisticRead(); - super.set(fromIndex, toIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - super.set(fromIndex, toIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - } - - @Override - public int nextSetBit(int fromIndex) { - long stamp = rwLock.tryOptimisticRead(); - int bit = super.nextSetBit(fromIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - bit = super.nextSetBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - return bit; - } - - @Override - public int nextClearBit(int fromIndex) { - long stamp = rwLock.tryOptimisticRead(); - int bit = super.nextClearBit(fromIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - bit = super.nextClearBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - return bit; - } - - @Override - public int previousSetBit(int fromIndex) { - long stamp = rwLock.tryOptimisticRead(); - int bit = super.previousSetBit(fromIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - bit = super.previousSetBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - return bit; - } - - @Override - public int previousClearBit(int fromIndex) { - long stamp = rwLock.tryOptimisticRead(); - int bit = super.previousClearBit(fromIndex); - if (!rwLock.validate(stamp)) { - stamp = rwLock.readLock(); - try { - bit = super.previousClearBit(fromIndex); - } finally { - rwLock.unlockRead(stamp); - } - } - return bit; - } - - @Override - public boolean isEmpty() { - long stamp = rwLock.tryOptimisticRead(); - boolean isEmpty = super.isEmpty(); - if (!rwLock.validate(stamp)) { - // Fallback to read lock - stamp = rwLock.readLock(); - try { - isEmpty = super.isEmpty(); - } finally { - rwLock.unlockRead(stamp); - } - } - return isEmpty; - } - - } - /** * Adds the specified range to this {@code RangeSet} (optional operation). That is, for equal range sets a and b, * the result of {@code a.add(range)} is that {@code a} will be the minimal range set for which both