From bd8e6bd88a898dddd1d59dc0c311071c081212a3 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 24 Nov 2022 00:30:16 +0800 Subject: [PATCH 1/2] [improve] [broker] Make dispatch rate limiter more precise --- .../service/AbstractBaseDispatcher.java | 87 ++++ .../AbstractDispatcherMultipleConsumers.java | 22 +- ...bstractDispatcherSingleActiveConsumer.java | 8 + ...PersistentDispatcherMultipleConsumers.java | 112 +---- ...sistentDispatcherSingleActiveConsumer.java | 83 +--- .../BatchedMessageDispatchThrottlingTest.java | 398 ++++++++++++++++++ 6 files changed, 542 insertions(+), 168 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java 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 02400f6cdeee3..a44696da02b6d 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 @@ -18,6 +18,7 @@ */ package org.apache.pulsar.broker.service; +import static org.apache.pulsar.broker.service.persistent.PersistentTopic.MESSAGE_RATE_BACKOFF_MS; import io.netty.buffer.ByteBuf; import io.prometheus.client.Gauge; import java.util.ArrayList; @@ -27,6 +28,7 @@ import java.util.Optional; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.LongAdder; +import javax.annotation.Nullable; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; @@ -46,6 +48,7 @@ import org.apache.pulsar.common.api.proto.ReplicatedSubscriptionsSnapshot; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.protocol.Markers; +import org.apache.pulsar.common.util.Codec; @Slf4j public abstract class AbstractBaseDispatcher extends EntryFilterSupport implements Dispatcher { @@ -381,4 +384,88 @@ public long getFilterRescheduledMsgCount() { protected final void updatePendingBytesToDispatch(long size) { PENDING_BYTES_TO_DISPATCH.inc(size); } + + /** + * Calculate messages count & bytes size to read by rate-limiters. + * @return left pair is messagesToRead, right pair is bytesToRead + */ + protected Pair calculateToReadByRateLimiter(int messageCountToReadDefault, + @Nullable ManagedCursor cursor, + Optional...rateLimiters){ + int messageCountToRead = messageCountToReadDefault; + long bytesToReadDefault = serviceConfig.getDispatcherMaxReadSizeBytes(); + + // 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 + boolean cursorActive = cursor == null /* NonDurableCursor has no backlog */ || cursor.isActive(); + if (!serviceConfig.isDispatchThrottlingOnNonBacklogConsumerEnabled() && cursorActive) { + return Pair.of(messageCountToRead, bytesToReadDefault); + } + + for (Optional rateLimiterOptional : rateLimiters){ + if (!rateLimiterOptional.isPresent()){ + continue; + } + DispatchRateLimiter rateLimiter = rateLimiterOptional.get(); + if (reachDispatchRateLimit(rateLimiter)) { + if (log.isDebugEnabled()) { + log.debug("[{}] message-read exceeded broker message-rate {}/{}, schedule after a {}", + buildName(cursor), rateLimiter.getDispatchRateOnMsg(), rateLimiter.getDispatchRateOnByte(), + MESSAGE_RATE_BACKOFF_MS); + } + return Pair.of(-1, -1L); + } else { + if (rateLimiter.getAvailableDispatchRateLimitOnMsg() > 0) { + messageCountToRead = + Math.min(messageCountToRead, (int) rateLimiter.getAvailableDispatchRateLimitOnMsg()); + } + if (rateLimiter.getAvailableDispatchRateLimitOnByte() > 0){ + bytesToReadDefault = + Math.min(bytesToReadDefault, rateLimiter.getAvailableDispatchRateLimitOnByte()); + } + } + } + // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException + messageCountToRead = Math.max(messageCountToRead, 1); + bytesToReadDefault = Math.max(bytesToReadDefault, 1); + return Pair.of(messageCountToRead, bytesToReadDefault); + } + + protected String buildName(@Nullable ManagedCursor cursor){ + // NonDurableCursor doesn't have cursor. + String cursorName = cursor == null || cursor.getName() == null ? "" : Codec.decode(cursor.getName()); + return subscription.getTopic().getName() + " / " + cursorName; + } + + protected int calculateEntryCountToReadIfEnabledBatch(Topic topic, int messageCountToRead){ + // if turn of precise dispatcher flow control, adjust the records to read. + if (!serviceConfig.isPreciseDispatcherFlowControl()) { + return messageCountToRead; + } + // If the consumer is new, the "consumer.getAvgMessagesPerEntry()" must be 0, + int avgMessagesPerEntry = Math.max(1, calculateAvgMessagesPerEntryInTopic(topic)); + return Math.min((int) Math.ceil(messageCountToRead * 1.0 / avgMessagesPerEntry), + serviceConfig.getDispatcherMaxReadBatchSize()); + } + + protected int calculateAvgMessagesPerEntry(){ + return 0; + } + + protected int calculateAvgMessagesPerEntryInTopic(Topic topic){ + int avgMessagesPerEntry = calculateAvgMessagesPerEntry(); + if (avgMessagesPerEntry > 0){ + return avgMessagesPerEntry; + } + for (Subscription sub : topic.getSubscriptions().values()){ + if (sub.getDispatcher() instanceof AbstractBaseDispatcher otherDispatcher){ + int avgMessagesPerEntryInOtherSub = otherDispatcher.calculateAvgMessagesPerEntry(); + if (avgMessagesPerEntryInOtherSub > 0){ + return avgMessagesPerEntryInOtherSub; + } + } + } + return 1; + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java index 1d01f8c3b00c7..4765c3bace810 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java @@ -237,7 +237,23 @@ private int getFirstConsumerIndexOfPriority(int targetPriority) { return -1; } - private static final Logger log = LoggerFactory.getLogger(PersistentStickyKeyDispatcherMultipleConsumers.class); - - + /** + * @return If the consumer knows, the correct value is returned; otherwise it returns 0. + */ + protected int calculateAvgMessagesPerEntry(){ + if (consumerList.isEmpty() || IS_CLOSED_UPDATER.get(this) == TRUE) { + return 0; + } + Consumer randomConsumer = null; + int nextConsumerIndex = random.nextInt(consumerList.size()); + for (int i = 0; i < consumerList.size(); i++){ + randomConsumer = consumerList.get(nextConsumerIndex); + if (randomConsumer.getAvgMessagesPerEntry() > 0){ + return randomConsumer.getAvgMessagesPerEntry(); + } + nextConsumerIndex++; + nextConsumerIndex = nextConsumerIndex % consumerList.size(); + } + return 0; + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java index 6380fb8384b04..e5c1c5941ac82 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java @@ -282,6 +282,14 @@ public boolean isConsumerConnected() { return ACTIVE_CONSUMER_UPDATER.get(this) != null; } + protected int calculateAvgMessagesPerEntry(){ + Consumer activeConsumer = getActiveConsumer(); + if (activeConsumer == null){ + return 0; + } + return Math.max(0, activeConsumer.getAvgMessagesPerEntry()); + } + private static final Logger log = LoggerFactory.getLogger(AbstractDispatcherSingleActiveConsumer.class); } 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 ca88f9751543b..a7d0f2961737b 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 @@ -269,33 +269,32 @@ public synchronized void readMoreEntries() { int currentTotalAvailablePermits = Math.max(totalAvailablePermits, firstAvailableConsumerPermits); if (currentTotalAvailablePermits > 0 && firstAvailableConsumerPermits > 0) { Pair calculateResult = calculateToRead(currentTotalAvailablePermits); - int messagesToRead = calculateResult.getLeft(); + int entryCountToRead = calculateResult.getLeft(); long bytesToRead = calculateResult.getRight(); - if (messagesToRead == -1 || bytesToRead == -1) { + if (entryCountToRead == -1 || bytesToRead == -1) { // Skip read as topic/dispatcher has exceed the dispatch rate or previous pending read hasn't complete. return; } + NavigableSet entriesToReplayNow = getMessagesToReplayNow(entryCountToRead); - NavigableSet messagesToReplayNow = getMessagesToReplayNow(messagesToRead); - - if (!messagesToReplayNow.isEmpty()) { + if (!entriesToReplayNow.isEmpty()) { if (log.isDebugEnabled()) { - log.debug("[{}] Schedule replay of {} messages for {} consumers", name, messagesToReplayNow.size(), + log.debug("[{}] Schedule replay of {} messages for {} consumers", name, entriesToReplayNow.size(), consumerList.size()); } havePendingReplayRead = true; - minReplayedPosition = messagesToReplayNow.first(); + minReplayedPosition = entriesToReplayNow.first(); Set deletedMessages = topic.isDelayedDeliveryEnabled() - ? asyncReplayEntriesInOrder(messagesToReplayNow) : asyncReplayEntries(messagesToReplayNow); + ? asyncReplayEntriesInOrder(entriesToReplayNow) : asyncReplayEntries(entriesToReplayNow); // clear already acked positions from replay bucket deletedMessages.forEach(position -> redeliveryMessages.remove(((PositionImpl) position).getLedgerId(), ((PositionImpl) position).getEntryId())); // if all the entries are acked-entries and cleared up from redeliveryMessages, try to read // next entries as readCompletedEntries-callback was never called - if ((messagesToReplayNow.size() - deletedMessages.size()) == 0) { + if ((entriesToReplayNow.size() - deletedMessages.size()) == 0) { havePendingReplayRead = false; readMoreEntriesAsync(); } @@ -306,7 +305,7 @@ public synchronized void readMoreEntries() { } } else if (!havePendingRead) { if (log.isDebugEnabled()) { - log.debug("[{}] Schedule read of {} messages for {} consumers", name, messagesToRead, + log.debug("[{}] Schedule read of {} messages for {} consumers", name, entryCountToRead, consumerList.size()); } havePendingRead = true; @@ -318,7 +317,7 @@ public synchronized void readMoreEntries() { minReplayedPosition = null; } - cursor.asyncReadEntriesOrWait(messagesToRead, bytesToRead, this, + cursor.asyncReadEntriesOrWait(entryCountToRead, bytesToRead, this, ReadType.Normal, topic.getMaxReadPosition()); } else { log.debug("[{}] Cannot schedule next read until previous one is done", name); @@ -345,95 +344,26 @@ protected void reScheduleRead() { } } - // left pair is messagesToRead, right pair is bytesToRead + // left pair is entries to read, right pair is bytesToRead protected Pair calculateToRead(int currentTotalAvailablePermits) { - int messagesToRead = Math.min(currentTotalAvailablePermits, readBatchSize); - long bytesToRead = serviceConfig.getDispatcherMaxReadSizeBytes(); - - Consumer c = getRandomConsumer(); - // if turn on precise dispatcher flow control, adjust the record to read - if (c != null && c.isPreciseDispatcherFlowControl()) { - int avgMessagesPerEntry = Math.max(1, c.getAvgMessagesPerEntry()); - messagesToRead = Math.min( - (int) Math.ceil(currentTotalAvailablePermits * 1.0 / avgMessagesPerEntry), - readBatchSize); + if (havePendingReplayRead) { + if (log.isDebugEnabled()) { + log.debug("[{}] Skipping replay while awaiting previous read to complete", name); + } + return Pair.of(-1, -1L); } - if (!isConsumerWritable()) { // If the connection is not currently writable, we issue the read request anyway, but for a single // message. The intent here is to keep use the request as a notification mechanism while avoiding to // read and dispatch a big batch of messages which will need to wait before getting written to the // socket. - messagesToRead = 1; - } - - // 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)) { - 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)) { - 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())) { - if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", - name, dispatchRateLimiter.get().getDispatchRateOnMsg(), - dispatchRateLimiter.get().getDispatchRateOnByte(), - MESSAGE_RATE_BACKOFF_MS); - } - return Pair.of(-1, -1L); - } else { - Pair calculateToRead = - updateMessagesToRead(dispatchRateLimiter.get(), messagesToRead, bytesToRead); - messagesToRead = calculateToRead.getLeft(); - bytesToRead = calculateToRead.getRight(); - } - } - } - - if (havePendingReplayRead) { - if (log.isDebugEnabled()) { - log.debug("[{}] Skipping replay while awaiting previous read to complete", name); - } - return Pair.of(-1, -1L); + currentTotalAvailablePermits = 1; } - // If messagesToRead is 0 or less, correct it to 1 to prevent IllegalArgumentException - messagesToRead = Math.max(messagesToRead, 1); - bytesToRead = Math.max(bytesToRead, 1); - return Pair.of(messagesToRead, bytesToRead); + Pair maxReadLimits = calculateToReadByRateLimiter(currentTotalAvailablePermits, cursor, + topic.getBrokerDispatchRateLimiter(), topic.getDispatchRateLimiter(), dispatchRateLimiter); + int entryCountToRead = calculateEntryCountToReadIfEnabledBatch(topic, maxReadLimits.getLeft()); + return Pair.of(entryCountToRead, maxReadLimits.getRight()); } 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 0421076ddf328..60eac1f91f202 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 @@ -337,26 +337,26 @@ protected void readMoreEntries(Consumer consumer) { } Pair calculateResult = calculateToRead(consumer); - int messagesToRead = calculateResult.getLeft(); + int entryCountToRead = calculateResult.getLeft(); long bytesToRead = calculateResult.getRight(); - if (-1 == messagesToRead || bytesToRead == -1) { + if (-1 == entryCountToRead || bytesToRead == -1) { // Skip read as topic/dispatcher has exceed the dispatch rate. return; } // Schedule read if (log.isDebugEnabled()) { - log.debug("[{}-{}] Schedule read of {} messages", name, consumer, messagesToRead); + log.debug("[{}-{}] Schedule read of {} entries", name, consumer, entryCountToRead); } havePendingRead = true; if (consumer.readCompacted()) { - topic.getCompactedTopic().asyncReadEntriesOrWait(cursor, messagesToRead, isFirstRead, + topic.getCompactedTopic().asyncReadEntriesOrWait(cursor, entryCountToRead, isFirstRead, this, consumer); } else { ReadEntriesCtx readEntriesCtx = ReadEntriesCtx.create(consumer, consumer.getConsumerEpoch()); - cursor.asyncReadEntriesOrWait(messagesToRead, + cursor.asyncReadEntriesOrWait(entryCountToRead, bytesToRead, this, readEntriesCtx, topic.getMaxReadPosition()); } } @@ -390,75 +390,10 @@ protected Pair calculateToRead(Consumer consumer) { // socket. availablePermits = 1; } - - int messagesToRead = Math.min(availablePermits, readBatchSize); - long bytesToRead = serviceConfig.getDispatcherMaxReadSizeBytes(); - // if turn of precise dispatcher flow control, adjust the records to read - if (consumer.isPreciseDispatcherFlowControl()) { - int avgMessagesPerEntry = Math.max(1, consumer.getAvgMessagesPerEntry()); - messagesToRead = Math.min((int) Math.ceil(availablePermits * 1.0 / avgMessagesPerEntry), readBatchSize); - } - - // 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)) { - 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)) { - 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())) { - if (log.isDebugEnabled()) { - log.debug("[{}] message-read exceeded subscription message-rate {}/{}, schedule after a {}", - name, dispatchRateLimiter.get().getDispatchRateOnMsg(), - dispatchRateLimiter.get().getDispatchRateOnByte(), - 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 - messagesToRead = Math.max(messagesToRead, 1); - bytesToRead = Math.max(bytesToRead, 1); - return Pair.of(messagesToRead, bytesToRead); + Pair maxReadLimits = calculateToReadByRateLimiter(availablePermits, cursor, + topic.getBrokerDispatchRateLimiter(), topic.getDispatchRateLimiter(), dispatchRateLimiter); + int entryCountToRead = calculateEntryCountToReadIfEnabledBatch(topic, maxReadLimits.getLeft()); + return Pair.of(entryCountToRead, maxReadLimits.getRight()); } @Override diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java new file mode 100644 index 0000000000000..946e73c4f8188 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java @@ -0,0 +1,398 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.api; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.fail; +import java.lang.reflect.Method; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.function.Predicate; +import lombok.AllArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; +import org.apache.pulsar.broker.BrokerTestUtil; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.client.impl.ConsumerImpl; +import org.apache.pulsar.client.impl.MessageIdImpl; +import org.apache.pulsar.client.impl.MultiTopicsConsumerImpl; +import org.apache.pulsar.common.policies.data.DispatchRate; +import org.apache.pulsar.common.policies.data.impl.DispatchRateImpl; +import org.testcontainers.shaded.org.awaitility.Awaitility; +import org.testcontainers.shaded.org.awaitility.reflect.WhiteboxImpl; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Test; + +@Slf4j +@Test(groups = "broker") +public class BatchedMessageDispatchThrottlingTest extends ProducerConsumerBase { + + private static final int MIN_ENTRY_COUNT_PER_DISPATCH = 5; + + private static final int MAX_ENTRY_COUNT_PER_DISPATCH = 120; + + @BeforeClass + @Override + protected void setup() throws Exception { + this.conf.setClusterName("test"); + super.internalSetup(); + super.producerBaseSetup(); + } + + @AfterClass(alwaysRun = true) + @Override + protected void cleanup() throws Exception { + super.internalCleanup(); + } + + @Override + protected void doInitConf() throws Exception { + super.doInitConf(); + conf.setPreciseDispatcherFlowControl(true); + conf.setDispatcherMaxReadBatchSize(MAX_ENTRY_COUNT_PER_DISPATCH); + conf.setDispatcherMinReadBatchSize(MIN_ENTRY_COUNT_PER_DISPATCH); + conf.setDispatchThrottlingOnNonBacklogConsumerEnabled(true); + } + + private void triggerTopicAndSubCreates(String topicName, String subName, SubscriptionType subscriptionType) + throws Exception{ + Consumer consumer = pulsarClient.newConsumer(Schema.JSON(String.class)) + .topic(topicName) + .readCompacted(false) + .receiverQueueSize(0) + .subscriptionName(subName) + .subscriptionType(subscriptionType) + .subscriptionMode(SubscriptionMode.Durable) + .subscribe(); + consumer.close(); + } + + private List> createConsumes(String topicName, String subName, int consumerCount, + SubscriptionType subscriptionType, int queueSize) throws Exception{ + List> consumers = new ArrayList<>(); + for (int i = 0; i < consumerCount; i++) { + consumers.add(pulsarClient.newConsumer(Schema.JSON(String.class)) + .topic(topicName) + .ackTimeout(Integer.MAX_VALUE, TimeUnit.SECONDS) + .ackTimeoutTickTime(Integer.MAX_VALUE, TimeUnit.SECONDS) + .readCompacted(false) + .receiverQueueSize(queueSize) + .subscriptionName(subName) + .subscriptionType(subscriptionType) + .subscriptionMode(SubscriptionMode.Durable) + .subscribe()); + } + return consumers; + } + + private void sendManyBatchedMessages(int msgCountPerEntry, int entryCount, String topicName) + throws Exception { + sendManyBatchedMessages(msgCountPerEntry, entryCount, topicName, "1"); + } + + private void sendManyBatchedMessages(int msgCountPerEntry, int entryCount, String topicName, String key) + throws Exception { + Producer producer = pulsarClient.newProducer(Schema.JSON(String.class)) + .topic(topicName) + .enableBatching(true) + .batchingMaxPublishDelay(Integer.MAX_VALUE, TimeUnit.SECONDS) + .batchingMaxMessages(Integer.MAX_VALUE) + .create(); + for (int i = 0; i < entryCount; i++){ + for (int j = 0; j < msgCountPerEntry; j++){ + producer.newMessage().key(key).value(String.format("entry-seq[%s], batch_index[%s]", i, j)).sendAsync(); + } + producer.flush(); + } + producer.close(); + } + + private ConsumerKnownHowManyMessagesPerBatch letOneConsumerKnowAvgMsgPerEntryAndReceiveOneEntry( + String topicName, String subName, SubscriptionType subType) throws Exception { + // To ensure just receive one entry, update rate limit to 1msg/s. + DispatchRate originalDispatchRate = admin.topicPolicies().getSubscriptionDispatchRate(topicName); + setSubscriptionDispatchRate(topicName, 1, -1, originalDispatchRate.getRatePeriodInSecond()); + // Receive one entry and rewind. + Consumer consumer = createConsumes(topicName, subName, 1, subType, 1).get(0); + waitForReceivedEntryCount(consumer, Duration.ofSeconds(3), entryCount -> entryCount == 1); + Set receivedEntryIds = collectIncomingQueueEntryIds(consumer); + // Reset dispatch rate. + setSubscriptionDispatchRate(topicName, originalDispatchRate.getDispatchThrottlingRateInMsg(), + originalDispatchRate.getDispatchThrottlingRateInByte(), originalDispatchRate.getRatePeriodInSecond()); + return new ConsumerKnownHowManyMessagesPerBatch(consumer, receivedEntryIds); + } + + @AllArgsConstructor + private static class ConsumerKnownHowManyMessagesPerBatch { + private Consumer consumer; + private Set receivedEntries; + } + + private void triggerManagedLedgerStatUpdate() throws Exception { + ManagedLedgerFactoryImpl managedLedgerFactory = + (ManagedLedgerFactoryImpl) pulsar.getBrokerService().getManagedLedgerFactory(); + Method method = ManagedLedgerFactoryImpl.class.getDeclaredMethod("refreshStats", new Class[]{}); + method.setAccessible(true); + method.invoke(managedLedgerFactory); + } + + private LinkedHashSet collectIncomingQueueEntryIds(Consumer consumer) { + return collectIncomingQueueEntryIds(Arrays.asList(consumer)); + } + + private LinkedHashSet collectIncomingQueueEntryIds(List> consumers) { + List>> allIncomingQueue = new ArrayList<>(); + for (Consumer consumer : consumers){ + if (consumer instanceof MultiTopicsConsumerImpl multiTopicsConsumer){ + List> consumerImplList = multiTopicsConsumer.getConsumers(); + for (Consumer consumerImpl : consumerImplList){ + BlockingQueue> incomingMessages = + WhiteboxImpl.getInternalState(consumerImpl, "incomingMessages"); + allIncomingQueue.add(incomingMessages); + } + } else if (consumer instanceof ConsumerImpl consumerImpl){ + BlockingQueue> incomingMessages = + WhiteboxImpl.getInternalState(consumerImpl, "incomingMessages"); + allIncomingQueue.add(incomingMessages); + } + } + + LinkedHashSet entryIds = new LinkedHashSet<>(); + for (BlockingQueue> incomingQueue : allIncomingQueue){ + ArrayList> messages = new ArrayList<>(); + while (incomingQueue.peek() != null){ + Message message = incomingQueue.poll(); + messages.add(message); + entryIds.add(((MessageIdImpl)message.getMessageId()).getEntryId()); + } + for (Message message : messages){ + incomingQueue.add(message); + } + } + return entryIds; + } + + private void assertWillNotReceiveMessagesAnyMore(final List> consumers, Duration waitTime, + int alreadyReceivedEntryCount){ + try { + waitForAnyConsumerHasReceivedAnyNewMsg(consumers, waitTime, alreadyReceivedEntryCount); + fail("Expected no messages received"); + } catch (Exception ex){ + // Expected no messages received. + } + } + + private void waitForAnyConsumerHasReceivedAnyNewMsg(final List> consumers, Duration maxWaitTime, + int alreadyReceivedMsgCount) { + Awaitility.await().atMost(maxWaitTime).until(() -> { + int currentReceivedEntryCount = collectIncomingQueueEntryIds(consumers).size(); + return currentReceivedEntryCount > alreadyReceivedMsgCount; + }); + } + + private void waitForReceivedEntryCount(final Consumer consumer, Duration maxWaitTime, + Predicate entryCountPredicate) { + waitForReceivedEntryCount(Arrays.asList(consumer), maxWaitTime, entryCountPredicate); + } + + private void waitForReceivedEntryCount(final List> consumers, Duration maxWaitTime, + Predicate entryCountPredicate) { + Awaitility.await().atMost(maxWaitTime).until(() -> { + return entryCountPredicate.test(collectIncomingQueueEntryIds(consumers).size()); + }); + } + + private void closeConsumers(Consumer...consumer) throws Exception { + closeConsumers(Arrays.asList(consumer)); + } + + private void closeConsumers(Collection consumers) throws Exception { + for (Consumer consumer : consumers){ + consumer.close(); + } + } + + private void rewind(Collection consumers) throws Exception { + if (consumers.isEmpty()){ + throw new IllegalArgumentException("consumers collection is empty"); + } + consumers.iterator().next().seek(MessageIdImpl.earliest); + } + + private void setSubscriptionDispatchRate(String topicName, int msgCount, long bytesSize, int ratePeriodInSecond) + throws PulsarAdminException { + DispatchRateImpl dispatchRate = new DispatchRateImpl(); + dispatchRate.setDispatchThrottlingRateInMsg(msgCount); + dispatchRate.setDispatchThrottlingRateInByte(bytesSize); + dispatchRate.setRatePeriodInSecond(ratePeriodInSecond); + admin.topicPolicies().setSubscriptionDispatchRate(topicName, dispatchRate); + waitDispatchRateLimitSetFinish(topicName, msgCount, bytesSize); + } + + private void waitDispatchRateLimitSetFinish(String topicName, long msgCount, long bytesSize) { + Awaitility.await().untilAsserted(() -> { + List subscriptions = pulsar.getBrokerService() + .getTopic(topicName, false).get().get().getSubscriptions().values(); + for (Subscription subscription : subscriptions){ + PersistentSubscription persistentSubscription = (PersistentSubscription) subscription; + Optional rateLimiterOptional = + persistentSubscription.getDispatcher().getRateLimiter(); + if (!rateLimiterOptional.isPresent()){ + fail("rate limiter is null."); + } + DispatchRateLimiter rateLimiter = rateLimiterOptional.get(); + assertEquals(rateLimiter.getDispatchRateOnMsg(), msgCount); + assertEquals(rateLimiter.getDispatchRateOnByte(), bytesSize); + } + }); + + } + + @DataProvider(name = "subTypes") + public Object[][] subTypes(){ + return new Object[][]{ + {SubscriptionType.Shared}, + {SubscriptionType.Key_Shared}, + {SubscriptionType.Failover} + }; + } + + /** + * A wrong logic verify: Consumers will receive overdose messages. this is a bug: when enabled batch rate limit no + * precise even if set {conf.preciseDispatcherFlowControl} to `true`. TODO will be fixed by PR: #18581. + * Because the broker could not know how many messages per entry then broker think that an entry as one message. + * so will receive {rateLimitMsgCountPerPeriod} entries. + */ + @Test(dataProvider = "subTypes") + public void testDispatchRateLimitAtFirstDispatch(SubscriptionType subType) throws Exception { + String topicName = "persistent://public/default/" + BrokerTestUtil.newUniqueName("tp"); + String subName = "sub"; + int ratePeriodInSecond = 60; + int msgCountPerEntry = 10; + int rateLimitMsgCountPerPeriod = 30; + // Create topic, subscription, then send many messages. + triggerTopicAndSubCreates(topicName, subName, subType); + setSubscriptionDispatchRate(topicName, rateLimitMsgCountPerPeriod, -1, ratePeriodInSecond); + sendManyBatchedMessages(msgCountPerEntry, 200, topicName); + triggerManagedLedgerStatUpdate(); + + // verify. + int expectedReceiveEntryCount = rateLimitMsgCountPerPeriod; + List> consumers = createConsumes(topicName, subName, 5, subType, 1000000); + waitForReceivedEntryCount(consumers, Duration.ofSeconds(3), + entryCount -> entryCount == expectedReceiveEntryCount); + int alreadyReceivedEntryCount = expectedReceiveEntryCount; + assertWillNotReceiveMessagesAnyMore(consumers, Duration.ofSeconds(3), alreadyReceivedEntryCount); + + // cleanup. + closeConsumers(consumers); + admin.topics().delete(topicName, false); + } + + /** + * After one active consumer received any messages, then the dispatcher rate limiter works correctly. + */ + @Test(dataProvider = "subTypes") + public void testDispatchRateLimitAfterFirstDispatch(SubscriptionType subType) throws Exception { + String topicName = "persistent://public/default/" + BrokerTestUtil.newUniqueName("tp"); + String subName = "sub"; + String subNameAnother = subName + "-another"; + int ratePeriodInSecond = 60; + int msgCountPerEntry = 10; + int rateLimitMsgCountPerPeriod = 30; + + // Create topic, subscription, rate limit rule. + triggerTopicAndSubCreates(topicName, subName, subType); + triggerTopicAndSubCreates(topicName, subNameAnother, subType); + setSubscriptionDispatchRate(topicName, rateLimitMsgCountPerPeriod, -1, ratePeriodInSecond); + + // send many messages. + sendManyBatchedMessages(msgCountPerEntry, 200, topicName); + triggerManagedLedgerStatUpdate(); + + // Let one active consumer(with another sub) know how many messages per batch. + ConsumerKnownHowManyMessagesPerBatch consumerWithAnotherSub = + letOneConsumerKnowAvgMsgPerEntryAndReceiveOneEntry(topicName, subNameAnother, subType); + + // verify. + int expectedReceiveEntryCount = rateLimitMsgCountPerPeriod / msgCountPerEntry; + List> consumers = createConsumes(topicName, subName, 10, subType, 1000000); + waitForReceivedEntryCount(consumers, Duration.ofSeconds(3), + entryCount -> entryCount == expectedReceiveEntryCount); + int alreadyReceivedEntryCount = expectedReceiveEntryCount; + assertWillNotReceiveMessagesAnyMore(consumers, Duration.ofSeconds(3), alreadyReceivedEntryCount); + + // cleanup. + closeConsumers(consumerWithAnotherSub.consumer); + closeConsumers(consumers); + admin.topics().delete(topicName, false); + } + + /** + * Differ {@link #testDispatchRateLimitAfterFirstDispatch(SubscriptionType)}: just has only one subscription. + */ + @Test + public void testDispatchRateLimitAfterFirstDispatchInOneSub() throws Exception { + String topicName = "persistent://public/default/" + BrokerTestUtil.newUniqueName("tp"); + String subName = "sub"; + SubscriptionType subType = SubscriptionType.Shared; + int ratePeriodInSecond = 60; + int msgCountPerEntry = 10; + int rateLimitMsgCountPerPeriod = 30; + + // Create topic, subscription, rate limit rule. + triggerTopicAndSubCreates(topicName, subName, subType); + setSubscriptionDispatchRate(topicName, rateLimitMsgCountPerPeriod, -1, ratePeriodInSecond); + + // send many messages. + sendManyBatchedMessages(msgCountPerEntry, 200, topicName); + triggerManagedLedgerStatUpdate(); + + // Let one active consumer know how many messages per batch. + ConsumerKnownHowManyMessagesPerBatch firstConsumer = + letOneConsumerKnowAvgMsgPerEntryAndReceiveOneEntry(topicName, subName, subType); + int firstConsumerReceivedEntryCount = firstConsumer.receivedEntries.size(); + + // verify. + int expectedReceiveEntryCount = rateLimitMsgCountPerPeriod / msgCountPerEntry - firstConsumerReceivedEntryCount; + List> consumers = createConsumes(topicName, subName, 10, subType, 1000000); + waitForReceivedEntryCount(consumers, Duration.ofSeconds(3), + entryCount -> entryCount == expectedReceiveEntryCount); + int alreadyReceivedEntryCount = expectedReceiveEntryCount; + assertWillNotReceiveMessagesAnyMore(consumers, Duration.ofSeconds(3), alreadyReceivedEntryCount); + + // cleanup. + closeConsumers(firstConsumer.consumer); + closeConsumers(consumers); + admin.topics().delete(topicName, false); + } +} From 55f3530ba1c46a8379d4c536f60f2599299eb52a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 24 Nov 2022 00:33:15 +0800 Subject: [PATCH 2/2] Make the dispatch rate limit more precise when subscription is created --- .../service/AbstractBaseDispatcher.java | 14 +++-- .../pulsar/broker/service/AbstractTopic.java | 4 ++ .../AvgMessagesPerEntryAccumulator.java | 51 +++++++++++++++++++ .../pulsar/broker/service/Consumer.java | 24 +++------ .../apache/pulsar/broker/service/Topic.java | 1 + .../nonpersistent/NonPersistentTopic.java | 6 +++ .../service/persistent/PersistentTopic.java | 5 ++ .../BatchedMessageDispatchThrottlingTest.java | 7 +-- 8 files changed, 87 insertions(+), 25 deletions(-) create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AvgMessagesPerEntryAccumulator.java 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 a44696da02b6d..7c14aeaa5b7b5 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 @@ -443,8 +443,15 @@ protected int calculateEntryCountToReadIfEnabledBatch(Topic topic, int messageCo if (!serviceConfig.isPreciseDispatcherFlowControl()) { return messageCountToRead; } - // If the consumer is new, the "consumer.getAvgMessagesPerEntry()" must be 0, - int avgMessagesPerEntry = Math.max(1, calculateAvgMessagesPerEntryInTopic(topic)); + // When calculating the "avgMessagesPerEntry," + // 1. look up the attribute avgMessagesPerEntry from the consumers under subscription. + // 2. if all consumers avgMessagesPerEntry is zero, look up from other subscriptions under this topic. + // 3. if still is zero, directly call "topic.getAvgMessagesPerEntryAccumulator()". + // 3. if still is zero too, just read one entry. + int avgMessagesPerEntry = calculateAvgMessagesPerEntryInTopic(topic); + if (avgMessagesPerEntry < 1){ + return 1; + } return Math.min((int) Math.ceil(messageCountToRead * 1.0 / avgMessagesPerEntry), serviceConfig.getDispatcherMaxReadBatchSize()); } @@ -466,6 +473,7 @@ protected int calculateAvgMessagesPerEntryInTopic(Topic topic){ } } } - return 1; + return Math.max(0, + Double.valueOf(topic.getAvgMessagesPerEntryAccumulator().getAvgMessagesPerEntry()).intValue()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 75b15c15df212..5efc7f009b069 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -150,6 +150,10 @@ public abstract class AbstractTopic implements Topic, TopicPolicyListener entryFilters; + @Getter + protected final AvgMessagesPerEntryAccumulator avgMessagesPerEntryAccumulator = + new AvgMessagesPerEntryAccumulator(); + public AbstractTopic(String topic, BrokerService brokerService) { this.topic = topic; this.brokerService = brokerService; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AvgMessagesPerEntryAccumulator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AvgMessagesPerEntryAccumulator.java new file mode 100644 index 0000000000000..ba34b562e0406 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AvgMessagesPerEntryAccumulator.java @@ -0,0 +1,51 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.service; + +import com.google.common.util.concurrent.AtomicDouble; + +/** + * It starts keep tracking the average messages per entry. + * The initial value is 0, when new value comes, it will update with + * avgMessagesPerEntry = avgMessagePerEntry * avgPercent + (1 - avgPercent) * new Value. + */ +public class AvgMessagesPerEntryAccumulator { + + private static final double avgPercent = 0.9; + + private final AtomicDouble avgMessagesPerEntry = new AtomicDouble(0); + + public double getAvgMessagesPerEntry(){ + return avgMessagesPerEntry.get(); + } + + public void setAvgMessagesPerEntry(double avgMessagesPerEntry){ + this.avgMessagesPerEntry.set(avgMessagesPerEntry); + } + + public void accumulate(int totalMessages, int totalEntries) { + if (avgMessagesPerEntry.get() < 1) { //valid avgMessagesPerEntry should always >= 1 + // set init value. + avgMessagesPerEntry.set(1.0 * totalMessages / totalEntries); + } else { + avgMessagesPerEntry.set(avgMessagesPerEntry.get() * avgPercent + + (1 - avgPercent) * totalMessages / totalEntries); + } + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index ac924725f7e81..e7e709cc86d98 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -121,12 +121,7 @@ public class Consumer { private final KeySharedMeta keySharedMeta; - /** - * It starts keep tracking the average messages per entry. - * The initial value is 0, when new value comes, it will update with - * avgMessagesPerEntry = avgMessagePerEntry * avgPercent + (1 - avgPercent) * new Value. - */ - private final AtomicDouble avgMessagesPerEntry = new AtomicDouble(0); + private AvgMessagesPerEntryAccumulator avgMessagesPerEntryAccumulator = new AvgMessagesPerEntryAccumulator(); private static final long [] EMPTY_ACK_SET = new long[0]; private static final double avgPercent = 0.9; @@ -319,14 +314,8 @@ public Future sendMessages(final List entries, EntryBatch } } - // calculate avg message per entry - if (avgMessagesPerEntry.get() < 1) { //valid avgMessagesPerEntry should always >= 1 - // set init value. - avgMessagesPerEntry.set(1.0 * totalMessages / totalEntries); - } else { - avgMessagesPerEntry.set(avgMessagesPerEntry.get() * avgPercent - + (1 - avgPercent) * totalMessages / totalEntries); - } + avgMessagesPerEntryAccumulator.accumulate(totalMessages, totalEntries); + getSubscription().getTopic().getAvgMessagesPerEntryAccumulator().accumulate(totalMessages, totalEntries); // reduce permit and increment unackedMsg count with total number of messages in batch-msgs int ackedCount = batchIndexesAcks == null ? 0 : batchIndexesAcks.getTotalAckedIndexCount(); @@ -334,7 +323,8 @@ public Future sendMessages(final List entries, EntryBatch if (log.isDebugEnabled()){ log.debug("[{}-{}] Added {} minus {} messages to MESSAGE_PERMITS_UPDATER in broker.service.Consumer" + " for consumerId: {}; avgMessagesPerEntry is {}", - topicName, subscription, ackedCount, totalMessages, consumerId, avgMessagesPerEntry.get()); + topicName, subscription, ackedCount, totalMessages, consumerId, + avgMessagesPerEntryAccumulator.getAvgMessagesPerEntry()); } incrementUnackedMessages(unackedMessages); Future writeAndFlushPromise = @@ -777,7 +767,7 @@ public int getAvailablePermits() { * return 0 if there is no entry dispatched yet. */ public int getAvgMessagesPerEntry() { - return (int) Math.round(avgMessagesPerEntry.get()); + return (int) Math.round(avgMessagesPerEntryAccumulator.getAvgMessagesPerEntry()); } public boolean isBlocked() { @@ -836,7 +826,7 @@ public void updateStats(ConsumerStatsImpl consumerStats) { } unackedMessages = consumerStats.unackedMessages; blockedConsumerOnUnackedMsgs = consumerStats.blockedConsumerOnUnackedMsgs; - avgMessagesPerEntry.set(consumerStats.avgMessagesPerEntry); + avgMessagesPerEntryAccumulator.setAvgMessagesPerEntry(consumerStats.avgMessagesPerEntry); } public ConsumerStatsImpl getStats() { 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 3949df92ceca5..2a19d05e38c71 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 @@ -366,4 +366,5 @@ default boolean isSystemTopic() { */ HierarchyTopicPolicies getHierarchyTopicPolicies(); + AvgMessagesPerEntryAccumulator getAvgMessagesPerEntryAccumulator(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index cf46103cc357b..1573f9c57310f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -37,6 +37,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLongFieldUpdater; +import lombok.Getter; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.Position; import org.apache.pulsar.broker.PulsarServerException; @@ -44,6 +45,7 @@ import org.apache.pulsar.broker.resources.NamespaceResources; import org.apache.pulsar.broker.service.AbstractReplicator; import org.apache.pulsar.broker.service.AbstractTopic; +import org.apache.pulsar.broker.service.AvgMessagesPerEntryAccumulator; import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; @@ -107,6 +109,10 @@ public class NonPersistentTopic extends AbstractTopic implements Topic, TopicPol AtomicLongFieldUpdater.newUpdater(NonPersistentTopic.class, "entriesAddedCounter"); private volatile long entriesAddedCounter = 0; + @Getter + protected final AvgMessagesPerEntryAccumulator avgMessagesPerEntryAccumulator = + new AvgMessagesPerEntryAccumulator(); + private volatile boolean migrated = false; private static final FastThreadLocal threadLocalTopicStats = new FastThreadLocal() { @Override 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 ea20d413484cf..35feca7e09dd3 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 @@ -81,6 +81,7 @@ import org.apache.pulsar.broker.resources.NamespaceResources.PartitionedTopicResources; import org.apache.pulsar.broker.service.AbstractReplicator; import org.apache.pulsar.broker.service.AbstractTopic; +import org.apache.pulsar.broker.service.AvgMessagesPerEntryAccumulator; import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.BrokerServiceException.AlreadyRunningException; @@ -229,6 +230,10 @@ protected TopicStatsHelper initialValue() { // Record the last time a data message (ie: not an internal Pulsar marker) is published on the topic private volatile long lastDataMessagePublishedTimestamp = 0; + @Getter + protected final AvgMessagesPerEntryAccumulator avgMessagesPerEntryAccumulator = + new AvgMessagesPerEntryAccumulator(); + private static class TopicStatsHelper { public double averageMsgSize; public double aggMsgRateIn; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java index 946e73c4f8188..a19dd72ff525e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BatchedMessageDispatchThrottlingTest.java @@ -287,10 +287,7 @@ public Object[][] subTypes(){ } /** - * A wrong logic verify: Consumers will receive overdose messages. this is a bug: when enabled batch rate limit no - * precise even if set {conf.preciseDispatcherFlowControl} to `true`. TODO will be fixed by PR: #18581. - * Because the broker could not know how many messages per entry then broker think that an entry as one message. - * so will receive {rateLimitMsgCountPerPeriod} entries. + * The first dispatch. */ @Test(dataProvider = "subTypes") public void testDispatchRateLimitAtFirstDispatch(SubscriptionType subType) throws Exception { @@ -306,7 +303,7 @@ public void testDispatchRateLimitAtFirstDispatch(SubscriptionType subType) throw triggerManagedLedgerStatUpdate(); // verify. - int expectedReceiveEntryCount = rateLimitMsgCountPerPeriod; + int expectedReceiveEntryCount = rateLimitMsgCountPerPeriod / msgCountPerEntry; List> consumers = createConsumes(topicName, subName, 5, subType, 1000000); waitForReceivedEntryCount(consumers, Duration.ofSeconds(3), entryCount -> entryCount == expectedReceiveEntryCount);