Skip to content

Commit aa7e718

Browse files
committed
[fix][broker] remove lock contention in delayed delivery stats read paths
1 parent 94de2b6 commit aa7e718

2 files changed

Lines changed: 20 additions & 31 deletions

File tree

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;
@@ -1372,13 +1372,12 @@ protected boolean isNormalReadAllowed() {
13721372
}
13731373

13741374

1375-
1376-
protected synchronized boolean shouldPauseDeliveryForDelayTracker() {
1377-
return delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().shouldPauseAllDeliveries();
1375+
protected boolean shouldPauseDeliveryForDelayTracker() {
1376+
return delayedDeliveryTracker.map(DelayedDeliveryTracker::shouldPauseAllDeliveries).orElse(false);
13781377
}
13791378

13801379
@Override
1381-
public synchronized long getNumberOfDelayedMessages() {
1380+
public long getNumberOfDelayedMessages() {
13821381
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getNumberOfDelayedMessages).orElse(0L);
13831382
}
13841383

@@ -1464,20 +1463,15 @@ public PersistentTopic getTopic() {
14641463
}
14651464

14661465

1467-
public synchronized long getDelayedTrackerMemoryUsage() {
1466+
public long getDelayedTrackerMemoryUsage() {
14681467
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getBufferMemoryUsage).orElse(0L);
14691468
}
14701469

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

14831477
@Override

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

Lines changed: 10 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -96,7 +96,7 @@ public class PersistentDispatcherMultipleConsumersClassic extends AbstractPersis
9696
protected final MessageRedeliveryController redeliveryMessages;
9797
protected final RedeliveryTracker redeliveryTracker;
9898

99-
private Optional<DelayedDeliveryTracker> delayedDeliveryTracker = Optional.empty();
99+
private volatile Optional<DelayedDeliveryTracker> delayedDeliveryTracker = Optional.empty();
100100

101101
protected volatile boolean havePendingRead = false;
102102
protected volatile boolean havePendingReplayRead = false;
@@ -1207,12 +1207,12 @@ protected boolean hasConsumersNeededNormalRead() {
12071207
return true;
12081208
}
12091209

1210-
protected synchronized boolean shouldPauseDeliveryForDelayTracker() {
1211-
return delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().shouldPauseAllDeliveries();
1210+
protected boolean shouldPauseDeliveryForDelayTracker() {
1211+
return delayedDeliveryTracker.map(DelayedDeliveryTracker::shouldPauseAllDeliveries).orElse(false);
12121212
}
12131213

12141214
@Override
1215-
public synchronized long getNumberOfDelayedMessages() {
1215+
public long getNumberOfDelayedMessages() {
12161216
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getNumberOfDelayedMessages).orElse(0L);
12171217
}
12181218

@@ -1291,20 +1291,15 @@ public PersistentTopic getTopic() {
12911291
}
12921292

12931293

1294-
public synchronized long getDelayedTrackerMemoryUsage() {
1294+
public long getDelayedTrackerMemoryUsage() {
12951295
return delayedDeliveryTracker.map(DelayedDeliveryTracker::getBufferMemoryUsage).orElse(0L);
12961296
}
12971297

1298-
public synchronized Map<String, TopicMetricBean> getBucketDelayedIndexStats() {
1299-
if (delayedDeliveryTracker.isEmpty()) {
1300-
return Collections.emptyMap();
1301-
}
1302-
1303-
if (delayedDeliveryTracker.get() instanceof BucketDelayedDeliveryTracker) {
1304-
return ((BucketDelayedDeliveryTracker) delayedDeliveryTracker.get()).genTopicMetricMap();
1305-
}
1306-
1307-
return Collections.emptyMap();
1298+
public Map<String, TopicMetricBean> getBucketDelayedIndexStats() {
1299+
return delayedDeliveryTracker
1300+
.filter(BucketDelayedDeliveryTracker.class::isInstance)
1301+
.map(tracker -> ((BucketDelayedDeliveryTracker) tracker).genTopicMetricMap())
1302+
.orElse(Collections.emptyMap());
13081303
}
13091304

13101305
@Override

0 commit comments

Comments
 (0)