From 3817d5cfaac71eaaf4c55599536de7f92cb59cd4 Mon Sep 17 00:00:00 2001 From: penghui Date: Sun, 4 Sep 2022 13:34:44 +0800 Subject: [PATCH 1/5] [improve][broker] Improve cursor.getNumberOfEntries if isUnackedRangesOpenCacheSetEnabled=true --- .../mledger/impl/ManagedCursorImpl.java | 49 ++++++++++++------- .../mledger/impl/RangeSetWrapper.java | 6 +++ .../systopic/PartitionedSystemTopicTest.java | 22 +++++++++ .../ConcurrentOpenLongPairRangeSet.java | 29 +++++++++++ .../util/collections/LongPairRangeSet.java | 11 +++++ .../ConcurrentOpenLongPairRangeSetTest.java | 27 ++++++++++ 6 files changed, 127 insertions(+), 17 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 8262118ee4787..b963dbce9a62b 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -1455,25 +1455,40 @@ protected long getNumberOfEntries(Range range) { lock.readLock().lock(); try { - individualDeletedMessages.forEach((r) -> { - try { - if (r.isConnected(range)) { - Range commonEntries = r.intersection(range); - long commonCount = ledger.getNumberOfEntries(commonEntries); - if (log.isDebugEnabled()) { - log.debug("[{}] [{}] Discounting {} entries for already deleted range {}", ledger.getName(), - name, commonCount, commonEntries); + if (config.isUnackedRangesOpenCacheSetEnabled()) { + Map cardinalityMap = individualDeletedMessages.cardinality( + range.lowerEndpoint().ledgerId, range.lowerEndpoint().entryId, + range.upperEndpoint().ledgerId, range.upperEndpoint().entryId); + deletedEntries.addAndGet(cardinalityMap.values().stream().mapToInt(v -> v).sum()); + deletedEntries.addAndGet( + ledger.ledgers.subMap(range.lowerEndpoint().ledgerId, true, + range.upperEndpoint().ledgerId, true) + .values() + .stream() + .filter(ledgerInfo -> !cardinalityMap.containsKey(ledgerInfo.getLedgerId())) + .mapToLong(LedgerInfo::getEntries).sum() + ); + } else { + individualDeletedMessages.forEach((r) -> { + try { + if (r.isConnected(range)) { + Range commonEntries = r.intersection(range); + long commonCount = ledger.getNumberOfEntries(commonEntries); + if (log.isDebugEnabled()) { + log.debug("[{}] [{}] Discounting {} entries for already deleted range {}", + ledger.getName(), name, commonCount, commonEntries); + } + deletedEntries.addAndGet(commonCount); + } + return true; + } finally { + if (r.lowerEndpoint() instanceof PositionImplRecyclable) { + ((PositionImplRecyclable) r.lowerEndpoint()).recycle(); + ((PositionImplRecyclable) r.upperEndpoint()).recycle(); } - deletedEntries.addAndGet(commonCount); - } - return true; - } finally { - if (r.lowerEndpoint() instanceof PositionImplRecyclable) { - ((PositionImplRecyclable) r.lowerEndpoint()).recycle(); - ((PositionImplRecyclable) r.upperEndpoint()).recycle(); } - } - }, recyclePositionRangeConverter); + }, recyclePositionRangeConverter); + } } finally { lock.readLock().unlock(); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java index e5337822bff81..4dd3eee7a9b6f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java @@ -24,6 +24,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.List; +import java.util.Map; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.pulsar.common.util.collections.ConcurrentOpenLongPairRangeSet; import org.apache.pulsar.common.util.collections.LongPairRangeSet; @@ -133,6 +134,11 @@ public Range lastRange() { return rangeSet.lastRange(); } + @Override + public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { + return rangeSet.cardinality(lowerKey, lowerValue, upperKey, upperValue); + } + @VisibleForTesting void add(Range range) { if (!(rangeSet instanceof ConcurrentOpenLongPairRangeSet)) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java index d4ed12573f3db..b930fa48f9240 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java @@ -235,4 +235,26 @@ private void testSetBacklogCausedCreatingProducerFailure() throws Exception { Assert.fail("failed to create producer"); } } + + @Test + public void testGetTopic() throws Exception { + final String ns = "prop/ns-test"; + admin.namespaces().createNamespace(ns, 2); + final String topicName = ns + "/topic-1"; + admin.topics().createNonPartitionedTopic(String.format("persistent://%s", topicName)); + Producer producer1 = pulsarClient.newProducer(Schema.STRING).topic(topicName).create(); + producer1.close(); + PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getTopic(topicName.toString(), false).get().get(); + persistentTopic.close().join(); + List topics = new ArrayList<>(pulsar.getBrokerService().getTopics().keys()); + topics.removeIf(item -> item.contains(SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME)); + Assert.assertEquals(topics.size(), 0); + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topicName) + .subscriptionName("sub-1") + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .subscriptionType(SubscriptionType.Shared) + .subscribe(); + } } 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 05d16c4b054e0..766ba7ef35750 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 @@ -23,7 +23,9 @@ import com.google.common.collect.Range; import java.util.ArrayList; import java.util.BitSet; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Map.Entry; import java.util.NavigableMap; import java.util.concurrent.ConcurrentSkipListMap; @@ -242,6 +244,33 @@ public Range lastRange() { return Range.openClosed(consumer.apply(lastSet.getKey(), lower), consumer.apply(lastSet.getKey(), upper)); } + @Override + public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { + NavigableMap subMap = rangeBitSetMap.subMap(lowerKey, true, upperKey, true); + Map v = new HashMap<>(); + subMap.forEach((ledgerId, bitset) -> { + boolean isLowerOrUpper = false; + BitSet lowerBitSet = null; + if (ledgerId == lowerKey) { + isLowerOrUpper = true; + BitSet temp = (BitSet) bitset.clone(); + temp.clear(0, (int) Math.max(0, lowerValue)); + lowerBitSet = temp; + v.put(ledgerId, temp.cardinality()); + } + if (ledgerId == upperKey) { + isLowerOrUpper = true; + BitSet temp = lowerBitSet == null ? (BitSet) bitset.clone() : lowerBitSet; + temp.clear((int) Math.min(upperValue + 1, temp.length()), temp.length()); + v.put(ledgerId, temp.cardinality()); + } + if (!isLowerOrUpper) { + v.put(ledgerId, bitset.cardinality()); + } + }); + return v; + } + @Override public int size() { if (updatedAfterCachedForSize) { 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 ba77ff4b839a8..e8a96a5fe0695 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 @@ -25,6 +25,7 @@ import java.util.Collection; import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.NoSuchElementException; import java.util.Set; import lombok.EqualsAndHashCode; @@ -125,6 +126,11 @@ public interface LongPairRangeSet> { */ Range lastRange(); + /** + * Return the number bit sets to true for each key from lower to upper. + */ + Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue); + /** * Represents a function that accepts two long arguments and produces a result. * @@ -296,6 +302,11 @@ public Range lastRange() { return list.get(list.size() - 1); } + @Override + public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { + throw new UnsupportedOperationException(); + } + @Override public int size() { return set.asRanges().size(); diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java index 2c0b8d3552c1b..4aa911df4533c 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -25,6 +25,7 @@ import java.util.ArrayList; import java.util.List; +import java.util.Map; import java.util.Set; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPair; @@ -460,4 +461,30 @@ private List> getConnectedRange(Set> gRanges) { gRangeConnected.add(lastRange); return gRangeConnected; } + + @Test + public void testCardinality() { + ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); + Map v = set.cardinality(0, 0, Integer.MAX_VALUE, Integer.MAX_VALUE); + assertEquals(v.size(), 0 ); + set.addOpenClosed(1, 0, 1, 20); + set.addOpenClosed(1, 30, 1, 90); + set.addOpenClosed(2, 0, 3, 30); + v = set.cardinality(1, 0, 1, 100); + assertEquals(v.size(), 1); + assertEquals((int) v.get(1L), 80); + v = set.cardinality(1, 11, 1, 100); + assertEquals(v.size(), 1); + assertEquals((int) v.get(1L), 70); + v = set.cardinality(1, 0, 1, 90); + assertEquals(v.size(), 1); + assertEquals((int) v.get(1L), 80); + v = set.cardinality(1, 0, 1, 80); + assertEquals(v.size(), 1); + assertEquals((int) v.get(1L), 70); + v = set.cardinality(1, 0, 3, 30); + assertEquals(v.size(), 2); + assertEquals((int) v.get(1L), 80); + assertEquals((int) v.get(3L), 31); + } } From 92ccc4046291eaf6c6972b26d543347d72259b1b Mon Sep 17 00:00:00 2001 From: penghui Date: Sun, 4 Sep 2022 14:15:39 +0800 Subject: [PATCH 2/5] Fix the case that cardinality map is empty --- .../mledger/impl/ManagedCursorImpl.java | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index b963dbce9a62b..4b11138a7017f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -1460,14 +1460,18 @@ protected long getNumberOfEntries(Range range) { range.lowerEndpoint().ledgerId, range.lowerEndpoint().entryId, range.upperEndpoint().ledgerId, range.upperEndpoint().entryId); deletedEntries.addAndGet(cardinalityMap.values().stream().mapToInt(v -> v).sum()); - deletedEntries.addAndGet( - ledger.ledgers.subMap(range.lowerEndpoint().ledgerId, true, - range.upperEndpoint().ledgerId, true) - .values() - .stream() - .filter(ledgerInfo -> !cardinalityMap.containsKey(ledgerInfo.getLedgerId())) - .mapToLong(LedgerInfo::getEntries).sum() - ); + if (cardinalityMap.isEmpty()) { + deletedEntries.addAndGet(ledger.getNumberOfEntries(range)); + } else { + deletedEntries.addAndGet( + ledger.ledgers.subMap(range.lowerEndpoint().ledgerId, true, + range.upperEndpoint().ledgerId, true) + .values() + .stream() + .filter(ledgerInfo -> !cardinalityMap.containsKey(ledgerInfo.getLedgerId())) + .mapToLong(LedgerInfo::getEntries).sum() + ); + } } else { individualDeletedMessages.forEach((r) -> { try { From 63c230e713c3b162b2c496129f5a14fe03303bca Mon Sep 17 00:00:00 2001 From: penghui Date: Mon, 5 Sep 2022 08:33:31 +0800 Subject: [PATCH 3/5] Fix the case that cardinality map is empty --- .../apache/bookkeeper/mledger/impl/ManagedCursorImpl.java | 6 ++---- .../util/collections/ConcurrentOpenLongPairRangeSet.java | 4 ++++ 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 4b11138a7017f..f09364c0d79c4 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -1459,10 +1459,8 @@ protected long getNumberOfEntries(Range range) { Map cardinalityMap = individualDeletedMessages.cardinality( range.lowerEndpoint().ledgerId, range.lowerEndpoint().entryId, range.upperEndpoint().ledgerId, range.upperEndpoint().entryId); - deletedEntries.addAndGet(cardinalityMap.values().stream().mapToInt(v -> v).sum()); - if (cardinalityMap.isEmpty()) { - deletedEntries.addAndGet(ledger.getNumberOfEntries(range)); - } else { + if (!cardinalityMap.isEmpty()) { + deletedEntries.addAndGet(cardinalityMap.values().stream().mapToInt(v -> v).sum()); deletedEntries.addAndGet( ledger.ledgers.subMap(range.lowerEndpoint().ledgerId, true, range.upperEndpoint().ledgerId, true) 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 766ba7ef35750..faad6c8956de1 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 @@ -23,6 +23,7 @@ import com.google.common.collect.Range; import java.util.ArrayList; import java.util.BitSet; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -247,6 +248,9 @@ public Range lastRange() { @Override public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { NavigableMap subMap = rangeBitSetMap.subMap(lowerKey, true, upperKey, true); + if (subMap.isEmpty()) { + return Collections.emptyMap(); + } Map v = new HashMap<>(); subMap.forEach((ledgerId, bitset) -> { boolean isLowerOrUpper = false; From d0fad7eefc6e36d8f69004986bc539d17aa7e49f Mon Sep 17 00:00:00 2001 From: penghui Date: Mon, 5 Sep 2022 11:18:07 +0800 Subject: [PATCH 4/5] Remove return the ledger ID --- .../mledger/impl/ManagedCursorImpl.java | 14 +------ .../mledger/impl/RangeSetWrapper.java | 3 +- .../systopic/PartitionedSystemTopicTest.java | 22 ----------- .../ConcurrentOpenLongPairRangeSet.java | 38 +++++++------------ .../util/collections/LongPairRangeSet.java | 7 ++-- .../ConcurrentOpenLongPairRangeSetTest.java | 21 ++++------ 6 files changed, 26 insertions(+), 79 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index f09364c0d79c4..59da6fc81da70 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -1456,20 +1456,10 @@ protected long getNumberOfEntries(Range range) { lock.readLock().lock(); try { if (config.isUnackedRangesOpenCacheSetEnabled()) { - Map cardinalityMap = individualDeletedMessages.cardinality( + int cardinality = individualDeletedMessages.cardinality( range.lowerEndpoint().ledgerId, range.lowerEndpoint().entryId, range.upperEndpoint().ledgerId, range.upperEndpoint().entryId); - if (!cardinalityMap.isEmpty()) { - deletedEntries.addAndGet(cardinalityMap.values().stream().mapToInt(v -> v).sum()); - deletedEntries.addAndGet( - ledger.ledgers.subMap(range.lowerEndpoint().ledgerId, true, - range.upperEndpoint().ledgerId, true) - .values() - .stream() - .filter(ledgerInfo -> !cardinalityMap.containsKey(ledgerInfo.getLedgerId())) - .mapToLong(LedgerInfo::getEntries).sum() - ); - } + deletedEntries.addAndGet(cardinality); } else { individualDeletedMessages.forEach((r) -> { try { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java index 4dd3eee7a9b6f..e39572698461e 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RangeSetWrapper.java @@ -24,7 +24,6 @@ import java.util.ArrayList; import java.util.Collection; import java.util.List; -import java.util.Map; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.pulsar.common.util.collections.ConcurrentOpenLongPairRangeSet; import org.apache.pulsar.common.util.collections.LongPairRangeSet; @@ -135,7 +134,7 @@ public Range lastRange() { } @Override - public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { + public int cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { return rangeSet.cardinality(lowerKey, lowerValue, upperKey, upperValue); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java index b930fa48f9240..d4ed12573f3db 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java @@ -235,26 +235,4 @@ private void testSetBacklogCausedCreatingProducerFailure() throws Exception { Assert.fail("failed to create producer"); } } - - @Test - public void testGetTopic() throws Exception { - final String ns = "prop/ns-test"; - admin.namespaces().createNamespace(ns, 2); - final String topicName = ns + "/topic-1"; - admin.topics().createNonPartitionedTopic(String.format("persistent://%s", topicName)); - Producer producer1 = pulsarClient.newProducer(Schema.STRING).topic(topicName).create(); - producer1.close(); - PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getTopic(topicName.toString(), false).get().get(); - persistentTopic.close().join(); - List topics = new ArrayList<>(pulsar.getBrokerService().getTopics().keys()); - topics.removeIf(item -> item.contains(SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME)); - Assert.assertEquals(topics.size(), 0); - @Cleanup - Consumer consumer = pulsarClient.newConsumer(Schema.STRING) - .topic(topicName) - .subscriptionName("sub-1") - .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) - .subscriptionType(SubscriptionType.Shared) - .subscribe(); - } } 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 faad6c8956de1..b13d384b1d6c6 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 @@ -23,10 +23,7 @@ import com.google.common.collect.Range; import java.util.ArrayList; import java.util.BitSet; -import java.util.Collections; -import java.util.HashMap; import java.util.List; -import java.util.Map; import java.util.Map.Entry; import java.util.NavigableMap; import java.util.concurrent.ConcurrentSkipListMap; @@ -246,33 +243,24 @@ public Range lastRange() { } @Override - public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { + public int cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { NavigableMap subMap = rangeBitSetMap.subMap(lowerKey, true, upperKey, true); - if (subMap.isEmpty()) { - return Collections.emptyMap(); - } - Map v = new HashMap<>(); + MutableInt v = new MutableInt(0); subMap.forEach((ledgerId, bitset) -> { - boolean isLowerOrUpper = false; - BitSet lowerBitSet = null; - if (ledgerId == lowerKey) { - isLowerOrUpper = true; + if (ledgerId == lowerKey || ledgerId == upperKey) { BitSet temp = (BitSet) bitset.clone(); - temp.clear(0, (int) Math.max(0, lowerValue)); - lowerBitSet = temp; - v.put(ledgerId, temp.cardinality()); - } - if (ledgerId == upperKey) { - isLowerOrUpper = true; - BitSet temp = lowerBitSet == null ? (BitSet) bitset.clone() : lowerBitSet; - temp.clear((int) Math.min(upperValue + 1, temp.length()), temp.length()); - v.put(ledgerId, temp.cardinality()); - } - if (!isLowerOrUpper) { - v.put(ledgerId, bitset.cardinality()); + if (ledgerId == lowerKey) { + temp.clear(0, (int) Math.max(0, lowerValue)); + } + if (ledgerId == upperKey) { + temp.clear((int) Math.min(upperValue + 1, temp.length()), temp.length()); + } + v.add(temp.cardinality()); + } else { + v.add(bitset.cardinality()); } }); - return v; + return v.intValue(); } @Override 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 e8a96a5fe0695..579a45aa62916 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 @@ -25,7 +25,6 @@ import java.util.Collection; import java.util.Iterator; import java.util.List; -import java.util.Map; import java.util.NoSuchElementException; import java.util.Set; import lombok.EqualsAndHashCode; @@ -127,9 +126,9 @@ public interface LongPairRangeSet> { Range lastRange(); /** - * Return the number bit sets to true for each key from lower to upper. + * Return the number bit sets to true from lower to upper. */ - Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue); + int cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue); /** * Represents a function that accepts two long arguments and produces a result. @@ -303,7 +302,7 @@ public Range lastRange() { } @Override - public Map cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { + public int cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { throw new UnsupportedOperationException(); } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java index 4aa911df4533c..5d9af2e02271e 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/collections/ConcurrentOpenLongPairRangeSetTest.java @@ -25,7 +25,6 @@ import java.util.ArrayList; import java.util.List; -import java.util.Map; import java.util.Set; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPair; @@ -465,26 +464,20 @@ private List> getConnectedRange(Set> gRanges) { @Test public void testCardinality() { ConcurrentOpenLongPairRangeSet set = new ConcurrentOpenLongPairRangeSet<>(consumer); - Map v = set.cardinality(0, 0, Integer.MAX_VALUE, Integer.MAX_VALUE); - assertEquals(v.size(), 0 ); + int v = set.cardinality(0, 0, Integer.MAX_VALUE, Integer.MAX_VALUE); + assertEquals(v, 0 ); set.addOpenClosed(1, 0, 1, 20); set.addOpenClosed(1, 30, 1, 90); set.addOpenClosed(2, 0, 3, 30); v = set.cardinality(1, 0, 1, 100); - assertEquals(v.size(), 1); - assertEquals((int) v.get(1L), 80); + assertEquals(v, 80); v = set.cardinality(1, 11, 1, 100); - assertEquals(v.size(), 1); - assertEquals((int) v.get(1L), 70); + assertEquals(v, 70); v = set.cardinality(1, 0, 1, 90); - assertEquals(v.size(), 1); - assertEquals((int) v.get(1L), 80); + assertEquals(v, 80); v = set.cardinality(1, 0, 1, 80); - assertEquals(v.size(), 1); - assertEquals((int) v.get(1L), 70); + assertEquals(v, 70); v = set.cardinality(1, 0, 3, 30); - assertEquals(v.size(), 2); - assertEquals((int) v.get(1L), 80); - assertEquals((int) v.get(3L), 31); + assertEquals(v, 80 + 31); } } From 1d54b0dd2ed17031abd4ff50b6614c4b788288d5 Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 6 Sep 2022 12:02:59 +0800 Subject: [PATCH 5/5] Address lari's comment --- .../collections/ConcurrentOpenLongPairRangeSet.java | 10 ++++++---- .../common/util/collections/LongPairRangeSet.java | 2 +- 2 files changed, 7 insertions(+), 5 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 b13d384b1d6c6..a71c5ceb8de79 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 @@ -246,13 +246,15 @@ public Range lastRange() { public int cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue) { NavigableMap subMap = rangeBitSetMap.subMap(lowerKey, true, upperKey, true); MutableInt v = new MutableInt(0); - subMap.forEach((ledgerId, bitset) -> { - if (ledgerId == lowerKey || ledgerId == upperKey) { + subMap.forEach((key, bitset) -> { + if (key == lowerKey || key == upperKey) { BitSet temp = (BitSet) bitset.clone(); - if (ledgerId == lowerKey) { + // Trim the bitset index which < lowerValue + if (key == lowerKey) { temp.clear(0, (int) Math.max(0, lowerValue)); } - if (ledgerId == upperKey) { + // Trim the bitset index which > upperValue + if (key == upperKey) { temp.clear((int) Math.min(upperValue + 1, temp.length()), temp.length()); } v.add(temp.cardinality()); 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 579a45aa62916..d804900ed420b 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 @@ -126,7 +126,7 @@ public interface LongPairRangeSet> { Range lastRange(); /** - * Return the number bit sets to true from lower to upper. + * Return the number bit sets to true from lower (inclusive) to upper (inclusive). */ int cardinality(long lowerKey, long lowerValue, long upperKey, long upperValue);