diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java index 489478df56236..452e6b13554a6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/Bucket.java @@ -81,20 +81,20 @@ boolean containsMessage(long ledgerId, long entryId) { if (bitSet == null) { return false; } - return bitSet.contains(entryId, entryId + 1); + return bitSet.contains(entryId); } void putIndexBit(long ledgerId, long entryId) { - delayedIndexBitMap.computeIfAbsent(ledgerId, k -> LongBitmaps.create()).add(entryId, entryId + 1); + delayedIndexBitMap.computeIfAbsent(ledgerId, k -> LongBitmaps.create()).add(entryId); } boolean removeIndexBit(long ledgerId, long entryId) { - boolean contained = false; LongBitmap bitSet = delayedIndexBitMap.get(ledgerId); - if (bitSet != null && bitSet.contains(entryId, entryId + 1)) { - contained = true; - bitSet.remove(entryId, entryId + 1); - + if (bitSet == null) { + return false; + } + boolean contained = bitSet.remove(entryId); + if (contained) { if (bitSet.isEmpty()) { delayedIndexBitMap.remove(ledgerId); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java index a91016bba0b56..ff0d4f5f19859 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/MutableBucket.java @@ -105,7 +105,7 @@ Pair createImmutableBucketAndAsyncPersistent( sharedQueue.add(timestamp, ledgerId, entryId); } - bitMap.computeIfAbsent(ledgerId, k -> LongBitmaps.create()).add(entryId, entryId + 1); + bitMap.computeIfAbsent(ledgerId, k -> LongBitmaps.create()).add(entryId); numMessages++; diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentRoaringBitmap.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentRoaringBitmap.java index 774227ed2b351..96ed33c2f8047 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentRoaringBitmap.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/ConcurrentRoaringBitmap.java @@ -116,32 +116,38 @@ public void add(long from, long to) { } @Override - public void remove(long value) { + public boolean remove(long value) { validateRange(value); long stamp = lock.writeLock(); try { - if (bitmap.checkedRemove((int) value)) { + boolean removed = bitmap.checkedRemove((int) value); + if (removed) { removesSinceTrim++; maybeTrim(); } + return removed; } finally { lock.unlockWrite(stamp); } } @Override - public void remove(long from, long to) { + public boolean remove(long from, long to) { if (to <= from) { - return; + return false; } validateRange(from); validateRange(to - 1); long stamp = lock.writeLock(); try { + long cardinalityBefore = bitmap.getLongCardinality(); bitmap.remove(from, to); - // Range size upper-bounds removals; clamp so a huge range can't overflow the counter. - removesSinceTrim = Math.min(removesSinceTrim + (to - from), TRIM_AFTER_REMOVES); - maybeTrim(); + long removedCount = cardinalityBefore - bitmap.getLongCardinality(); + if (removedCount > 0) { + removesSinceTrim += removedCount; + maybeTrim(); + } + return removedCount > 0; } finally { lock.unlockWrite(stamp); } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongBitmap.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongBitmap.java index 53432eae603dc..e08f05f1bc596 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongBitmap.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/collections/LongBitmap.java @@ -83,8 +83,9 @@ public interface LongBitmap { * * @param value value to remove * @throws IllegalArgumentException if value is outside the supported range + * @return {@code true} if removed, otherwise {@code false} */ - void remove(long value); + boolean remove(long value); /** * Removes all values in the half-open range {@code [from, to)}. @@ -94,8 +95,9 @@ public interface LongBitmap { * @param from inclusive lower bound * @param to exclusive upper bound * @throws IllegalArgumentException if the range exceeds the supported value range + * @return {@code true} if removed, otherwise {@code false} */ - void remove(long from, long to); + boolean remove(long from, long to); /** * Returns whether the bitmap contains the given value.