Skip to content

Commit 30b06a0

Browse files
committed
[fix][broker] Remove lock contention in delayed delivery stats read paths
1 parent 6d0d99b commit 30b06a0

6 files changed

Lines changed: 99 additions & 73 deletions

File tree

pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,9 @@ public class InMemoryDelayedDeliveryTracker extends AbstractDelayedDeliveryTrack
7171
// Count of delayed messages in the tracker.
7272
private final AtomicLong delayedMessagesCount = new AtomicLong(0);
7373

74+
// Cached memory usage of the delayed message bitmaps, maintained via delta on each mutation.
75+
private final AtomicLong memoryUsage = new AtomicLong(0);
76+
7477
InMemoryDelayedDeliveryTracker(AbstractPersistentDispatcherMultipleConsumers dispatcher, Timer timer,
7578
long tickTimeMillis,
7679
boolean isDelayedDeliveryDeliverAtTimeStrict,
@@ -148,7 +151,9 @@ public boolean addMessage(long ledgerId, long entryId, long deliverAt) {
148151
boolean isNew = !bitmap.contains(entryId);
149152

150153
if (isNew) {
154+
long oldSize = bitmap.getLongSizeInBytes();
151155
bitmap.addLong(entryId);
156+
memoryUsage.addAndGet(bitmap.getLongSizeInBytes() - oldSize);
152157
delayedMessagesCount.incrementAndGet();
153158
}
154159

@@ -233,14 +238,17 @@ public NavigableSet<Position> getScheduledMessages(int maxMessages) {
233238
});
234239
n -= cardinalityInt;
235240
delayedMessagesCount.addAndGet(-cardinalityInt);
241+
memoryUsage.addAndGet(-entryIds.getLongSizeInBytes());
236242
ledgerIdToDelete.add(ledgerId);
237243
} else {
244+
long oldSize = entryIds.getLongSizeInBytes();
238245
Roaring64Bitmap entryIdsToRemove = new Roaring64Bitmap();
239246
entryIds.stream().limit(n).forEach(entryId -> {
240247
positions.add(PositionFactory.create(ledgerId, entryId));
241248
entryIdsToRemove.addLong(entryId);
242249
});
243250
entryIds.andNot(entryIdsToRemove);
251+
memoryUsage.addAndGet(entryIds.getLongSizeInBytes() - oldSize);
244252
delayedMessagesCount.addAndGet(-n);
245253
n = 0;
246254
}
@@ -277,6 +285,7 @@ public NavigableSet<Position> getScheduledMessages(int maxMessages) {
277285
public CompletableFuture<Void> clear() {
278286
this.delayedMessageMap.clear();
279287
this.delayedMessagesCount.set(0);
288+
this.memoryUsage.set(0);
280289
return CompletableFuture.completedFuture(null);
281290
}
282291

@@ -293,9 +302,7 @@ public long getNumberOfDelayedMessages() {
293302
*/
294303
@Override
295304
public long getBufferMemoryUsage() {
296-
return delayedMessageMap.values().stream().mapToLong(
297-
ledgerMap -> ledgerMap.values().stream().mapToLong(
298-
Roaring64Bitmap::getLongSizeInBytes).sum()).sum();
305+
return memoryUsage.get();
299306
}
300307

301308
@Override

pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java

Lines changed: 64 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,6 @@
3232
import java.util.ArrayList;
3333
import java.util.Collections;
3434
import java.util.HashMap;
35-
import java.util.Iterator;
3635
import java.util.List;
3736
import java.util.Map;
3837
import java.util.NavigableSet;
@@ -111,6 +110,10 @@ public static record SnapshotKey(long ledgerId, long entryId) {}
111110
@VisibleForTesting
112111
private final RangeMap<Long, ImmutableBucket> immutableBuckets;
113112

113+
private final AtomicLong bucketsCount = new AtomicLong(0);
114+
115+
private final AtomicLong totalSnapshotLengthBytes = new AtomicLong(0);
116+
114117
private final ConcurrentHashMap<SnapshotKey, ImmutableBucket> snapshotSegmentLastIndexMap;
115118

116119
private final BucketDelayedMessageIndexStats stats;
@@ -245,15 +248,19 @@ private synchronized long recoverBucketSnapshot() throws RecoverDelayedDeliveryT
245248
for (Map.Entry<Range<Long>, ImmutableBucket> mapEntry : toBeDeletedBucketMap.entrySet()) {
246249
Range<Long> key = mapEntry.getKey();
247250
ImmutableBucket immutableBucket = mapEntry.getValue();
248-
immutableBucketMap.remove(key);
251+
removeBucket(key);
249252
// delete asynchronously without waiting for completion
250253
immutableBucket.asyncDeleteBucketSnapshot(stats);
251254
}
252255

253256
MutableLong numberDelayedMessages = new MutableLong(0);
254-
immutableBucketMap.values().forEach(bucket -> {
257+
long totalLength = 0;
258+
for (ImmutableBucket bucket : immutableBucketMap.values()) {
255259
numberDelayedMessages.add(bucket.numberBucketDelayedMessages);
256-
});
260+
totalLength += bucket.getSnapshotLength();
261+
}
262+
totalSnapshotLengthBytes.set(totalLength);
263+
bucketsCount.set(immutableBuckets.asMapOfRanges().size());
257264

258265
log.info()
259266
.attr("buckets", immutableBucketMap.size())
@@ -292,10 +299,10 @@ private CompletableFuture<List<DelayedIndex>> handleRecoverBucketSnapshotEntry(I
292299

293300
private synchronized void putAndCleanOverlapRange(Range<Long> range, ImmutableBucket immutableBucket,
294301
Map<Range<Long>, ImmutableBucket> toBeDeletedBucketMap) {
295-
RangeMap<Long, ImmutableBucket> subRangeMap = immutableBuckets.subRangeMap(range);
302+
Map<Range<Long>, ImmutableBucket> subRangeMap = immutableBuckets.subRangeMap(range).asMapOfRanges();
296303
boolean canPut = false;
297-
if (!subRangeMap.asMapOfRanges().isEmpty()) {
298-
for (Map.Entry<Range<Long>, ImmutableBucket> rangeEntry : subRangeMap.asMapOfRanges().entrySet()) {
304+
if (!subRangeMap.isEmpty()) {
305+
for (Map.Entry<Range<Long>, ImmutableBucket> rangeEntry : subRangeMap.entrySet()) {
299306
if (range.encloses(rangeEntry.getKey())) {
300307
toBeDeletedBucketMap.put(rangeEntry.getKey(), rangeEntry.getValue());
301308
canPut = true;
@@ -306,7 +313,7 @@ private synchronized void putAndCleanOverlapRange(Range<Long> range, ImmutableBu
306313
}
307314

308315
if (canPut) {
309-
immutableBuckets.put(range, immutableBucket);
316+
putBucket(range, immutableBucket);
310317
}
311318
}
312319

@@ -333,7 +340,7 @@ private void afterCreateImmutableBucket(Pair<ImmutableBucket, DelayedIndex> immu
333340
long startTime) {
334341
if (immutableBucketDelayedIndexPair != null) {
335342
ImmutableBucket immutableBucket = immutableBucketDelayedIndexPair.getLeft();
336-
immutableBuckets.put(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId),
343+
putBucket(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId),
337344
immutableBucket);
338345

339346
DelayedIndex lastDelayedIndex = immutableBucketDelayedIndexPair.getRight();
@@ -345,7 +352,12 @@ private void afterCreateImmutableBucket(Pair<ImmutableBucket, DelayedIndex> immu
345352
CompletableFuture<Long> future = createFuture.handle((bucketId, ex) -> {
346353
if (ex == null) {
347354
immutableBucket.setSnapshotSegments(null);
348-
immutableBucket.asyncUpdateSnapshotLength();
355+
immutableBucket.asyncUpdateSnapshotLength()
356+
.thenAccept(newLength -> {
357+
synchronized (BucketDelayedDeliveryTracker.this) {
358+
updateBucketSnapshotLength(immutableBucket, newLength);
359+
}
360+
});
349361
log.info()
350362
.attr("bucketKey", immutableBucket.bucketKey())
351363
.log("Create bucket snapshot finish, bucketKey");
@@ -375,7 +387,7 @@ private void afterCreateImmutableBucket(Pair<ImmutableBucket, DelayedIndex> immu
375387
});
376388

377389
immutableBucket.setCurrentSegmentEntryId(immutableBucket.lastSegmentEntryId);
378-
immutableBuckets.asMapOfRanges().remove(
390+
removeBucket(
379391
Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId));
380392
snapshotSegmentLastIndexMap.remove(
381393
new SnapshotKey(lastDelayedIndex.getLedgerId(), lastDelayedIndex.getEntryId()));
@@ -413,7 +425,7 @@ public synchronized boolean addMessage(long ledgerId, long entryId, long deliver
413425
afterCreateImmutableBucket(immutableBucketDelayedIndexPair, createStartTime);
414426
lastMutableBucket.resetLastMutableBucketRange();
415427

416-
if (maxNumBuckets > 0 && immutableBuckets.asMapOfRanges().size() > maxNumBuckets
428+
if (maxNumBuckets > 0 && bucketsCount.get() > maxNumBuckets
417429
&& (trimFuture == null || trimFuture.isDone())) {
418430
trimFuture = asyncTrimImmutableBuckets()
419431
.thenCompose(ignore -> asyncMergeBucketSnapshot())
@@ -483,11 +495,10 @@ private synchronized List<ImmutableBucket> selectMergedBuckets(final List<Immuta
483495
}
484496

485497
private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot() {
486-
List<ImmutableBucket> immutableBucketList = immutableBuckets.asMapOfRanges().values().stream().toList();
487-
if (maxNumBuckets <= 0 || immutableBucketList.size() <= maxNumBuckets) {
498+
if (maxNumBuckets <= 0 || bucketsCount.get() <= maxNumBuckets) {
488499
return CompletableFuture.completedFuture(null);
489500
}
490-
501+
List<ImmutableBucket> immutableBucketList = immutableBuckets.asMapOfRanges().values().stream().toList();
491502
List<ImmutableBucket> toBeMergeImmutableBuckets = selectMergedBuckets(immutableBucketList, MAX_MERGE_NUM);
492503

493504
if (toBeMergeImmutableBuckets.isEmpty()) {
@@ -522,7 +533,7 @@ private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot() {
522533
} else {
523534
log.info()
524535
.attr("bucketKeys", bucketsStr)
525-
.attr("bucketNum", immutableBuckets.asMapOfRanges().size())
536+
.attr("bucketNum", bucketsCount.get())
526537
.log("Merge bucket snapshot finish");
527538

528539
stats.recordSuccessEvent(BucketDelayedMessageIndexStats.Type.merge,
@@ -592,8 +603,7 @@ private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot(List<Immut
592603
});
593604

594605
for (ImmutableBucket bucket : buckets) {
595-
immutableBuckets.asMapOfRanges()
596-
.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
606+
removeBucket(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
597607
}
598608
}
599609
});
@@ -696,8 +706,7 @@ public synchronized NavigableSet<Position> getScheduledMessages(int maxMessages)
696706
synchronized (BucketDelayedDeliveryTracker.this) {
697707
this.snapshotSegmentLastIndexMap.remove(snapshotKey);
698708
if (CollectionUtils.isEmpty(indexList)) {
699-
immutableBuckets.asMapOfRanges()
700-
.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
709+
removeBucket(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
701710
bucket.asyncDeleteBucketSnapshot(stats);
702711
return;
703712
}
@@ -817,14 +826,16 @@ public CompletableFuture<Void> closeAsync() {
817826
}
818827

819828
private CompletableFuture<Void> cleanImmutableBuckets() {
829+
Map<Range<Long>, ImmutableBucket> bucketsToDelete =
830+
new HashMap<>(immutableBuckets.asMapOfRanges());
831+
820832
List<CompletableFuture<Void>> futures = new ArrayList<>();
821-
Iterator<ImmutableBucket> iterator = immutableBuckets.asMapOfRanges().values().iterator();
822-
while (iterator.hasNext()) {
823-
ImmutableBucket bucket = iterator.next();
824-
futures.add(bucket.clear(stats));
833+
bucketsToDelete.forEach((range, bucket) -> {
834+
removeBucket(range);
825835
numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages());
826-
iterator.remove();
827-
}
836+
futures.add(bucket.clear(stats));
837+
});
838+
828839
return FutureUtil.waitForAll(futures);
829840
}
830841

@@ -847,13 +858,9 @@ public synchronized boolean containsMessage(long ledgerId, long entryId) {
847858
}
848859

849860
public Map<String, TopicMetricBean> genTopicMetricMap() {
850-
stats.recordNumOfBuckets(immutableBuckets.asMapOfRanges().size() + 1);
861+
stats.recordNumOfBuckets((int) (bucketsCount.get() + 1));
851862
stats.recordDelayedMessageIndexLoaded(this.sharedBucketPriorityQueue.size() + this.lastMutableBucket.size());
852-
MutableLong totalSnapshotLength = new MutableLong();
853-
immutableBuckets.asMapOfRanges().values().forEach(immutableBucket -> {
854-
totalSnapshotLength.add(immutableBucket.getSnapshotLength());
855-
});
856-
stats.recordBucketSnapshotSizeBytes(totalSnapshotLength.longValue());
863+
stats.recordBucketSnapshotSizeBytes(totalSnapshotLengthBytes.get());
857864
return stats.genTopicMetricMap();
858865
}
859866

@@ -897,7 +904,7 @@ private CompletableFuture<Void> deleteBucketSnapshot(String ledgerName,
897904
}
898905
synchronized (this) {
899906
snapshotSegmentLastIndexMap.entrySet().removeIf(entry -> entry.getValue() == bucket);
900-
immutableBuckets.remove(range);
907+
removeBucket(range);
901908
numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages());
902909
}
903910
return null;
@@ -909,4 +916,28 @@ private Long firstActiveLedgerId() {
909916
Position mdp = cursor.getMarkDeletedPosition();
910917
return mdp == null ? null : mdp.getLedgerId();
911918
}
919+
920+
private void putBucket(Range<Long> range, ImmutableBucket bucket) {
921+
long removedLength = immutableBuckets.subRangeMap(range).asMapOfRanges().values().stream()
922+
.mapToLong(ImmutableBucket::getSnapshotLength)
923+
.sum();
924+
925+
immutableBuckets.put(range, bucket);
926+
bucketsCount.set(immutableBuckets.asMapOfRanges().size());
927+
totalSnapshotLengthBytes.addAndGet(bucket.getSnapshotLength() - removedLength);
928+
}
929+
930+
private void removeBucket(Range<Long> range) {
931+
ImmutableBucket removed = immutableBuckets.asMapOfRanges().remove(range);
932+
if (removed != null) {
933+
bucketsCount.set(immutableBuckets.asMapOfRanges().size());
934+
totalSnapshotLengthBytes.addAndGet(-removed.getSnapshotLength());
935+
}
936+
}
937+
938+
private void updateBucketSnapshotLength(ImmutableBucket bucket, long newLength) {
939+
long oldLength = bucket.getSnapshotLength();
940+
bucket.setSnapshotLength(newLength);
941+
totalSnapshotLengthBytes.addAndGet(newLength - oldLength);
942+
}
912943
}

pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedMessageIndexStats.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ enum Type {
5959
public BucketDelayedMessageIndexStats() {
6060
}
6161

62-
public Map<String, TopicMetricBean> genTopicMetricMap() {
62+
public synchronized Map<String, TopicMetricBean> genTopicMetricMap() {
6363
Map<String, TopicMetricBean> metrics = new HashMap<>();
6464

6565
metrics.put(BUCKET_TOTAL_NAME,

pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/ImmutableBucket.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -130,19 +130,19 @@ private CompletableFuture<List<DelayedIndex>> asyncLoadNextBucketSnapshotEntry(b
130130
.log("Failed to get bucket snapshot segment");
131131
}
132132
}), BucketSnapshotPersistenceException.class, MaxRetryTimes)
133-
.thenApply(bucketSnapshotSegments -> {
133+
.thenCompose(bucketSnapshotSegments -> {
134134
if (CollectionUtils.isEmpty(bucketSnapshotSegments)) {
135-
return Collections.emptyList();
135+
return CompletableFuture.completedFuture(Collections.emptyList());
136136
}
137137

138138
SnapshotSegment snapshotSegment =
139139
bucketSnapshotSegments.get(0);
140140
List<DelayedIndex> indexList = snapshotSegment.getIndexesList();
141141
this.setCurrentSegmentEntryId(nextSegmentEntryId);
142142
if (isRecover) {
143-
this.asyncUpdateSnapshotLength();
143+
return this.asyncUpdateSnapshotLength().thenApply(__ -> indexList);
144144
}
145-
return indexList;
145+
return CompletableFuture.completedFuture(indexList);
146146
});
147147
});
148148
}

pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java

Lines changed: 10 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,7 @@ public class PersistentDispatcherMultipleConsumers extends AbstractPersistentDis
9494
protected final MessageRedeliveryController redeliveryMessages;
9595
protected final RedeliveryTracker redeliveryTracker;
9696

97-
private Optional<DelayedDeliveryTracker> delayedDeliveryTracker = Optional.empty();
97+
private volatile Optional<DelayedDeliveryTracker> delayedDeliveryTracker = Optional.empty();
9898

9999
protected volatile boolean havePendingRead = false;
100100
protected volatile boolean havePendingReplayRead = false;
@@ -1374,13 +1374,12 @@ protected boolean isNormalReadAllowed() {
13741374
}
13751375

13761376

1377-
1378-
protected synchronized boolean shouldPauseDeliveryForDelayTracker() {
1379-
return delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().shouldPauseAllDeliveries();
1377+
protected boolean shouldPauseDeliveryForDelayTracker() {
1378+
return delayedDeliveryTracker.map(DelayedDeliveryTracker::shouldPauseAllDeliveries).orElse(false);
13801379
}
13811380

13821381
@Override
1383-
public synchronized long getNumberOfDelayedMessages() {
1382+
public long getNumberOfDelayedMessages() {
13841383
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getNumberOfDelayedMessages).orElse(0L);
13851384
}
13861385

@@ -1466,20 +1465,15 @@ public PersistentTopic getTopic() {
14661465
}
14671466

14681467

1469-
public synchronized long getDelayedTrackerMemoryUsage() {
1468+
public long getDelayedTrackerMemoryUsage() {
14701469
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getBufferMemoryUsage).orElse(0L);
14711470
}
14721471

1473-
public synchronized Map<String, TopicMetricBean> getBucketDelayedIndexStats() {
1474-
if (delayedDeliveryTracker.isEmpty()) {
1475-
return Collections.emptyMap();
1476-
}
1477-
1478-
if (delayedDeliveryTracker.get() instanceof BucketDelayedDeliveryTracker) {
1479-
return ((BucketDelayedDeliveryTracker) delayedDeliveryTracker.get()).genTopicMetricMap();
1480-
}
1481-
1482-
return Collections.emptyMap();
1472+
public Map<String, TopicMetricBean> getBucketDelayedIndexStats() {
1473+
return delayedDeliveryTracker
1474+
.filter(BucketDelayedDeliveryTracker.class::isInstance)
1475+
.map(tracker -> ((BucketDelayedDeliveryTracker) tracker).genTopicMetricMap())
1476+
.orElse(Collections.emptyMap());
14831477
}
14841478

14851479
@Override

0 commit comments

Comments
 (0)