From 0593a32072bc106c358c00c4f8a2579b76836253 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Mon, 19 Jun 2023 21:56:12 +0800 Subject: [PATCH 1/2] [fix][broker]fix the publish latency spike issue with large number of producers --- .../pulsar/broker/service/AbstractTopic.java | 20 ++++++++++--------- 1 file changed, 11 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 1371019be41dc..fb9cc0ea3ef6f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -35,6 +35,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.concurrent.atomic.AtomicLongFieldUpdater; import java.util.concurrent.atomic.LongAdder; import java.util.concurrent.locks.ReentrantReadWriteLock; @@ -131,7 +132,11 @@ public abstract class AbstractTopic implements Topic, TopicPolicyListener RATE_LIMITED_UPDATER = AtomicLongFieldUpdater.newUpdater(AbstractTopic.class, "publishRateLimitedTimes"); - protected volatile long publishRateLimitedTimes = 0; + protected volatile long publishRateLimitedTimes = 0L; + + private static final AtomicIntegerFieldUpdater USER_CREATED_PRODUCER_COUNTER_UPDATER = + AtomicIntegerFieldUpdater.newUpdater(AbstractTopic.class, "userCreatedProducerCount"); + private volatile int userCreatedProducerCount = 0; protected volatile Optional topicEpoch = Optional.empty(); private volatile boolean hasExclusiveProducer; @@ -447,14 +452,8 @@ protected boolean isProducersExceeded(Producer producer) { return false; } Integer maxProducers = topicPolicies.getMaxProducersPerTopic().get(); - if (maxProducers != null && maxProducers > 0 && maxProducers <= getUserCreatedProducersSize()) { - return true; - } - return false; - } - - private long getUserCreatedProducersSize() { - return producers.values().stream().filter(p -> !p.isRemote()).count(); + return maxProducers != null && maxProducers > 0 + && maxProducers <= USER_CREATED_PRODUCER_COUNTER_UPDATER.get(this); } protected void registerTopicPolicyListener() { @@ -988,6 +987,8 @@ protected void internalAddProducer(Producer producer) throws BrokerServiceExcept Producer existProducer = producers.putIfAbsent(producer.getProducerName(), producer); if (existProducer != null) { tryOverwriteOldProducer(existProducer, producer); + } else { + USER_CREATED_PRODUCER_COUNTER_UPDATER.incrementAndGet(this); } } @@ -1020,6 +1021,7 @@ public void removeProducer(Producer producer) { checkArgument(producer.getTopic() == this); if (producers.remove(producer.getProducerName(), producer)) { + USER_CREATED_PRODUCER_COUNTER_UPDATER.decrementAndGet(this); handleProducerRemoved(producer); } } From f601751f61bc3bdcfb6b2c9a098d2a6e1227980c Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Mon, 19 Jun 2023 22:05:01 +0800 Subject: [PATCH 2/2] Add isRemote check --- .../org/apache/pulsar/broker/service/AbstractTopic.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index fb9cc0ea3ef6f..827069045509b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -987,7 +987,7 @@ protected void internalAddProducer(Producer producer) throws BrokerServiceExcept Producer existProducer = producers.putIfAbsent(producer.getProducerName(), producer); if (existProducer != null) { tryOverwriteOldProducer(existProducer, producer); - } else { + } else if (!producer.isRemote()) { USER_CREATED_PRODUCER_COUNTER_UPDATER.incrementAndGet(this); } } @@ -1021,7 +1021,9 @@ public void removeProducer(Producer producer) { checkArgument(producer.getTopic() == this); if (producers.remove(producer.getProducerName(), producer)) { - USER_CREATED_PRODUCER_COUNTER_UPDATER.decrementAndGet(this); + if (!producer.isRemote()) { + USER_CREATED_PRODUCER_COUNTER_UPDATER.decrementAndGet(this); + } handleProducerRemoved(producer); } }