From ef98a135e526036911807351248bb44f2bf19d31 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Thu, 15 Jul 2021 17:51:40 +0800 Subject: [PATCH 01/15] support broker level dispatch rate limiter --- conf/broker.conf | 8 ++ .../pulsar/broker/ServiceConfiguration.java | 14 +++ .../service/AbstractBaseDispatcher.java | 23 +++++ .../pulsar/broker/service/BrokerService.java | 18 +++- .../apache/pulsar/broker/service/Topic.java | 4 + ...PersistentDispatcherMultipleConsumers.java | 5 + ...sistentDispatcherSingleActiveConsumer.java | 5 + .../persistent/DispatchRateLimiter.java | 34 ++++++- ...PersistentDispatcherMultipleConsumers.java | 61 ++++++------ ...sistentDispatcherSingleActiveConsumer.java | 95 +++++++++--------- ...tStickyKeyDispatcherMultipleConsumers.java | 5 + .../service/persistent/PersistentTopic.java | 5 + ...criptionMessageDispatchThrottlingTest.java | 96 +++++++++++++++++++ 13 files changed, 291 insertions(+), 82 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index eaabbd6955abc..307f9a06f079b 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -323,6 +323,14 @@ brokerPublisherThrottlingMaxMessageRate=0 # (Disable byte rate limit with value 0) brokerPublisherThrottlingMaxByteRate=0 +# Max Rate(in 1 seconds) of Message allowed to dispatch from a broker if broker dispatch rate limiting enabled +# (Disable message rate limit with value 0) +brokerDispatchThrottlingMaxMessageRate=0 + +# Max Rate(in 1 seconds) of Byte allowed to dispatch from a broker if broker dispatch rate limiting enabled. +# (Disable byte rate limit with value 0) +brokerDispatchThrottlingMaxByteRate=0 + # Max Rate(in 1 seconds) of Message allowed to publish for a topic if topic publish rate limiting enabled # (Disable byte rate limit with value 0) maxPublishRatePerTopicInMessages=0 diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 27ea82aca25d1..15d74948fdb40 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -644,6 +644,20 @@ public class ServiceConfiguration implements PulsarConfiguration { + "when broker publish rate limiting enabled. (Disable byte rate limit with value 0)" ) private long brokerPublisherThrottlingMaxByteRate = 0; + @FieldContext( + category = CATEGORY_SERVER, + dynamic = true, + doc = "Max Rate(in 1 seconds) of Message allowed to dispatch from a broker " + + "when broker dispatch rate limiting enabled. (Disable message rate limit with value 0)" + ) + private int brokerDispatchThrottlingMaxMessageRate = 0; + @FieldContext( + category = CATEGORY_SERVER, + dynamic = true, + doc = "Max Rate(in 1 seconds) of Byte allowed to dispatch from a broker " + + "when broker dispatch rate limiting enabled. (Disable byte rate limit with value 0)" + ) + private long brokerDispatchThrottlingMaxByteRate = 0; @FieldContext( category = CATEGORY_SERVER, dynamic = true, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index f98cfe59c0813..16339514805fa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -27,9 +27,11 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; +import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.CommandAck.AckType; @@ -236,6 +238,27 @@ public void resetCloseFuture() { // noop } + protected abstract void reScheduleRead(); + + protected boolean reachDispatchRateLimit(DispatchRateLimiter dispatchRateLimiter, + MutablePair calculateToRead) { + if (dispatchRateLimiter.isDispatchRateLimitingEnabled()) { + if (!dispatchRateLimiter.hasMessageDispatchPermit()) { + reScheduleRead(); + return true; + } else { + // update messagesToRead according to available dispatch rate limit. + Pair calculateResult = computeReadLimits(calculateToRead.getLeft(), + (int) dispatchRateLimiter.getAvailableDispatchRateLimitOnMsg(), + calculateToRead.getRight(), + dispatchRateLimiter.getAvailableDispatchRateLimitOnByte()); + calculateToRead.setLeft(calculateResult.getLeft()); + calculateToRead.setRight(calculateResult.getRight()); + } + } + return false; + } + protected static Pair computeReadLimits(int messagesToRead, int availablePermitsOnMsg, long bytesToRead, long availablePermitsOnByte) { if (availablePermitsOnMsg > 0) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index bd3c1e9b9bc19..e25880fddf3e1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -227,6 +227,7 @@ public class BrokerService implements Closeable { private ScheduledExecutorService brokerPublishRateLimiterMonitor; private ScheduledExecutorService deduplicationSnapshotMonitor; protected volatile PublishRateLimiter brokerPublishRateLimiter = PublishRateLimiter.DISABLED_RATE_LIMITER; + protected volatile DispatchRateLimiter brokerDispatchRateLimiter = null; private DistributedIdGenerator producerNameGenerator; @@ -450,6 +451,7 @@ public void start() throws Exception { this.startConsumedLedgersMonitor(); this.startBacklogQuotaChecker(); this.updateBrokerPublisherThrottlingMaxRate(); + this.updateBrokerDispatchThrottlingMaxRate(); this.startCheckReplicationPolicies(); this.startDeduplicationSnapshotMonitor(); } @@ -2009,7 +2011,13 @@ private void updateConfigurationAndRegisterListeners() { registerConfigurationListener("brokerPublisherThrottlingMaxByteRate", (brokerPublisherThrottlingMaxByteRate) -> updateBrokerPublisherThrottlingMaxRate()); - + // add listener to notify broker dispatch-rate dynamic config + registerConfigurationListener("brokerDispatchThrottlingMaxMessageRate", + (brokerPublisherThrottlingMaxMessageRate) -> + updateBrokerDispatchThrottlingMaxRate()); + registerConfigurationListener("brokerDispatchThrottlingMaxByteRate", + (brokerPublisherThrottlingMaxByteRate) -> + updateBrokerDispatchThrottlingMaxRate()); // add listener to notify topic publish-rate monitoring if (!preciseTopicPublishRateLimitingEnable) { registerConfigurationListener("topicPublisherThrottlingTickTimeMillis", @@ -2021,6 +2029,14 @@ private void updateConfigurationAndRegisterListeners() { // add more listeners here } + private void updateBrokerDispatchThrottlingMaxRate() { + if (brokerDispatchRateLimiter == null) { + brokerDispatchRateLimiter = new DispatchRateLimiter(this); + } else { + brokerDispatchRateLimiter.updateDispatchRate(); + } + } + private void updateBrokerPublisherThrottlingMaxRate() { int currentMaxMessageRate = pulsar.getConfiguration().getBrokerPublisherThrottlingMaxMessageRate(); long currentMaxByteRate = pulsar.getConfiguration().getBrokerPublisherThrottlingMaxByteRate(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index d02b54a8a6340..53b35de8c0cad 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -246,6 +246,10 @@ default Optional getDispatchRateLimiter() { return Optional.empty(); } + default Optional getBrokerDispatchRateLimiter() { + return Optional.empty(); + } + default boolean isSystemTopic() { return false; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java index 5688683359280..c1094e04e3cf8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java @@ -218,6 +218,11 @@ public boolean isConsumerAvailable(Consumer consumer) { return consumer != null && consumer.getAvailablePermits() > 0 && consumer.isWritable(); } + @Override + protected void reScheduleRead() { + // No-op + } + private static final Logger log = LoggerFactory.getLogger(NonPersistentDispatcherMultipleConsumers.class); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java index 6094ab71df2cc..8ee66b3f2980b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java @@ -110,4 +110,9 @@ protected void readMoreEntries(Consumer consumer) { protected void cancelPendingRead() { // No-op } + + @Override + protected void reScheduleRead() { + // No-op + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java index 696995d51ccb3..96cbf53a92bae 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java @@ -39,7 +39,8 @@ public class DispatchRateLimiter { public enum Type { TOPIC, SUBSCRIPTION, - REPLICATOR + REPLICATOR, + BROKER } private final PersistentTopic topic; @@ -62,6 +63,16 @@ public DispatchRateLimiter(PersistentTopic topic, Type type) { updateDispatchRate(); } + public DispatchRateLimiter(BrokerService brokerService) { + this.topic = null; + this.topicName = null; + this.brokerService = brokerService; + this.type = Type.BROKER; + this.subscriptionRelativeRatelimiterOnMessage = -1; + this.subscriptionRelativeRatelimiterOnByte = -1; + updateDispatchRate(); + } + /** * returns available msg-permit if msg-dispatch-throttling is enabled else it returns -1. * @@ -139,6 +150,10 @@ private DispatchRate createDispatchRate() { dispatchThrottlingRateInMsg = config.getDispatchThrottlingRatePerReplicatorInMsg(); dispatchThrottlingRateInByte = config.getDispatchThrottlingRatePerReplicatorInByte(); break; + case BROKER: + dispatchThrottlingRateInMsg = config.getBrokerDispatchThrottlingMaxMessageRate(); + dispatchThrottlingRateInByte = config.getBrokerDispatchThrottlingMaxByteRate(); + break; default: dispatchThrottlingRateInMsg = -1; dispatchThrottlingRateInByte = -1; @@ -149,7 +164,7 @@ private DispatchRate createDispatchRate() { .dispatchThrottlingRateInMsg(dispatchThrottlingRateInMsg) .dispatchThrottlingRateInByte(dispatchThrottlingRateInByte) .ratePeriodInSecond(1) - .relativeToPublishRate(config.isDispatchThrottlingRateRelativeToPublishRate()) + .relativeToPublishRate(type != Type.BROKER && config.isDispatchThrottlingRateRelativeToPublishRate()) .build(); } @@ -169,12 +184,20 @@ public void updateDispatchRate() { } updateDispatchRate(dispatchRate.get()); - log.info("[{}] configured {} message-dispatch rate at broker {}", this.topicName, type, dispatchRate.get()); - } + if (type == Type.BROKER) { + log.info("configured broker message-dispatch rate {}", dispatchRate.get()); + } else { + log.info("[{}] configured {} message-dispatch rate at broker {}", this.topicName, type, dispatchRate.get()); + } +} public static boolean isDispatchRateNeeded(BrokerService brokerService, Optional policies, String topicName, Type type) { final ServiceConfiguration serviceConfig = brokerService.pulsar().getConfiguration(); + if (type == Type.BROKER) { + return brokerService.getBrokerDispatchRateLimiter().isDispatchRateLimitingEnabled(); + } + Optional dispatchRate = getTopicPolicyDispatchRate(brokerService, topicName, type); if (dispatchRate.isPresent()) { return true; @@ -313,6 +336,9 @@ public static DispatchRateImpl getPoliciesDispatchRate(final String cluster, * @return */ public DispatchRate getPoliciesDispatchRate(BrokerService brokerService) { + if (type == Type.BROKER) { + return null; + } final String cluster = brokerService.pulsar().getConfiguration().getClusterName(); final Optional policies = getPolicies(brokerService, topicName); return getPoliciesDispatchRate(cluster, policies, type); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 06a5c4cad518d..6e64ecd958515 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -38,6 +38,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; @@ -279,6 +280,12 @@ public synchronized void readMoreEntries() { } } + @Override + protected void reScheduleRead() { + topic.getBrokerService().executor().schedule(() -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, + TimeUnit.MILLISECONDS); + } + // left pair is messagesToRead, right pair is bytesToRead protected Pair calculateToRead(int currentTotalAvailablePermits) { int messagesToRead = Math.min(currentTotalAvailablePermits, readBatchSize); @@ -300,55 +307,47 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) messagesToRead = 1; } + MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); + // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { - if (topic.getDispatchRateLimiter().isPresent() - && topic.getDispatchRateLimiter().get().isDispatchRateLimitingEnabled()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); + if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (log.isDebugEnabled()) { + log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, + brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), + MESSAGE_RATE_BACKOFF_MS); + } + return Pair.of(-1, -1L); + } + } + + if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (!topicRateLimiter.hasMessageDispatchPermit()) { + if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, - TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - Pair calculateResult = computeReadLimits(messagesToRead, - (int) topicRateLimiter.getAvailableDispatchRateLimitOnMsg(), - bytesToRead, topicRateLimiter.getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); - } } - if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - if (!dispatchRateLimiter.get().hasMessageDispatchPermit()) { + if (dispatchRateLimiter.isPresent()) { + if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded subscription message-rate {}/{}," - + " schedule after a {}", name, + log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, - TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - Pair calculateResult = computeReadLimits(messagesToRead, - (int) dispatchRateLimiter.get().getAvailableDispatchRateLimitOnMsg(), - bytesToRead, dispatchRateLimiter.get().getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); } } - } if (havePendingReplayRead) { @@ -359,7 +358,9 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - return Pair.of(Math.max(messagesToRead, 1), Math.max(bytesToRead, 1)); + calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); + calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); + return calculateToRead; } protected Set asyncReplayEntries(Set positions) { @@ -572,6 +573,10 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, + totalBytesSent); + } if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, totalBytesSent); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index fa7ac03a18221..e5cf2500dc4ce 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -34,6 +34,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer; import org.apache.pulsar.broker.service.Consumer; @@ -215,6 +216,11 @@ protected void dispatchEntriesToConsumer(Consumer currentConsumer, List e if (future.isSuccess()) { // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit( + sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes()); + } + if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes()); @@ -354,6 +360,23 @@ protected void readMoreEntries(Consumer consumer) { } } + @Override + protected void reScheduleRead() { + topic.getBrokerService().executor().schedule(() -> { + Consumer currentConsumer = ACTIVE_CONSUMER_UPDATER.get(this); + if (currentConsumer != null && !havePendingRead) { + readMoreEntries(currentConsumer); + } else { + if (log.isDebugEnabled()) { + log.debug("[{}] Skipping read retry for topic: Current Consumer {}," + + " havePendingRead {}", + topic.getName(), currentConsumer, havePendingRead); + } + } + }, MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); + + } + protected Pair calculateToRead(Consumer consumer) { int availablePermits = consumer.getAvailablePermits(); if (!consumer.isWritable()) { @@ -372,79 +395,53 @@ protected Pair calculateToRead(Consumer consumer) { messagesToRead = Math.min((int) Math.ceil(availablePermits * 1.0 / avgMessagesPerEntry), readBatchSize); } + MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); + // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { - if (topic.getDispatchRateLimiter().isPresent() - && topic.getDispatchRateLimiter().get().isDispatchRateLimitingEnabled()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); + if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (log.isDebugEnabled()) { + log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, + brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), + MESSAGE_RATE_BACKOFF_MS); + } + return Pair.of(-1, -1L); + } + } + + if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (!topicRateLimiter.hasMessageDispatchPermit()) { + if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> { - Consumer currentConsumer = ACTIVE_CONSUMER_UPDATER.get(this); - if (currentConsumer != null && !havePendingRead) { - readMoreEntries(currentConsumer); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Skipping read retry for topic: Current Consumer {}," - + " havePendingRead {}", - topic.getName(), currentConsumer, havePendingRead); - } - } - }, MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - - Pair calculateResult = computeReadLimits(messagesToRead, - (int) topicRateLimiter.getAvailableDispatchRateLimitOnMsg(), - bytesToRead, topicRateLimiter.getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); - } } - if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - if (!dispatchRateLimiter.get().hasMessageDispatchPermit()) { + if (dispatchRateLimiter.isPresent()) { + if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded subscription message-rate {}/{}," - + " schedule after a {}", - name, dispatchRateLimiter.get().getDispatchRateOnMsg(), + log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, + dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> { - Consumer currentConsumer = ACTIVE_CONSUMER_UPDATER.get(this); - if (currentConsumer != null && !havePendingRead) { - readMoreEntries(currentConsumer); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Skipping read retry: Current Consumer {}, havePendingRead {}", - topic.getName(), currentConsumer, havePendingRead); - } - } - }, MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - - Pair calculateResult = computeReadLimits(messagesToRead, - (int) dispatchRateLimiter.get().getAvailableDispatchRateLimitOnMsg(), - bytesToRead, dispatchRateLimiter.get().getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); } } } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - return Pair.of(Math.max(messagesToRead, 1), Math.max(bytesToRead, 1)); + calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); + calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); + return calculateToRead; } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index f3bcbf2b2e9c9..7bbe30ee40dca 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -252,6 +252,11 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, + totalBytesSent); + } + if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, totalBytesSent); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 0f403a87ca130..9cf29e86e4198 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -2676,6 +2676,11 @@ public Optional getDispatchRateLimiter() { return this.dispatchRateLimiter; } + @Override + public Optional getBrokerDispatchRateLimiter() { + return Optional.ofNullable(this.brokerService.getBrokerDispatchRateLimiter()); + } + public Optional getSubscribeRateLimiter() { return this.subscribeRateLimiter; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index 65b4a716afd65..a6c127fc45a46 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -310,6 +310,102 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT log.info("-- Exiting {} test --", methodName); } + /** + * verify broker level rate-limiting should throttle message-dispatching based on byte-rate + * + *
+     *  1. broker level dispatch-byte-rate = 1000 bytes/sec
+     *  2. start two consumers for two topics
+     *  3. send 15 msgs to each of the two topics: each msgs with 100 byte, total 3000 byte
+     *  4. it should take up to 2 second to receive all messages of the two topic
+     * 
+ * + * @param subscription + * @throws Exception + */ + @Test(dataProvider = "subscriptions", timeOut = 8000) + public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionType subscription) throws Exception { + log.info("-- Starting {} test --", methodName); + + final String namespace1 = "my-property/throttling_ns1"; + final String topicName1 = BrokerTestUtil.newUniqueName("persistent://" + namespace1 + "/throttlingAll"); + final String namespace2 = "my-property/throttling_ns2"; + final String topicName2 = BrokerTestUtil.newUniqueName("persistent://" + namespace2 + "/throttlingAll"); + final String subName = "my-subscriber-name-" + subscription; + + final int byteRate = 1000; + admin.brokers().updateDynamicConfiguration("brokerDispatchThrottlingMaxByteRate", "" + byteRate); + admin.namespaces().createNamespace(namespace1, Sets.newHashSet("test")); + admin.namespaces().createNamespace(namespace2, Sets.newHashSet("test")); + + final int numProducedMessagesEachTopic = 15; + final int numProducedMessages = numProducedMessagesEachTopic * 2; + final CountDownLatch latch = new CountDownLatch(numProducedMessages); + final AtomicInteger totalReceived = new AtomicInteger(0); + // enable throttling for nonBacklog consumers + conf.setDispatchThrottlingOnNonBacklogConsumerEnabled(true); + + Consumer consumer1 = pulsarClient.newConsumer().topic(topicName1).subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(subscription).messageListener((c1, msg) -> { + Assert.assertNotNull(msg, "Message cannot be null"); + String receivedMessage = new String(msg.getData()); + log.debug("Received message [{}] in topic1", receivedMessage); + totalReceived.incrementAndGet(); + latch.countDown(); + }).subscribe(); + + Consumer consumer2 = pulsarClient.newConsumer().topic(topicName2).subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(subscription).messageListener((c1, msg) -> { + Assert.assertNotNull(msg, "Message cannot be null"); + String receivedMessage = new String(msg.getData()); + log.debug("Received message [{}] in topic2", receivedMessage); + totalReceived.incrementAndGet(); + latch.countDown(); + }).subscribe(); + + Producer producer1 = pulsarClient.newProducer().topic(topicName1).create(); + Producer producer2 = pulsarClient.newProducer().topic(topicName2).create(); + + boolean isMessageRateUpdate = false; + DispatchRateLimiter dispatchRateLimiter; + int retry = 5; + for (int i = 0; i < retry; i++) { + dispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); + if (dispatchRateLimiter != null + && dispatchRateLimiter.getDispatchRateOnByte() > 0) { + isMessageRateUpdate = true; + break; + } else { + if (i != retry - 1) { + Thread.sleep(100); + } + } + } + Assert.assertTrue(isMessageRateUpdate); + + long start = System.currentTimeMillis(); + // Asynchronously produce messages + for (int i = 0; i < numProducedMessagesEachTopic; i++) { + producer1.send(new byte[byteRate / 10]); + producer2.send(new byte[byteRate / 10]); + } + latch.await(); + Assert.assertEquals(totalReceived.get(), numProducedMessages, 10); + long end = System.currentTimeMillis(); + log.info("-- time to receive all messages: {} ", end - start); + + // first 10 messages, which equals receiverQueueSize, will not wait. + Assert.assertTrue((end - start) >= 2000); + + consumer1.close(); + consumer2.close(); + producer1.close(); + producer2.close(); + log.info("-- Exiting {} test --", methodName); + } + /** * verify message-rate on multiple consumers with shared-subscription * From 940c67e5733a6794c0a73a6175a4329721aba392 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Thu, 15 Jul 2021 18:10:34 +0800 Subject: [PATCH 02/15] correct debug log --- .../persistent/PersistentDispatcherMultipleConsumers.java | 4 ++-- .../persistent/PersistentDispatcherSingleActiveConsumer.java | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 6e64ecd958515..44ce69290c940 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -340,8 +340,8 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) if (dispatchRateLimiter.isPresent()) { if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, - dispatchRateLimiter.get().getDispatchRateOnMsg(), + log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", + name, dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index e5cf2500dc4ce..00f036eff2e13 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -428,8 +428,8 @@ protected Pair calculateToRead(Consumer consumer) { if (dispatchRateLimiter.isPresent()) { if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, - dispatchRateLimiter.get().getDispatchRateOnMsg(), + log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", + name, dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } From ca4896626116812313235c16a3ee515b68dbe534 Mon Sep 17 00:00:00 2001 From: WangJialing <65590138+wangjialing218@users.noreply.github.com> Date: Fri, 16 Jul 2021 09:38:17 +0800 Subject: [PATCH 03/15] Update pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../client/api/SubscriptionMessageDispatchThrottlingTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index a6c127fc45a46..6d30618a67a96 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -311,7 +311,7 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT } /** - * verify broker level rate-limiting should throttle message-dispatching based on byte-rate + * Verify whether the broker level rate-limiting is throttle message-dispatching based on byte-rate or not * *
      *  1. broker level dispatch-byte-rate = 1000 bytes/sec

From cc935e1f740d77de6723e8caa08a71b3bd3787a3 Mon Sep 17 00:00:00 2001
From: WangJialing <65590138+wangjialing218@users.noreply.github.com>
Date: Fri, 16 Jul 2021 09:38:22 +0800
Subject: [PATCH 04/15] Update
 pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java

Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com>
---
 .../api/SubscriptionMessageDispatchThrottlingTest.java    | 8 ++++----
 1 file changed, 4 insertions(+), 4 deletions(-)

diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java
index 6d30618a67a96..a545547e4a4a3 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java
@@ -314,10 +314,10 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT
      * Verify whether the broker level rate-limiting is throttle message-dispatching based on byte-rate or not
      *
      * 
-     *  1. broker level dispatch-byte-rate = 1000 bytes/sec
-     *  2. start two consumers for two topics
-     *  3. send 15 msgs to each of the two topics: each msgs with 100 byte, total 3000 byte
-     *  4. it should take up to 2 second to receive all messages of the two topic
+     *  1. Broker level dispatch-byte-rate is equal to 1000 bytes per second.
+     *  2. Start two consumers for two topics.
+     *  3. Send 15 msgs to each of the two topics. Each msgs with 100 bytes, thus 3000 bytes in total.
+     *  4. It should take up to 2 seconds to receive all messages of the two topics.
      * 
* * @param subscription From 0fb920d9176f4e4a33a3b94e95c0597f038c4c84 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Fri, 16 Jul 2021 10:22:11 +0800 Subject: [PATCH 05/15] correct member name --- .../java/org/apache/pulsar/broker/service/BrokerService.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index e25880fddf3e1..caf20ce78d40a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2013,10 +2013,10 @@ private void updateConfigurationAndRegisterListeners() { updateBrokerPublisherThrottlingMaxRate()); // add listener to notify broker dispatch-rate dynamic config registerConfigurationListener("brokerDispatchThrottlingMaxMessageRate", - (brokerPublisherThrottlingMaxMessageRate) -> + (brokerDispatchThrottlingMaxMessageRate) -> updateBrokerDispatchThrottlingMaxRate()); registerConfigurationListener("brokerDispatchThrottlingMaxByteRate", - (brokerPublisherThrottlingMaxByteRate) -> + (brokerDispatchThrottlingMaxByteRate) -> updateBrokerDispatchThrottlingMaxRate()); // add listener to notify topic publish-rate monitoring if (!preciseTopicPublishRateLimitingEnable) { From 9d51e98dfc73b3823531c2c67fdb0e37ee4b947e Mon Sep 17 00:00:00 2001 From: wangjialing Date: Mon, 11 Oct 2021 14:34:18 +0800 Subject: [PATCH 06/15] update after review --- .../service/AbstractBaseDispatcher.java | 20 ++- .../persistent/DispatchRateLimiter.java | 2 +- ...PersistentDispatcherMultipleConsumers.java | 30 ++-- ...sistentDispatcherSingleActiveConsumer.java | 30 ++-- ...criptionMessageDispatchThrottlingTest.java | 134 ++++++++++++++++-- 5 files changed, 172 insertions(+), 44 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 16339514805fa..d9b65f1340ee4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -27,7 +27,6 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; @@ -240,25 +239,24 @@ public void resetCloseFuture() { protected abstract void reScheduleRead(); - protected boolean reachDispatchRateLimit(DispatchRateLimiter dispatchRateLimiter, - MutablePair calculateToRead) { + protected boolean reachDispatchRateLimit(DispatchRateLimiter dispatchRateLimiter) { if (dispatchRateLimiter.isDispatchRateLimitingEnabled()) { if (!dispatchRateLimiter.hasMessageDispatchPermit()) { reScheduleRead(); return true; - } else { - // update messagesToRead according to available dispatch rate limit. - Pair calculateResult = computeReadLimits(calculateToRead.getLeft(), - (int) dispatchRateLimiter.getAvailableDispatchRateLimitOnMsg(), - calculateToRead.getRight(), - dispatchRateLimiter.getAvailableDispatchRateLimitOnByte()); - calculateToRead.setLeft(calculateResult.getLeft()); - calculateToRead.setRight(calculateResult.getRight()); } } return false; } + protected Pair updateMessagesToRead(DispatchRateLimiter dispatchRateLimiter, + int messagesToRead, long bytesToRead) { + // update messagesToRead according to available dispatch rate limit. + return computeReadLimits(messagesToRead, + (int) dispatchRateLimiter.getAvailableDispatchRateLimitOnMsg(), + bytesToRead, dispatchRateLimiter.getAvailableDispatchRateLimitOnByte()); + } + protected static Pair computeReadLimits(int messagesToRead, int availablePermitsOnMsg, long bytesToRead, long availablePermitsOnByte) { if (availablePermitsOnMsg > 0) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java index 96cbf53a92bae..d0e2d0f474aa0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java @@ -336,7 +336,7 @@ public static DispatchRateImpl getPoliciesDispatchRate(final String cluster, * @return */ public DispatchRate getPoliciesDispatchRate(BrokerService brokerService) { - if (type == Type.BROKER) { + if (topicName == null) { return null; } final String cluster = brokerService.pulsar().getConfiguration().getClusterName(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 44ce69290c940..547f262ff0277 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -38,7 +38,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; @@ -307,38 +306,46 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) messagesToRead = 1; } - MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); - // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { if (topic.getBrokerDispatchRateLimiter().isPresent()) { DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); - if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(brokerRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(brokerRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(topicRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(topicRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (dispatchRateLimiter.isPresent()) { - if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { + if (reachDispatchRateLimit(dispatchRateLimiter.get())) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", name, dispatchRateLimiter.get().getDispatchRateOnMsg(), @@ -346,6 +353,11 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(dispatchRateLimiter.get(), messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } } @@ -358,9 +370,9 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); - calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); - return calculateToRead; + messagesToRead = Math.max(messagesToRead, 1); + bytesToRead = Math.max(bytesToRead, 1); + return Pair.of(messagesToRead, bytesToRead); } protected Set asyncReplayEntries(Set positions) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 00f036eff2e13..90fcd5237c5a2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -34,7 +34,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer; import org.apache.pulsar.broker.service.Consumer; @@ -395,38 +394,46 @@ protected Pair calculateToRead(Consumer consumer) { messagesToRead = Math.min((int) Math.ceil(availablePermits * 1.0 / avgMessagesPerEntry), readBatchSize); } - MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); - // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { if (topic.getBrokerDispatchRateLimiter().isPresent()) { DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); - if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(brokerRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(brokerRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(topicRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(topicRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (dispatchRateLimiter.isPresent()) { - if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { + if (reachDispatchRateLimit(dispatchRateLimiter.get())) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", name, dispatchRateLimiter.get().getDispatchRateOnMsg(), @@ -434,14 +441,19 @@ protected Pair calculateToRead(Consumer consumer) { MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(dispatchRateLimiter.get(), messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); - calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); - return calculateToRead; + messagesToRead = Math.max(messagesToRead, 1); + bytesToRead = Math.max(bytesToRead, 1); + return Pair.of(messagesToRead, bytesToRead); } @Override diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index a545547e4a4a3..8c6eb63ff0a1f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -20,6 +20,7 @@ import com.google.common.collect.Sets; +import java.time.Duration; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -32,6 +33,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.common.policies.data.DispatchRate; +import org.awaitility.Awaitility; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -310,6 +312,118 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT log.info("-- Exiting {} test --", methodName); } + private void testDispatchRate(SubscriptionType subscription, + int brokerRate, int topicRate, int subRate, int expectRate) throws Exception { + + final String namespace = "my-property/throttling_ns"; + final String topicName = BrokerTestUtil.newUniqueName("persistent://" + namespace + "/throttlingAll"); + final String subName = "my-subscriber-name-" + subscription; + + DispatchRate subscriptionDispatchRate = DispatchRate.builder() + .dispatchThrottlingRateInMsg(-1) + .dispatchThrottlingRateInByte(subRate) + .ratePeriodInSecond(1) + .build(); + DispatchRate topicDispatchRate = DispatchRate.builder() + .dispatchThrottlingRateInMsg(-1) + .dispatchThrottlingRateInByte(topicRate) + .ratePeriodInSecond(1) + .build(); + admin.namespaces().createNamespace(namespace, Sets.newHashSet("test")); + admin.namespaces().setSubscriptionDispatchRate(namespace, subscriptionDispatchRate); + admin.namespaces().setDispatchRate(namespace, topicDispatchRate); + admin.brokers().updateDynamicConfiguration("brokerDispatchThrottlingMaxByteRate", "" + brokerRate); + + final int numProducedMessages = 30; + final CountDownLatch latch = new CountDownLatch(numProducedMessages); + final AtomicInteger totalReceived = new AtomicInteger(0); + // enable throttling for nonBacklog consumers + conf.setDispatchThrottlingOnNonBacklogConsumerEnabled(true); + + Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(subscription).messageListener((c1, msg) -> { + Assert.assertNotNull(msg, "Message cannot be null"); + String receivedMessage = new String(msg.getData()); + log.debug("Received message [{}] in the listener", receivedMessage); + totalReceived.incrementAndGet(); + latch.countDown(); + }).subscribe(); + + Producer producer = pulsarClient.newProducer().topic(topicName).create(); + PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); + + DispatchRateLimiter subRateLimiter = null; + Dispatcher subDispatcher = topic.getSubscription(subName).getDispatcher(); + if (subDispatcher instanceof PersistentDispatcherMultipleConsumers) { + subRateLimiter = subDispatcher.getRateLimiter().get(); + } else if (subDispatcher instanceof PersistentDispatcherSingleActiveConsumer) { + subRateLimiter = subDispatcher.getRateLimiter().get(); + } else { + Assert.fail("Should only have PersistentDispatcher in this test"); + } + final DispatchRateLimiter subDispatchRateLimiter = subRateLimiter; + Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> { + DispatchRateLimiter brokerDispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); + Assert.assertTrue(brokerDispatchRateLimiter != null + && brokerDispatchRateLimiter.getDispatchRateOnByte() > 0); + DispatchRateLimiter topicDispatchRateLimiter = topic.getDispatchRateLimiter().orElse(null); + Assert.assertTrue(topicDispatchRateLimiter != null + && topicDispatchRateLimiter.getDispatchRateOnByte() > 0); + Assert.assertTrue(subDispatchRateLimiter != null + && subDispatchRateLimiter.getDispatchRateOnByte() > 0); + }); + + Assert.assertEquals(admin.namespaces().getSubscriptionDispatchRate(namespace) + .getDispatchThrottlingRateInByte(), subRate); + Assert.assertEquals(admin.namespaces().getDispatchRate(namespace) + .getDispatchThrottlingRateInByte(), topicRate); + + long start = System.currentTimeMillis(); + // Asynchronously produce messages + for (int i = 0; i < numProducedMessages; i++) { + producer.send(new byte[expectRate / 10]); + } + latch.await(); + Assert.assertEquals(totalReceived.get(), numProducedMessages, 10); + long end = System.currentTimeMillis(); + log.info("-- end - start: {} ", end - start); + + // first 10 messages, which equals receiverQueueSize, will not wait. + Assert.assertTrue((end - start) >= 2500); + Assert.assertTrue((end - start) <= 8000); + + consumer.close(); + producer.close(); + admin.topics().delete(topicName, true); + admin.namespaces().deleteNamespace(namespace); + } + + /** + * Verify whether rate-limiting works well when different levels rate-limiting enabled. + * + *
+     *  1. Set broker level, topic level and subscription level dispatch-byte-rate with different limit rate value.
+     *  2. Start one consumer for one topics.
+     *  3. the expect dispatch rate should be the minimum value of different limit rate.
+     * 
+ * + * @param subscription + * @throws Exception + */ + @Test(dataProvider = "subscriptions") + public void testMultiLevelDispatch(SubscriptionType subscription) throws Exception { + log.info("-- Starting {} test --", methodName); + + testDispatchRate(subscription, 1000, 5000, 10000, 1000); + + testDispatchRate(subscription, 10000, 1000, 5000, 1000); + + testDispatchRate(subscription, 5000, 10000, 1000, 1000); + + log.info("-- Exiting {} test --", methodName); + } + /** * Verify whether the broker level rate-limiting is throttle message-dispatching based on byte-rate or not * @@ -370,20 +484,12 @@ public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(Subscri boolean isMessageRateUpdate = false; DispatchRateLimiter dispatchRateLimiter; - int retry = 5; - for (int i = 0; i < retry; i++) { - dispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); - if (dispatchRateLimiter != null - && dispatchRateLimiter.getDispatchRateOnByte() > 0) { - isMessageRateUpdate = true; - break; - } else { - if (i != retry - 1) { - Thread.sleep(100); - } - } - } - Assert.assertTrue(isMessageRateUpdate); + + Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> { + DispatchRateLimiter rateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); + Assert.assertTrue(rateLimiter != null + && rateLimiter.getDispatchRateOnByte() > 0); + }); long start = System.currentTimeMillis(); // Asynchronously produce messages From e4ead19bc09e5961319e07c699a03bff3b7b2517 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Thu, 15 Jul 2021 17:51:40 +0800 Subject: [PATCH 07/15] support broker level dispatch rate limiter --- conf/broker.conf | 8 ++ .../pulsar/broker/ServiceConfiguration.java | 14 +++ .../service/AbstractBaseDispatcher.java | 23 +++++ .../pulsar/broker/service/BrokerService.java | 18 +++- .../apache/pulsar/broker/service/Topic.java | 4 + ...PersistentDispatcherMultipleConsumers.java | 5 + ...sistentDispatcherSingleActiveConsumer.java | 5 + ...PersistentDispatcherMultipleConsumers.java | 61 ++++++------ ...sistentDispatcherSingleActiveConsumer.java | 95 +++++++++--------- ...tStickyKeyDispatcherMultipleConsumers.java | 5 + .../service/persistent/PersistentTopic.java | 5 + ...criptionMessageDispatchThrottlingTest.java | 96 +++++++++++++++++++ 12 files changed, 261 insertions(+), 78 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 1fc578a287f6f..5d28a6122aa10 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -327,6 +327,14 @@ brokerPublisherThrottlingMaxMessageRate=0 # (Disable byte rate limit with value 0) brokerPublisherThrottlingMaxByteRate=0 +# Max Rate(in 1 seconds) of Message allowed to dispatch from a broker if broker dispatch rate limiting enabled +# (Disable message rate limit with value 0) +brokerDispatchThrottlingMaxMessageRate=0 + +# Max Rate(in 1 seconds) of Byte allowed to dispatch from a broker if broker dispatch rate limiting enabled. +# (Disable byte rate limit with value 0) +brokerDispatchThrottlingMaxByteRate=0 + # Max Rate(in 1 seconds) of Message allowed to publish for a topic if topic publish rate limiting enabled # (Disable byte rate limit with value 0) maxPublishRatePerTopicInMessages=0 diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 4c7ee850ad195..4f0458441319e 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -656,6 +656,20 @@ public class ServiceConfiguration implements PulsarConfiguration { + "when broker publish rate limiting enabled. (Disable byte rate limit with value 0)" ) private long brokerPublisherThrottlingMaxByteRate = 0; + @FieldContext( + category = CATEGORY_SERVER, + dynamic = true, + doc = "Max Rate(in 1 seconds) of Message allowed to dispatch from a broker " + + "when broker dispatch rate limiting enabled. (Disable message rate limit with value 0)" + ) + private int brokerDispatchThrottlingMaxMessageRate = 0; + @FieldContext( + category = CATEGORY_SERVER, + dynamic = true, + doc = "Max Rate(in 1 seconds) of Byte allowed to dispatch from a broker " + + "when broker dispatch rate limiting enabled. (Disable byte rate limit with value 0)" + ) + private long brokerDispatchThrottlingMaxByteRate = 0; @FieldContext( category = CATEGORY_SERVER, dynamic = true, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index f98cfe59c0813..16339514805fa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -27,9 +27,11 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; +import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.CommandAck.AckType; @@ -236,6 +238,27 @@ public void resetCloseFuture() { // noop } + protected abstract void reScheduleRead(); + + protected boolean reachDispatchRateLimit(DispatchRateLimiter dispatchRateLimiter, + MutablePair calculateToRead) { + if (dispatchRateLimiter.isDispatchRateLimitingEnabled()) { + if (!dispatchRateLimiter.hasMessageDispatchPermit()) { + reScheduleRead(); + return true; + } else { + // update messagesToRead according to available dispatch rate limit. + Pair calculateResult = computeReadLimits(calculateToRead.getLeft(), + (int) dispatchRateLimiter.getAvailableDispatchRateLimitOnMsg(), + calculateToRead.getRight(), + dispatchRateLimiter.getAvailableDispatchRateLimitOnByte()); + calculateToRead.setLeft(calculateResult.getLeft()); + calculateToRead.setRight(calculateResult.getRight()); + } + } + return false; + } + protected static Pair computeReadLimits(int messagesToRead, int availablePermitsOnMsg, long bytesToRead, long availablePermitsOnByte) { if (availablePermitsOnMsg > 0) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index c3a241d6077ee..091d45a9dec2b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -230,6 +230,7 @@ public class BrokerService implements Closeable { private ScheduledExecutorService brokerPublishRateLimiterMonitor; private ScheduledExecutorService deduplicationSnapshotMonitor; protected volatile PublishRateLimiter brokerPublishRateLimiter = PublishRateLimiter.DISABLED_RATE_LIMITER; + protected volatile DispatchRateLimiter brokerDispatchRateLimiter = null; private DistributedIdGenerator producerNameGenerator; @@ -462,6 +463,7 @@ public void start() throws Exception { this.startConsumedLedgersMonitor(); this.startBacklogQuotaChecker(); this.updateBrokerPublisherThrottlingMaxRate(); + this.updateBrokerDispatchThrottlingMaxRate(); this.startCheckReplicationPolicies(); this.startDeduplicationSnapshotMonitor(); } @@ -2018,7 +2020,13 @@ private void updateConfigurationAndRegisterListeners() { registerConfigurationListener("brokerPublisherThrottlingMaxByteRate", (brokerPublisherThrottlingMaxByteRate) -> updateBrokerPublisherThrottlingMaxRate()); - + // add listener to notify broker dispatch-rate dynamic config + registerConfigurationListener("brokerDispatchThrottlingMaxMessageRate", + (brokerPublisherThrottlingMaxMessageRate) -> + updateBrokerDispatchThrottlingMaxRate()); + registerConfigurationListener("brokerDispatchThrottlingMaxByteRate", + (brokerPublisherThrottlingMaxByteRate) -> + updateBrokerDispatchThrottlingMaxRate()); // add listener to notify topic publish-rate monitoring if (!preciseTopicPublishRateLimitingEnable) { registerConfigurationListener("topicPublisherThrottlingTickTimeMillis", @@ -2030,6 +2038,14 @@ private void updateConfigurationAndRegisterListeners() { // add more listeners here } + private void updateBrokerDispatchThrottlingMaxRate() { + if (brokerDispatchRateLimiter == null) { + brokerDispatchRateLimiter = new DispatchRateLimiter(this); + } else { + brokerDispatchRateLimiter.updateDispatchRate(); + } + } + private void updateBrokerPublisherThrottlingMaxRate() { int currentMaxMessageRate = pulsar.getConfiguration().getBrokerPublisherThrottlingMaxMessageRate(); long currentMaxByteRate = pulsar.getConfiguration().getBrokerPublisherThrottlingMaxByteRate(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index d02b54a8a6340..53b35de8c0cad 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -246,6 +246,10 @@ default Optional getDispatchRateLimiter() { return Optional.empty(); } + default Optional getBrokerDispatchRateLimiter() { + return Optional.empty(); + } + default boolean isSystemTopic() { return false; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java index 5688683359280..c1094e04e3cf8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java @@ -218,6 +218,11 @@ public boolean isConsumerAvailable(Consumer consumer) { return consumer != null && consumer.getAvailablePermits() > 0 && consumer.isWritable(); } + @Override + protected void reScheduleRead() { + // No-op + } + private static final Logger log = LoggerFactory.getLogger(NonPersistentDispatcherMultipleConsumers.class); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java index 6094ab71df2cc..8ee66b3f2980b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java @@ -110,4 +110,9 @@ protected void readMoreEntries(Consumer consumer) { protected void cancelPendingRead() { // No-op } + + @Override + protected void reScheduleRead() { + // No-op + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 06a5c4cad518d..6e64ecd958515 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -38,6 +38,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; @@ -279,6 +280,12 @@ public synchronized void readMoreEntries() { } } + @Override + protected void reScheduleRead() { + topic.getBrokerService().executor().schedule(() -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, + TimeUnit.MILLISECONDS); + } + // left pair is messagesToRead, right pair is bytesToRead protected Pair calculateToRead(int currentTotalAvailablePermits) { int messagesToRead = Math.min(currentTotalAvailablePermits, readBatchSize); @@ -300,55 +307,47 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) messagesToRead = 1; } + MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); + // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { - if (topic.getDispatchRateLimiter().isPresent() - && topic.getDispatchRateLimiter().get().isDispatchRateLimitingEnabled()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); + if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (log.isDebugEnabled()) { + log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, + brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), + MESSAGE_RATE_BACKOFF_MS); + } + return Pair.of(-1, -1L); + } + } + + if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (!topicRateLimiter.hasMessageDispatchPermit()) { + if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, - TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - Pair calculateResult = computeReadLimits(messagesToRead, - (int) topicRateLimiter.getAvailableDispatchRateLimitOnMsg(), - bytesToRead, topicRateLimiter.getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); - } } - if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - if (!dispatchRateLimiter.get().hasMessageDispatchPermit()) { + if (dispatchRateLimiter.isPresent()) { + if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded subscription message-rate {}/{}," - + " schedule after a {}", name, + log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> readMoreEntries(), MESSAGE_RATE_BACKOFF_MS, - TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - Pair calculateResult = computeReadLimits(messagesToRead, - (int) dispatchRateLimiter.get().getAvailableDispatchRateLimitOnMsg(), - bytesToRead, dispatchRateLimiter.get().getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); } } - } if (havePendingReplayRead) { @@ -359,7 +358,9 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - return Pair.of(Math.max(messagesToRead, 1), Math.max(bytesToRead, 1)); + calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); + calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); + return calculateToRead; } protected Set asyncReplayEntries(Set positions) { @@ -572,6 +573,10 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, + totalBytesSent); + } if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, totalBytesSent); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index fa7ac03a18221..e5cf2500dc4ce 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -34,6 +34,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer; import org.apache.pulsar.broker.service.Consumer; @@ -215,6 +216,11 @@ protected void dispatchEntriesToConsumer(Consumer currentConsumer, List e if (future.isSuccess()) { // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit( + sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes()); + } + if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes()); @@ -354,6 +360,23 @@ protected void readMoreEntries(Consumer consumer) { } } + @Override + protected void reScheduleRead() { + topic.getBrokerService().executor().schedule(() -> { + Consumer currentConsumer = ACTIVE_CONSUMER_UPDATER.get(this); + if (currentConsumer != null && !havePendingRead) { + readMoreEntries(currentConsumer); + } else { + if (log.isDebugEnabled()) { + log.debug("[{}] Skipping read retry for topic: Current Consumer {}," + + " havePendingRead {}", + topic.getName(), currentConsumer, havePendingRead); + } + } + }, MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); + + } + protected Pair calculateToRead(Consumer consumer) { int availablePermits = consumer.getAvailablePermits(); if (!consumer.isWritable()) { @@ -372,79 +395,53 @@ protected Pair calculateToRead(Consumer consumer) { messagesToRead = Math.min((int) Math.ceil(availablePermits * 1.0 / avgMessagesPerEntry), readBatchSize); } + MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); + // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { - if (topic.getDispatchRateLimiter().isPresent() - && topic.getDispatchRateLimiter().get().isDispatchRateLimitingEnabled()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); + if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (log.isDebugEnabled()) { + log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, + brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), + MESSAGE_RATE_BACKOFF_MS); + } + return Pair.of(-1, -1L); + } + } + + if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (!topicRateLimiter.hasMessageDispatchPermit()) { + if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> { - Consumer currentConsumer = ACTIVE_CONSUMER_UPDATER.get(this); - if (currentConsumer != null && !havePendingRead) { - readMoreEntries(currentConsumer); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Skipping read retry for topic: Current Consumer {}," - + " havePendingRead {}", - topic.getName(), currentConsumer, havePendingRead); - } - } - }, MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - - Pair calculateResult = computeReadLimits(messagesToRead, - (int) topicRateLimiter.getAvailableDispatchRateLimitOnMsg(), - bytesToRead, topicRateLimiter.getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); - } } - if (dispatchRateLimiter.isPresent() && dispatchRateLimiter.get().isDispatchRateLimitingEnabled()) { - if (!dispatchRateLimiter.get().hasMessageDispatchPermit()) { + if (dispatchRateLimiter.isPresent()) { + if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded subscription message-rate {}/{}," - + " schedule after a {}", - name, dispatchRateLimiter.get().getDispatchRateOnMsg(), + log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, + dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } - topic.getBrokerService().executor().schedule(() -> { - Consumer currentConsumer = ACTIVE_CONSUMER_UPDATER.get(this); - if (currentConsumer != null && !havePendingRead) { - readMoreEntries(currentConsumer); - } else { - if (log.isDebugEnabled()) { - log.debug("[{}] Skipping read retry: Current Consumer {}, havePendingRead {}", - topic.getName(), currentConsumer, havePendingRead); - } - } - }, MESSAGE_RATE_BACKOFF_MS, TimeUnit.MILLISECONDS); return Pair.of(-1, -1L); - } else { - - Pair calculateResult = computeReadLimits(messagesToRead, - (int) dispatchRateLimiter.get().getAvailableDispatchRateLimitOnMsg(), - bytesToRead, dispatchRateLimiter.get().getAvailableDispatchRateLimitOnByte()); - - messagesToRead = calculateResult.getLeft(); - bytesToRead = calculateResult.getRight(); } } } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - return Pair.of(Math.max(messagesToRead, 1), Math.max(bytesToRead, 1)); + calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); + calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); + return calculateToRead; } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index f3bcbf2b2e9c9..7bbe30ee40dca 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -252,6 +252,11 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { + if (topic.getBrokerDispatchRateLimiter().isPresent()) { + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, + totalBytesSent); + } + if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, totalBytesSent); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index d1f5a4e682929..56c44c603d60a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -2676,6 +2676,11 @@ public Optional getDispatchRateLimiter() { return this.dispatchRateLimiter; } + @Override + public Optional getBrokerDispatchRateLimiter() { + return Optional.ofNullable(this.brokerService.getBrokerDispatchRateLimiter()); + } + public Optional getSubscribeRateLimiter() { return this.subscribeRateLimiter; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index 65b4a716afd65..a6c127fc45a46 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -310,6 +310,102 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT log.info("-- Exiting {} test --", methodName); } + /** + * verify broker level rate-limiting should throttle message-dispatching based on byte-rate + * + *
+     *  1. broker level dispatch-byte-rate = 1000 bytes/sec
+     *  2. start two consumers for two topics
+     *  3. send 15 msgs to each of the two topics: each msgs with 100 byte, total 3000 byte
+     *  4. it should take up to 2 second to receive all messages of the two topic
+     * 
+ * + * @param subscription + * @throws Exception + */ + @Test(dataProvider = "subscriptions", timeOut = 8000) + public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionType subscription) throws Exception { + log.info("-- Starting {} test --", methodName); + + final String namespace1 = "my-property/throttling_ns1"; + final String topicName1 = BrokerTestUtil.newUniqueName("persistent://" + namespace1 + "/throttlingAll"); + final String namespace2 = "my-property/throttling_ns2"; + final String topicName2 = BrokerTestUtil.newUniqueName("persistent://" + namespace2 + "/throttlingAll"); + final String subName = "my-subscriber-name-" + subscription; + + final int byteRate = 1000; + admin.brokers().updateDynamicConfiguration("brokerDispatchThrottlingMaxByteRate", "" + byteRate); + admin.namespaces().createNamespace(namespace1, Sets.newHashSet("test")); + admin.namespaces().createNamespace(namespace2, Sets.newHashSet("test")); + + final int numProducedMessagesEachTopic = 15; + final int numProducedMessages = numProducedMessagesEachTopic * 2; + final CountDownLatch latch = new CountDownLatch(numProducedMessages); + final AtomicInteger totalReceived = new AtomicInteger(0); + // enable throttling for nonBacklog consumers + conf.setDispatchThrottlingOnNonBacklogConsumerEnabled(true); + + Consumer consumer1 = pulsarClient.newConsumer().topic(topicName1).subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(subscription).messageListener((c1, msg) -> { + Assert.assertNotNull(msg, "Message cannot be null"); + String receivedMessage = new String(msg.getData()); + log.debug("Received message [{}] in topic1", receivedMessage); + totalReceived.incrementAndGet(); + latch.countDown(); + }).subscribe(); + + Consumer consumer2 = pulsarClient.newConsumer().topic(topicName2).subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(subscription).messageListener((c1, msg) -> { + Assert.assertNotNull(msg, "Message cannot be null"); + String receivedMessage = new String(msg.getData()); + log.debug("Received message [{}] in topic2", receivedMessage); + totalReceived.incrementAndGet(); + latch.countDown(); + }).subscribe(); + + Producer producer1 = pulsarClient.newProducer().topic(topicName1).create(); + Producer producer2 = pulsarClient.newProducer().topic(topicName2).create(); + + boolean isMessageRateUpdate = false; + DispatchRateLimiter dispatchRateLimiter; + int retry = 5; + for (int i = 0; i < retry; i++) { + dispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); + if (dispatchRateLimiter != null + && dispatchRateLimiter.getDispatchRateOnByte() > 0) { + isMessageRateUpdate = true; + break; + } else { + if (i != retry - 1) { + Thread.sleep(100); + } + } + } + Assert.assertTrue(isMessageRateUpdate); + + long start = System.currentTimeMillis(); + // Asynchronously produce messages + for (int i = 0; i < numProducedMessagesEachTopic; i++) { + producer1.send(new byte[byteRate / 10]); + producer2.send(new byte[byteRate / 10]); + } + latch.await(); + Assert.assertEquals(totalReceived.get(), numProducedMessages, 10); + long end = System.currentTimeMillis(); + log.info("-- time to receive all messages: {} ", end - start); + + // first 10 messages, which equals receiverQueueSize, will not wait. + Assert.assertTrue((end - start) >= 2000); + + consumer1.close(); + consumer2.close(); + producer1.close(); + producer2.close(); + log.info("-- Exiting {} test --", methodName); + } + /** * verify message-rate on multiple consumers with shared-subscription * From 3eb0739c0617a1d5be8baf41024d08ea50e904b3 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Thu, 15 Jul 2021 18:10:34 +0800 Subject: [PATCH 08/15] correct debug log --- .../persistent/PersistentDispatcherMultipleConsumers.java | 4 ++-- .../persistent/PersistentDispatcherSingleActiveConsumer.java | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 6e64ecd958515..44ce69290c940 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -340,8 +340,8 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) if (dispatchRateLimiter.isPresent()) { if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, - dispatchRateLimiter.get().getDispatchRateOnMsg(), + log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", + name, dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index e5cf2500dc4ce..00f036eff2e13 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -428,8 +428,8 @@ protected Pair calculateToRead(Consumer consumer) { if (dispatchRateLimiter.isPresent()) { if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, - dispatchRateLimiter.get().getDispatchRateOnMsg(), + log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", + name, dispatchRateLimiter.get().getDispatchRateOnMsg(), dispatchRateLimiter.get().getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } From 90f4a87b9d08e8dcdf9928424e3f1dbb2b9990b4 Mon Sep 17 00:00:00 2001 From: WangJialing <65590138+wangjialing218@users.noreply.github.com> Date: Fri, 16 Jul 2021 09:38:17 +0800 Subject: [PATCH 09/15] Update pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../client/api/SubscriptionMessageDispatchThrottlingTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index a6c127fc45a46..6d30618a67a96 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -311,7 +311,7 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT } /** - * verify broker level rate-limiting should throttle message-dispatching based on byte-rate + * Verify whether the broker level rate-limiting is throttle message-dispatching based on byte-rate or not * *
      *  1. broker level dispatch-byte-rate = 1000 bytes/sec

From e6e02ac1d2dd88a9dada18a5d0efc066374349af Mon Sep 17 00:00:00 2001
From: WangJialing <65590138+wangjialing218@users.noreply.github.com>
Date: Fri, 16 Jul 2021 09:38:22 +0800
Subject: [PATCH 10/15] Update
 pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java

Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com>
---
 .../api/SubscriptionMessageDispatchThrottlingTest.java    | 8 ++++----
 1 file changed, 4 insertions(+), 4 deletions(-)

diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java
index 6d30618a67a96..a545547e4a4a3 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java
@@ -314,10 +314,10 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT
      * Verify whether the broker level rate-limiting is throttle message-dispatching based on byte-rate or not
      *
      * 
-     *  1. broker level dispatch-byte-rate = 1000 bytes/sec
-     *  2. start two consumers for two topics
-     *  3. send 15 msgs to each of the two topics: each msgs with 100 byte, total 3000 byte
-     *  4. it should take up to 2 second to receive all messages of the two topic
+     *  1. Broker level dispatch-byte-rate is equal to 1000 bytes per second.
+     *  2. Start two consumers for two topics.
+     *  3. Send 15 msgs to each of the two topics. Each msgs with 100 bytes, thus 3000 bytes in total.
+     *  4. It should take up to 2 seconds to receive all messages of the two topics.
      * 
* * @param subscription From 43102e1fa12e562c9ef15d4edbc01b13276bac84 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Fri, 16 Jul 2021 10:22:11 +0800 Subject: [PATCH 11/15] correct member name --- .../java/org/apache/pulsar/broker/service/BrokerService.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 091d45a9dec2b..1c023f326f8de 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2022,10 +2022,10 @@ private void updateConfigurationAndRegisterListeners() { updateBrokerPublisherThrottlingMaxRate()); // add listener to notify broker dispatch-rate dynamic config registerConfigurationListener("brokerDispatchThrottlingMaxMessageRate", - (brokerPublisherThrottlingMaxMessageRate) -> + (brokerDispatchThrottlingMaxMessageRate) -> updateBrokerDispatchThrottlingMaxRate()); registerConfigurationListener("brokerDispatchThrottlingMaxByteRate", - (brokerPublisherThrottlingMaxByteRate) -> + (brokerDispatchThrottlingMaxByteRate) -> updateBrokerDispatchThrottlingMaxRate()); // add listener to notify topic publish-rate monitoring if (!preciseTopicPublishRateLimitingEnable) { From 920f365ea8da75221f5755d28e3de3abc3cf4e1b Mon Sep 17 00:00:00 2001 From: wangjialing Date: Mon, 11 Oct 2021 14:34:18 +0800 Subject: [PATCH 12/15] update after review --- .../service/AbstractBaseDispatcher.java | 20 ++- ...PersistentDispatcherMultipleConsumers.java | 30 ++-- ...sistentDispatcherSingleActiveConsumer.java | 30 ++-- ...criptionMessageDispatchThrottlingTest.java | 134 ++++++++++++++++-- 4 files changed, 171 insertions(+), 43 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java index 16339514805fa..d9b65f1340ee4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java @@ -27,7 +27,6 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; @@ -240,25 +239,24 @@ public void resetCloseFuture() { protected abstract void reScheduleRead(); - protected boolean reachDispatchRateLimit(DispatchRateLimiter dispatchRateLimiter, - MutablePair calculateToRead) { + protected boolean reachDispatchRateLimit(DispatchRateLimiter dispatchRateLimiter) { if (dispatchRateLimiter.isDispatchRateLimitingEnabled()) { if (!dispatchRateLimiter.hasMessageDispatchPermit()) { reScheduleRead(); return true; - } else { - // update messagesToRead according to available dispatch rate limit. - Pair calculateResult = computeReadLimits(calculateToRead.getLeft(), - (int) dispatchRateLimiter.getAvailableDispatchRateLimitOnMsg(), - calculateToRead.getRight(), - dispatchRateLimiter.getAvailableDispatchRateLimitOnByte()); - calculateToRead.setLeft(calculateResult.getLeft()); - calculateToRead.setRight(calculateResult.getRight()); } } return false; } + protected Pair updateMessagesToRead(DispatchRateLimiter dispatchRateLimiter, + int messagesToRead, long bytesToRead) { + // update messagesToRead according to available dispatch rate limit. + return computeReadLimits(messagesToRead, + (int) dispatchRateLimiter.getAvailableDispatchRateLimitOnMsg(), + bytesToRead, dispatchRateLimiter.getAvailableDispatchRateLimitOnByte()); + } + protected static Pair computeReadLimits(int messagesToRead, int availablePermitsOnMsg, long bytesToRead, long availablePermitsOnByte) { if (availablePermitsOnMsg > 0) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 44ce69290c940..547f262ff0277 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -38,7 +38,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; @@ -307,38 +306,46 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) messagesToRead = 1; } - MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); - // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { if (topic.getBrokerDispatchRateLimiter().isPresent()) { DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); - if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(brokerRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(brokerRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(topicRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(topicRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (dispatchRateLimiter.isPresent()) { - if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { + if (reachDispatchRateLimit(dispatchRateLimiter.get())) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", name, dispatchRateLimiter.get().getDispatchRateOnMsg(), @@ -346,6 +353,11 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(dispatchRateLimiter.get(), messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } } @@ -358,9 +370,9 @@ protected Pair calculateToRead(int currentTotalAvailablePermits) } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); - calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); - return calculateToRead; + messagesToRead = Math.max(messagesToRead, 1); + bytesToRead = Math.max(bytesToRead, 1); + return Pair.of(messagesToRead, bytesToRead); } protected Set asyncReplayEntries(Set positions) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 00f036eff2e13..90fcd5237c5a2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -34,7 +34,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.util.SafeRun; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer; import org.apache.pulsar.broker.service.Consumer; @@ -395,38 +394,46 @@ protected Pair calculateToRead(Consumer consumer) { messagesToRead = Math.min((int) Math.ceil(availablePermits * 1.0 / avgMessagesPerEntry), readBatchSize); } - MutablePair calculateToRead = MutablePair.of(messagesToRead, bytesToRead); - // throttle only if: (1) cursor is not active (or flag for throttle-nonBacklogConsumer is enabled) bcz // active-cursor reads message from cache rather from bookkeeper (2) if topic has reached message-rate // threshold: then schedule the read after MESSAGE_RATE_BACKOFF_MS if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { if (topic.getBrokerDispatchRateLimiter().isPresent()) { DispatchRateLimiter brokerRateLimiter = topic.getBrokerDispatchRateLimiter().get(); - if (reachDispatchRateLimit(brokerRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(brokerRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", name, brokerRateLimiter.getDispatchRateOnMsg(), brokerRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(brokerRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (topic.getDispatchRateLimiter().isPresent()) { DispatchRateLimiter topicRateLimiter = topic.getDispatchRateLimiter().get(); - if (reachDispatchRateLimit(topicRateLimiter, calculateToRead)) { + if (reachDispatchRateLimit(topicRateLimiter)) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded topic message-rate {}/{}, schedule after a {}", name, topicRateLimiter.getDispatchRateOnMsg(), topicRateLimiter.getDispatchRateOnByte(), MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(topicRateLimiter, messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } if (dispatchRateLimiter.isPresent()) { - if (reachDispatchRateLimit(dispatchRateLimiter.get(), calculateToRead)) { + if (reachDispatchRateLimit(dispatchRateLimiter.get())) { if (log.isDebugEnabled()) { log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", name, dispatchRateLimiter.get().getDispatchRateOnMsg(), @@ -434,14 +441,19 @@ protected Pair calculateToRead(Consumer consumer) { MESSAGE_RATE_BACKOFF_MS); } return Pair.of(-1, -1L); + } else { + Pair calculateToRead = + updateMessagesToRead(dispatchRateLimiter.get(), messagesToRead, bytesToRead); + messagesToRead = calculateToRead.getLeft(); + bytesToRead = calculateToRead.getRight(); } } } // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - calculateToRead.setLeft(Math.max(calculateToRead.getLeft(), 1)); - calculateToRead.setRight(Math.max(calculateToRead.getRight(), 1)); - return calculateToRead; + messagesToRead = Math.max(messagesToRead, 1); + bytesToRead = Math.max(bytesToRead, 1); + return Pair.of(messagesToRead, bytesToRead); } @Override diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index a545547e4a4a3..8c6eb63ff0a1f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -20,6 +20,7 @@ import com.google.common.collect.Sets; +import java.time.Duration; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -32,6 +33,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.common.policies.data.DispatchRate; +import org.awaitility.Awaitility; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -310,6 +312,118 @@ public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionT log.info("-- Exiting {} test --", methodName); } + private void testDispatchRate(SubscriptionType subscription, + int brokerRate, int topicRate, int subRate, int expectRate) throws Exception { + + final String namespace = "my-property/throttling_ns"; + final String topicName = BrokerTestUtil.newUniqueName("persistent://" + namespace + "/throttlingAll"); + final String subName = "my-subscriber-name-" + subscription; + + DispatchRate subscriptionDispatchRate = DispatchRate.builder() + .dispatchThrottlingRateInMsg(-1) + .dispatchThrottlingRateInByte(subRate) + .ratePeriodInSecond(1) + .build(); + DispatchRate topicDispatchRate = DispatchRate.builder() + .dispatchThrottlingRateInMsg(-1) + .dispatchThrottlingRateInByte(topicRate) + .ratePeriodInSecond(1) + .build(); + admin.namespaces().createNamespace(namespace, Sets.newHashSet("test")); + admin.namespaces().setSubscriptionDispatchRate(namespace, subscriptionDispatchRate); + admin.namespaces().setDispatchRate(namespace, topicDispatchRate); + admin.brokers().updateDynamicConfiguration("brokerDispatchThrottlingMaxByteRate", "" + brokerRate); + + final int numProducedMessages = 30; + final CountDownLatch latch = new CountDownLatch(numProducedMessages); + final AtomicInteger totalReceived = new AtomicInteger(0); + // enable throttling for nonBacklog consumers + conf.setDispatchThrottlingOnNonBacklogConsumerEnabled(true); + + Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subName) + .receiverQueueSize(10) + .subscriptionType(subscription).messageListener((c1, msg) -> { + Assert.assertNotNull(msg, "Message cannot be null"); + String receivedMessage = new String(msg.getData()); + log.debug("Received message [{}] in the listener", receivedMessage); + totalReceived.incrementAndGet(); + latch.countDown(); + }).subscribe(); + + Producer producer = pulsarClient.newProducer().topic(topicName).create(); + PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); + + DispatchRateLimiter subRateLimiter = null; + Dispatcher subDispatcher = topic.getSubscription(subName).getDispatcher(); + if (subDispatcher instanceof PersistentDispatcherMultipleConsumers) { + subRateLimiter = subDispatcher.getRateLimiter().get(); + } else if (subDispatcher instanceof PersistentDispatcherSingleActiveConsumer) { + subRateLimiter = subDispatcher.getRateLimiter().get(); + } else { + Assert.fail("Should only have PersistentDispatcher in this test"); + } + final DispatchRateLimiter subDispatchRateLimiter = subRateLimiter; + Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> { + DispatchRateLimiter brokerDispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); + Assert.assertTrue(brokerDispatchRateLimiter != null + && brokerDispatchRateLimiter.getDispatchRateOnByte() > 0); + DispatchRateLimiter topicDispatchRateLimiter = topic.getDispatchRateLimiter().orElse(null); + Assert.assertTrue(topicDispatchRateLimiter != null + && topicDispatchRateLimiter.getDispatchRateOnByte() > 0); + Assert.assertTrue(subDispatchRateLimiter != null + && subDispatchRateLimiter.getDispatchRateOnByte() > 0); + }); + + Assert.assertEquals(admin.namespaces().getSubscriptionDispatchRate(namespace) + .getDispatchThrottlingRateInByte(), subRate); + Assert.assertEquals(admin.namespaces().getDispatchRate(namespace) + .getDispatchThrottlingRateInByte(), topicRate); + + long start = System.currentTimeMillis(); + // Asynchronously produce messages + for (int i = 0; i < numProducedMessages; i++) { + producer.send(new byte[expectRate / 10]); + } + latch.await(); + Assert.assertEquals(totalReceived.get(), numProducedMessages, 10); + long end = System.currentTimeMillis(); + log.info("-- end - start: {} ", end - start); + + // first 10 messages, which equals receiverQueueSize, will not wait. + Assert.assertTrue((end - start) >= 2500); + Assert.assertTrue((end - start) <= 8000); + + consumer.close(); + producer.close(); + admin.topics().delete(topicName, true); + admin.namespaces().deleteNamespace(namespace); + } + + /** + * Verify whether rate-limiting works well when different levels rate-limiting enabled. + * + *
+     *  1. Set broker level, topic level and subscription level dispatch-byte-rate with different limit rate value.
+     *  2. Start one consumer for one topics.
+     *  3. the expect dispatch rate should be the minimum value of different limit rate.
+     * 
+ * + * @param subscription + * @throws Exception + */ + @Test(dataProvider = "subscriptions") + public void testMultiLevelDispatch(SubscriptionType subscription) throws Exception { + log.info("-- Starting {} test --", methodName); + + testDispatchRate(subscription, 1000, 5000, 10000, 1000); + + testDispatchRate(subscription, 10000, 1000, 5000, 1000); + + testDispatchRate(subscription, 5000, 10000, 1000, 1000); + + log.info("-- Exiting {} test --", methodName); + } + /** * Verify whether the broker level rate-limiting is throttle message-dispatching based on byte-rate or not * @@ -370,20 +484,12 @@ public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(Subscri boolean isMessageRateUpdate = false; DispatchRateLimiter dispatchRateLimiter; - int retry = 5; - for (int i = 0; i < retry; i++) { - dispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); - if (dispatchRateLimiter != null - && dispatchRateLimiter.getDispatchRateOnByte() > 0) { - isMessageRateUpdate = true; - break; - } else { - if (i != retry - 1) { - Thread.sleep(100); - } - } - } - Assert.assertTrue(isMessageRateUpdate); + + Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> { + DispatchRateLimiter rateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter(); + Assert.assertTrue(rateLimiter != null + && rateLimiter.getDispatchRateOnByte() > 0); + }); long start = System.currentTimeMillis(); // Asynchronously produce messages From 99a776ff3f4b060c5695f00bea0120ddae22e922 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Thu, 21 Oct 2021 16:55:45 +0800 Subject: [PATCH 13/15] rebase code for #12294 --- .../persistent/PersistentDispatcherMultipleConsumers.java | 3 +-- .../persistent/PersistentDispatcherSingleActiveConsumer.java | 4 ++-- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index f58b09661c528..3d776194e6eb4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -589,8 +589,7 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { long permits = dispatchThrottlingOnBatchMessageEnabled ? totalEntries : totalMessagesSent; if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { if (topic.getBrokerDispatchRateLimiter().isPresent()) { - topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(totalMessagesSent, - totalBytesSent); + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(permits, totalBytesSent); } if (topic.getDispatchRateLimiter().isPresent()) { topic.getDispatchRateLimiter().get().tryDispatchPermit(permits, totalBytesSent); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index b5a599b77b9a5..58e925238b11e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -218,8 +218,8 @@ protected void dispatchEntriesToConsumer(Consumer currentConsumer, List e // acquire message-dispatch permits for already delivered messages if (serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() || !cursor.isActive()) { if (topic.getBrokerDispatchRateLimiter().isPresent()) { - topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit( - sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes()); + topic.getBrokerDispatchRateLimiter().get().tryDispatchPermit(permits, + sendMessageInfo.getTotalBytes()); } if (topic.getDispatchRateLimiter().isPresent()) { From a975212bc9ad7edb71bb12cd18d93d1599353d3a Mon Sep 17 00:00:00 2001 From: wangjialing Date: Thu, 11 Nov 2021 10:21:20 +0800 Subject: [PATCH 14/15] rename configurations --- conf/broker.conf | 16 ++++++++-------- .../pulsar/broker/ServiceConfiguration.java | 12 ++++++------ .../pulsar/broker/service/BrokerService.java | 8 ++++---- .../service/persistent/DispatchRateLimiter.java | 4 ++-- ...ubscriptionMessageDispatchThrottlingTest.java | 4 ++-- 5 files changed, 22 insertions(+), 22 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 10c35c04c58c7..0457b3116d902 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -327,14 +327,6 @@ brokerPublisherThrottlingMaxMessageRate=0 # (Disable byte rate limit with value 0) brokerPublisherThrottlingMaxByteRate=0 -# Max Rate(in 1 seconds) of Message allowed to dispatch from a broker if broker dispatch rate limiting enabled -# (Disable message rate limit with value 0) -brokerDispatchThrottlingMaxMessageRate=0 - -# Max Rate(in 1 seconds) of Byte allowed to dispatch from a broker if broker dispatch rate limiting enabled. -# (Disable byte rate limit with value 0) -brokerDispatchThrottlingMaxByteRate=0 - # Max Rate(in 1 seconds) of Message allowed to publish for a topic if topic publish rate limiting enabled # (Disable byte rate limit with value 0) maxPublishRatePerTopicInMessages=0 @@ -352,6 +344,14 @@ subscribeThrottlingRatePerConsumer=0 # Rate period for {subscribeThrottlingRatePerConsumer}. Default is 30s. subscribeRatePeriodPerConsumerInSecond=30 +# Default messages per second dispatch throttling-limit for whole broker. Using a value of 0, is disabling default +# message dispatch-throttling +dispatchThrottlingRateInMsg=0 + +# Default bytes per second dispatch throttling-limit for whole broker. Using a value of 0, is disabling +# default message-byte dispatch-throttling +dispatchThrottlingRateInByte=0 + # Default messages per second dispatch throttling-limit for every topic. Using a value of 0, is disabling default # message dispatch-throttling dispatchThrottlingRatePerTopicInMsg=0 diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index e5435daf1cddd..867b941e5c2d5 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -659,17 +659,17 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( category = CATEGORY_SERVER, dynamic = true, - doc = "Max Rate(in 1 seconds) of Message allowed to dispatch from a broker " - + "when broker dispatch rate limiting enabled. (Disable message rate limit with value 0)" + doc = "Default messages per second dispatch throttling-limit for whole broker. " + + "Using a value of 0, is disabling default message-byte dispatch-throttling" ) - private int brokerDispatchThrottlingMaxMessageRate = 0; + private int dispatchThrottlingRateInMsg = 0; @FieldContext( category = CATEGORY_SERVER, dynamic = true, - doc = "Max Rate(in 1 seconds) of Byte allowed to dispatch from a broker " - + "when broker dispatch rate limiting enabled. (Disable byte rate limit with value 0)" + doc = "Default bytes per second dispatch throttling-limit for whole broker. " + + "Using a value of 0, is disabling default message-byte dispatch-throttling" ) - private long brokerDispatchThrottlingMaxByteRate = 0; + private long dispatchThrottlingRateInByte = 0; @FieldContext( category = CATEGORY_SERVER, dynamic = true, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 1c023f326f8de..17cc94e3ded0f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2021,11 +2021,11 @@ private void updateConfigurationAndRegisterListeners() { (brokerPublisherThrottlingMaxByteRate) -> updateBrokerPublisherThrottlingMaxRate()); // add listener to notify broker dispatch-rate dynamic config - registerConfigurationListener("brokerDispatchThrottlingMaxMessageRate", - (brokerDispatchThrottlingMaxMessageRate) -> + registerConfigurationListener("dispatchThrottlingRateInMsg", + (dispatchThrottlingRateInMsg) -> updateBrokerDispatchThrottlingMaxRate()); - registerConfigurationListener("brokerDispatchThrottlingMaxByteRate", - (brokerDispatchThrottlingMaxByteRate) -> + registerConfigurationListener("dispatchThrottlingRateInByte", + (dispatchThrottlingRateInByte) -> updateBrokerDispatchThrottlingMaxRate()); // add listener to notify topic publish-rate monitoring if (!preciseTopicPublishRateLimitingEnable) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java index e788d85089abc..bae4cec27e61e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java @@ -152,8 +152,8 @@ private DispatchRate createDispatchRate() { dispatchThrottlingRateInByte = config.getDispatchThrottlingRatePerReplicatorInByte(); break; case BROKER: - dispatchThrottlingRateInMsg = config.getBrokerDispatchThrottlingMaxMessageRate(); - dispatchThrottlingRateInByte = config.getBrokerDispatchThrottlingMaxByteRate(); + dispatchThrottlingRateInMsg = config.getDispatchThrottlingRateInMsg(); + dispatchThrottlingRateInByte = config.getDispatchThrottlingRateInByte(); break; default: dispatchThrottlingRateInMsg = -1; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index 8c6eb63ff0a1f..82fac25d41a42 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -332,7 +332,7 @@ private void testDispatchRate(SubscriptionType subscription, admin.namespaces().createNamespace(namespace, Sets.newHashSet("test")); admin.namespaces().setSubscriptionDispatchRate(namespace, subscriptionDispatchRate); admin.namespaces().setDispatchRate(namespace, topicDispatchRate); - admin.brokers().updateDynamicConfiguration("brokerDispatchThrottlingMaxByteRate", "" + brokerRate); + admin.brokers().updateDynamicConfiguration("dispatchThrottlingRateInByte", "" + brokerRate); final int numProducedMessages = 30; final CountDownLatch latch = new CountDownLatch(numProducedMessages); @@ -448,7 +448,7 @@ public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(Subscri final String subName = "my-subscriber-name-" + subscription; final int byteRate = 1000; - admin.brokers().updateDynamicConfiguration("brokerDispatchThrottlingMaxByteRate", "" + byteRate); + admin.brokers().updateDynamicConfiguration("dispatchThrottlingRateInByte", "" + byteRate); admin.namespaces().createNamespace(namespace1, Sets.newHashSet("test")); admin.namespaces().createNamespace(namespace2, Sets.newHashSet("test")); From fa1c04bbe36cfe96c483307e3959a3dbf33f7de7 Mon Sep 17 00:00:00 2001 From: wangjialing Date: Tue, 18 Jan 2022 14:03:05 +0800 Subject: [PATCH 15/15] fix compile error after rebase master --- .../pulsar/broker/service/persistent/DispatchRateLimiter.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java index 3f8d6bb32f736..528f2d8ce02df 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java @@ -65,8 +65,6 @@ public DispatchRateLimiter(BrokerService brokerService) { this.topicName = null; this.brokerService = brokerService; this.type = Type.BROKER; - this.subscriptionRelativeRatelimiterOnMessage = -1; - this.subscriptionRelativeRatelimiterOnByte = -1; updateDispatchRate(); }