Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,9 @@ public class InMemoryDelayedDeliveryTracker extends AbstractDelayedDeliveryTrack
// Count of delayed messages in the tracker.
private final AtomicLong delayedMessagesCount = new AtomicLong(0);

// Cached memory usage of the delayed message bitmaps, maintained via delta on each mutation.
private final AtomicLong memoryUsage = new AtomicLong(0);

InMemoryDelayedDeliveryTracker(AbstractPersistentDispatcherMultipleConsumers dispatcher, Timer timer,
long tickTimeMillis,
boolean isDelayedDeliveryDeliverAtTimeStrict,
Expand Down Expand Up @@ -144,7 +147,11 @@ public boolean addMessage(long ledgerId, long entryId, long deliverAt) {

LongBitmap bitmap = delayedMessageMap.computeIfAbsent(timestamp, k -> new Long2ObjectRBTreeMap<>())
.computeIfAbsent(ledgerId, k -> LongBitmaps.create());

long oldSize = bitmap.serializedSize();
if (bitmap.checkedAdd(entryId)) {
long newSize = bitmap.serializedSize();
memoryUsage.addAndGet(newSize - oldSize);
delayedMessagesCount.incrementAndGet();
}

Expand Down Expand Up @@ -222,9 +229,12 @@ public NavigableSet<Position> getScheduledMessages(int maxMessages) {
long ledgerId = ledgerEntry.getLongKey();
LongBitmap entryIds = ledgerEntry.getValue();
long cardinality = entryIds.cardinality();
long oldSize = entryIds.serializedSize();
long drained = entryIds.drainTo(n, entryId -> {
positions.add(PositionFactory.create(ledgerId, entryId));
});
long newSize = entryIds.serializedSize();
memoryUsage.addAndGet(newSize - oldSize);
delayedMessagesCount.addAndGet(-drained);
n -= drained;
if (drained == cardinality) {
Expand Down Expand Up @@ -264,6 +274,7 @@ public NavigableSet<Position> getScheduledMessages(int maxMessages) {
public CompletableFuture<Void> clear() {
this.delayedMessageMap.clear();
this.delayedMessagesCount.set(0);
this.memoryUsage.set(0);
return CompletableFuture.completedFuture(null);
}

Expand All @@ -279,8 +290,7 @@ public long getNumberOfDelayedMessages() {
*/
@Override
public long getBufferMemoryUsage() {
return delayedMessageMap.values().stream().mapToLong(
ledgerMap -> ledgerMap.values().stream().mapToLong(LongBitmap::serializedSize).sum()).sum();
return memoryUsage.get();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.NavigableSet;
Expand Down Expand Up @@ -111,6 +110,14 @@ public static record SnapshotKey(long ledgerId, long entryId) {}
@VisibleForTesting
private final RangeMap<Long, ImmutableBucket> immutableBuckets;

@Getter
@VisibleForTesting
private final AtomicLong bucketsCount = new AtomicLong(0);

@Getter
@VisibleForTesting
private final AtomicLong totalSnapshotLengthBytes = new AtomicLong(0);

private final ConcurrentHashMap<SnapshotKey, ImmutableBucket> snapshotSegmentLastIndexMap;

private final BucketDelayedMessageIndexStats stats;
Expand Down Expand Up @@ -245,15 +252,19 @@ private synchronized long recoverBucketSnapshot() throws RecoverDelayedDeliveryT
for (Map.Entry<Range<Long>, ImmutableBucket> mapEntry : toBeDeletedBucketMap.entrySet()) {
Range<Long> key = mapEntry.getKey();
ImmutableBucket immutableBucket = mapEntry.getValue();
immutableBucketMap.remove(key);
removeBucket(key);
// delete asynchronously without waiting for completion
immutableBucket.asyncDeleteBucketSnapshot(stats);
}

MutableLong numberDelayedMessages = new MutableLong(0);
immutableBucketMap.values().forEach(bucket -> {
long totalLength = 0;
for (ImmutableBucket bucket : immutableBucketMap.values()) {
numberDelayedMessages.add(bucket.numberBucketDelayedMessages);
});
totalLength += bucket.getSnapshotLength();
}
totalSnapshotLengthBytes.set(totalLength);
bucketsCount.set(immutableBuckets.asMapOfRanges().size());

log.info()
.attr("buckets", immutableBucketMap.size())
Expand Down Expand Up @@ -292,12 +303,16 @@ private CompletableFuture<List<DelayedIndex>> handleRecoverBucketSnapshotEntry(I

private synchronized void putAndCleanOverlapRange(Range<Long> range, ImmutableBucket immutableBucket,
Map<Range<Long>, ImmutableBucket> toBeDeletedBucketMap) {
RangeMap<Long, ImmutableBucket> subRangeMap = immutableBuckets.subRangeMap(range);
Map<Range<Long>, ImmutableBucket> subRangeMap = immutableBuckets.subRangeMap(range).asMapOfRanges();
boolean canPut = false;
if (!subRangeMap.asMapOfRanges().isEmpty()) {
for (Map.Entry<Range<Long>, ImmutableBucket> rangeEntry : subRangeMap.asMapOfRanges().entrySet()) {
if (range.encloses(rangeEntry.getKey())) {
toBeDeletedBucketMap.put(rangeEntry.getKey(), rangeEntry.getValue());
if (!subRangeMap.isEmpty()) {
for (Map.Entry<Range<Long>, ImmutableBucket> rangeEntry : subRangeMap.entrySet()) {
// Use original key instead of truncated key for encloses check
ImmutableBucket bucket = rangeEntry.getValue();
Range<Long> originalKey = Range.closed(bucket.startLedgerId, bucket.endLedgerId);

if (range.encloses(originalKey)) {
toBeDeletedBucketMap.put(originalKey, bucket);
canPut = true;
}
}
Expand All @@ -306,7 +321,7 @@ private synchronized void putAndCleanOverlapRange(Range<Long> range, ImmutableBu
}

if (canPut) {
immutableBuckets.put(range, immutableBucket);
putBucket(range, immutableBucket);
}
}

Expand All @@ -333,7 +348,7 @@ private void afterCreateImmutableBucket(Pair<ImmutableBucket, DelayedIndex> immu
long startTime) {
if (immutableBucketDelayedIndexPair != null) {
ImmutableBucket immutableBucket = immutableBucketDelayedIndexPair.getLeft();
immutableBuckets.put(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId),
putBucket(Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId),
immutableBucket);

DelayedIndex lastDelayedIndex = immutableBucketDelayedIndexPair.getRight();
Expand All @@ -345,7 +360,12 @@ private void afterCreateImmutableBucket(Pair<ImmutableBucket, DelayedIndex> immu
CompletableFuture<Long> future = createFuture.handle((bucketId, ex) -> {
if (ex == null) {
immutableBucket.setSnapshotSegments(null);
immutableBucket.asyncUpdateSnapshotLength();
immutableBucket.asyncUpdateSnapshotLength()
.thenAccept(newLength -> {
synchronized (BucketDelayedDeliveryTracker.this) {
updateBucketSnapshotLength(immutableBucket, newLength);
}
});
log.info()
.attr("bucketKey", immutableBucket.bucketKey())
.log("Create bucket snapshot finish, bucketKey");
Expand Down Expand Up @@ -375,7 +395,7 @@ private void afterCreateImmutableBucket(Pair<ImmutableBucket, DelayedIndex> immu
});

immutableBucket.setCurrentSegmentEntryId(immutableBucket.lastSegmentEntryId);
immutableBuckets.asMapOfRanges().remove(
removeBucket(
Range.closed(immutableBucket.startLedgerId, immutableBucket.endLedgerId));
snapshotSegmentLastIndexMap.remove(
new SnapshotKey(lastDelayedIndex.getLedgerId(), lastDelayedIndex.getEntryId()));
Expand Down Expand Up @@ -413,7 +433,7 @@ public synchronized boolean addMessage(long ledgerId, long entryId, long deliver
afterCreateImmutableBucket(immutableBucketDelayedIndexPair, createStartTime);
lastMutableBucket.resetLastMutableBucketRange();

if (maxNumBuckets > 0 && immutableBuckets.asMapOfRanges().size() > maxNumBuckets
if (maxNumBuckets > 0 && bucketsCount.get() > maxNumBuckets
&& (trimFuture == null || trimFuture.isDone())) {
trimFuture = asyncTrimImmutableBuckets()
.thenCompose(ignore -> asyncMergeBucketSnapshot())
Expand Down Expand Up @@ -483,11 +503,10 @@ private synchronized List<ImmutableBucket> selectMergedBuckets(final List<Immuta
}

private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot() {
List<ImmutableBucket> immutableBucketList = immutableBuckets.asMapOfRanges().values().stream().toList();
if (maxNumBuckets <= 0 || immutableBucketList.size() <= maxNumBuckets) {
if (maxNumBuckets <= 0 || bucketsCount.get() <= maxNumBuckets) {
return CompletableFuture.completedFuture(null);
}

List<ImmutableBucket> immutableBucketList = immutableBuckets.asMapOfRanges().values().stream().toList();
List<ImmutableBucket> toBeMergeImmutableBuckets = selectMergedBuckets(immutableBucketList, MAX_MERGE_NUM);

if (toBeMergeImmutableBuckets.isEmpty()) {
Expand Down Expand Up @@ -522,7 +541,7 @@ private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot() {
} else {
log.info()
.attr("bucketKeys", bucketsStr)
.attr("bucketNum", immutableBuckets.asMapOfRanges().size())
.attr("bucketNum", bucketsCount.get())
.log("Merge bucket snapshot finish");

stats.recordSuccessEvent(BucketDelayedMessageIndexStats.Type.merge,
Expand Down Expand Up @@ -590,8 +609,7 @@ private synchronized CompletableFuture<Void> asyncMergeBucketSnapshot(List<Immut
});

for (ImmutableBucket bucket : buckets) {
immutableBuckets.asMapOfRanges()
.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
removeBucket(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
}
}
});
Expand Down Expand Up @@ -699,8 +717,7 @@ public synchronized NavigableSet<Position> getScheduledMessages(int maxMessages)
synchronized (BucketDelayedDeliveryTracker.this) {
this.snapshotSegmentLastIndexMap.remove(snapshotKey);
if (CollectionUtils.isEmpty(indexList)) {
immutableBuckets.asMapOfRanges()
.remove(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
removeBucket(Range.closed(bucket.startLedgerId, bucket.endLedgerId));
bucket.asyncDeleteBucketSnapshot(stats);
return;
}
Expand Down Expand Up @@ -820,14 +837,16 @@ public CompletableFuture<Void> closeAsync() {
}

private CompletableFuture<Void> cleanImmutableBuckets() {
Map<Range<Long>, ImmutableBucket> bucketsToDelete =
new HashMap<>(immutableBuckets.asMapOfRanges());

List<CompletableFuture<Void>> futures = new ArrayList<>();
Iterator<ImmutableBucket> iterator = immutableBuckets.asMapOfRanges().values().iterator();
while (iterator.hasNext()) {
ImmutableBucket bucket = iterator.next();
futures.add(bucket.clear(stats));
bucketsToDelete.forEach((range, bucket) -> {
removeBucket(range);
numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages());
iterator.remove();
}
futures.add(bucket.clear(stats));
});

return FutureUtil.waitForAll(futures);
}

Expand All @@ -850,13 +869,9 @@ public synchronized boolean containsMessage(long ledgerId, long entryId) {
}

public Map<String, TopicMetricBean> genTopicMetricMap() {
stats.recordNumOfBuckets(immutableBuckets.asMapOfRanges().size() + 1);
stats.recordNumOfBuckets((int) (bucketsCount.get() + 1));
stats.recordDelayedMessageIndexLoaded(this.sharedBucketPriorityQueue.size() + this.lastMutableBucket.size());
Comment thread
nodece marked this conversation as resolved.
Comment thread
nodece marked this conversation as resolved.
MutableLong totalSnapshotLength = new MutableLong();
immutableBuckets.asMapOfRanges().values().forEach(immutableBucket -> {
totalSnapshotLength.add(immutableBucket.getSnapshotLength());
});
stats.recordBucketSnapshotSizeBytes(totalSnapshotLength.longValue());
stats.recordBucketSnapshotSizeBytes(totalSnapshotLengthBytes.get());
return stats.genTopicMetricMap();
}

Expand Down Expand Up @@ -900,7 +915,7 @@ private CompletableFuture<Void> deleteBucketSnapshot(String ledgerName,
}
synchronized (this) {
snapshotSegmentLastIndexMap.entrySet().removeIf(entry -> entry.getValue() == bucket);
immutableBuckets.remove(range);
removeBucket(range);
Comment thread
nodece marked this conversation as resolved.
numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages());
}
return null;
Expand All @@ -912,4 +927,35 @@ private Long firstActiveLedgerId() {
Position mdp = cursor.getMarkDeletedPosition();
return mdp == null ? null : mdp.getLedgerId();
}

private void putBucket(Range<Long> range, ImmutableBucket bucket) {
long removedLength = immutableBuckets.subRangeMap(range).asMapOfRanges().values().stream()
.mapToLong(ImmutableBucket::getSnapshotLength)
.sum();

immutableBuckets.put(range, bucket);
bucketsCount.set(immutableBuckets.asMapOfRanges().size());
totalSnapshotLengthBytes.addAndGet(bucket.getSnapshotLength() - removedLength);
}

private void removeBucket(Range<Long> range) {
// Use exact key matching - all callers should provide exact keys
ImmutableBucket bucket = immutableBuckets.asMapOfRanges().get(range);

if (bucket != null) {
// Remove even if snapshot length is 0 (for newly created buckets)
immutableBuckets.asMapOfRanges().remove(range);
bucketsCount.set(immutableBuckets.asMapOfRanges().size());
totalSnapshotLengthBytes.addAndGet(-bucket.getSnapshotLength());
}
}

private void updateBucketSnapshotLength(ImmutableBucket bucket, long newLength) {
if (!immutableBuckets.asMapOfRanges().containsValue(bucket)) {
return;
}
long oldLength = bucket.getSnapshotLength();
bucket.setSnapshotLength(newLength);
totalSnapshotLengthBytes.addAndGet(newLength - oldLength);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ enum Type {
public BucketDelayedMessageIndexStats() {
}

public Map<String, TopicMetricBean> genTopicMetricMap() {
public synchronized Map<String, TopicMetricBean> genTopicMetricMap() {
Map<String, TopicMetricBean> metrics = new HashMap<>();

metrics.put(BUCKET_TOTAL_NAME,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,19 +131,21 @@ private CompletableFuture<List<DelayedIndex>> asyncLoadNextBucketSnapshotEntry(b
.log("Failed to get bucket snapshot segment");
}
}), BucketSnapshotPersistenceException.class, MaxRetryTimes)
.thenApply(bucketSnapshotSegments -> {
.thenCompose(bucketSnapshotSegments -> {
if (CollectionUtils.isEmpty(bucketSnapshotSegments)) {
return Collections.emptyList();
return CompletableFuture.completedFuture(Collections.emptyList());
}

SnapshotSegment snapshotSegment =
bucketSnapshotSegments.get(0);
List<DelayedIndex> indexList = snapshotSegment.getIndexesList();
this.setCurrentSegmentEntryId(nextSegmentEntryId);
if (isRecover) {
this.asyncUpdateSnapshotLength();
return this.asyncUpdateSnapshotLength()
.thenAccept(this::setSnapshotLength)
.thenApply(__ -> indexList);
}
return indexList;
return CompletableFuture.completedFuture(indexList);
});
});
}
Expand Down Expand Up @@ -245,8 +247,6 @@ protected CompletableFuture<Long> asyncUpdateSnapshotLength() {
.attr("bucketKey", bucketKey())
.exception(ex)
.log("Failed to get snapshot length");
} else {
setSnapshotLength(length);
}
});
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ public class PersistentDispatcherMultipleConsumers extends AbstractPersistentDis
protected final MessageRedeliveryController redeliveryMessages;
protected final RedeliveryTracker redeliveryTracker;

private Optional<DelayedDeliveryTracker> delayedDeliveryTracker = Optional.empty();
private volatile Optional<DelayedDeliveryTracker> delayedDeliveryTracker = Optional.empty();

protected volatile boolean havePendingRead = false;
protected volatile boolean havePendingReplayRead = false;
Expand Down Expand Up @@ -1374,13 +1374,12 @@ protected boolean isNormalReadAllowed() {
}



protected synchronized boolean shouldPauseDeliveryForDelayTracker() {
return delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().shouldPauseAllDeliveries();
protected boolean shouldPauseDeliveryForDelayTracker() {
return delayedDeliveryTracker.map(DelayedDeliveryTracker::shouldPauseAllDeliveries).orElse(false);
}

@Override
public synchronized long getNumberOfDelayedMessages() {
public long getNumberOfDelayedMessages() {
Comment thread
nodece marked this conversation as resolved.
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getNumberOfDelayedMessages).orElse(0L);
}

Expand Down Expand Up @@ -1466,20 +1465,15 @@ public PersistentTopic getTopic() {
}


public synchronized long getDelayedTrackerMemoryUsage() {
public long getDelayedTrackerMemoryUsage() {
Comment thread
nodece marked this conversation as resolved.
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getBufferMemoryUsage).orElse(0L);
}

public synchronized Map<String, TopicMetricBean> getBucketDelayedIndexStats() {
if (delayedDeliveryTracker.isEmpty()) {
return Collections.emptyMap();
}

if (delayedDeliveryTracker.get() instanceof BucketDelayedDeliveryTracker) {
return ((BucketDelayedDeliveryTracker) delayedDeliveryTracker.get()).genTopicMetricMap();
}

return Collections.emptyMap();
public Map<String, TopicMetricBean> getBucketDelayedIndexStats() {
return delayedDeliveryTracker
.filter(BucketDelayedDeliveryTracker.class::isInstance)
.map(tracker -> ((BucketDelayedDeliveryTracker) tracker).genTopicMetricMap())
Comment thread
nodece marked this conversation as resolved.
.orElse(Collections.emptyMap());
}

@Override
Expand Down
Loading
Loading