From 68d15b5822e16d245ed8ffff7b5197795157d570 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Sat, 1 Feb 2020 00:53:07 +0800 Subject: [PATCH 1/6] Add maxMessagePublishBufferSizeInMB configuration to avoid broker OOM --- conf/broker.conf | 8 ++ .../pulsar/broker/ServiceConfiguration.java | 11 +++ .../pulsar/broker/service/AbstractTopic.java | 3 +- .../pulsar/broker/service/BrokerService.java | 61 ++++++++++++ .../pulsar/broker/service/Producer.java | 12 +-- .../pulsar/broker/service/ServerCnx.java | 14 +-- .../apache/pulsar/broker/service/Topic.java | 2 + .../MessagePublishBufferThrottleTest.java | 97 +++++++++++++++++++ .../naming/ServiceConfigurationTest.java | 1 + .../configurations/pulsar_broker_test.conf | 1 + 10 files changed, 196 insertions(+), 14 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java diff --git a/conf/broker.conf b/conf/broker.conf index a98aa16ad84dc..9e445428b0985 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -324,6 +324,14 @@ replicatedSubscriptionsSnapshotTimeoutSeconds=30 # Max number of snapshot to be cached per subscription. replicatedSubscriptionsSnapshotMaxCachedPerSubscription=10 +# Max memory size for broker handling messages sending from producers. +# If the processing message size exceed this value, broker will stop read data +# from the connection. The processing messages means messages are sends to broker +# but broker have not send response to client, usually waiting to write to bookies. +# It's shared across all the topics running in the same broker. +# Use -1 to disable the memory limitation. Default is 1/5 of direct memory. +maxMessagePublishBufferSizeInMB= + ### --- Authentication --- ### # Role names that are treated as "proxy roles". If the broker sees a request with #role as proxyRoles - it will demand to see a valid original principal. 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 d44f84644e920..49fd2287e6b4f 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 @@ -604,6 +604,17 @@ public class ServiceConfiguration implements PulsarConfiguration { doc = "Max number of snapshot to be cached per subscription.") private int replicatedSubscriptionsSnapshotMaxCachedPerSubscription = 10; + @FieldContext( + category = CATEGORY_SERVER, + doc = "Max memory size for broker handling messages sending from producers.\n\n" + + " If the processing message size exceed this value, broker will stop read data" + + " from the connection. The processing messages means messages are sends to broker" + + " but broker have not send response to client, usually waiting to write to bookies.\n\n" + + " It's shared across all the topics running in the same broker.\n\n" + + " Use -1 to disable the memory limitation. Default is 1/5 of direct memory.\n\n") + private int maxMessagePublishBufferSizeInMB = Math.max(64, + (int) (PlatformDependent.maxDirectMemory() / 5 / (1024 * 1024))); + /**** --- Messaging Protocols --- ****/ @FieldContext( 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 173e3d074e900..1275a59586f0c 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 @@ -292,7 +292,8 @@ public void resetBrokerPublishCountAndEnableReadIfRequired(boolean doneBrokerRes /** * it sets cnx auto-readable if producer's cnx is disabled due to publish-throttling */ - protected void enableProducerRead() { + @Override + public void enableProducerRead() { if (producers != null) { producers.values().forEach(producer -> producer.getCnx().enableCnxAutoRead()); } 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 46f322ff3729b..d2a13e3c957a2 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 @@ -27,6 +27,7 @@ import static org.apache.pulsar.broker.cache.LocalZooKeeperCacheService.LOCAL_POLICIES_ROOT; import static org.apache.pulsar.broker.web.PulsarWebResource.joinPath; +import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.common.collect.Queues; @@ -53,11 +54,13 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ExecutionException; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.LongAdder; import java.util.concurrent.locks.ReadWriteLock; @@ -216,8 +219,21 @@ public class BrokerService implements Closeable, ZooKeeperCacheListener 0 ? + pulsar.getConfiguration().getMaxMessagePublishBufferSizeInMB() * 1024 * 1024 : -1; + this.resumeProducerReadMessagePublishBufferSize = this.maxMessagePublishBufferSize / 2; + this.currentMessagePublishBufferSize = new AtomicLong(0); this.managedLedgerFactory = pulsar.getManagedLedgerFactory(); this.topics = new ConcurrentOpenHashMap<>(); this.replicationClients = new ConcurrentOpenHashMap<>(); @@ -2011,4 +2027,49 @@ public Optional getListenPortTls() { return Optional.empty(); } } + + private void enableTopicsAutoRead() { + topics.values().forEach(future -> { + if (future.isDone() && !future.isCompletedExceptionally()) { + try { + future.get().ifPresent(Topic::enableProducerRead); + } catch (InterruptedException | ExecutionException e) { + // no-op + } + } + }); + } + + @VisibleForTesting + boolean increasePublishBufferSizeAndCheckStopRead(int msgSize) { + if (maxMessagePublishBufferSize < 0) { + return false; + } + if (currentMessagePublishBufferSize.addAndGet(msgSize) >= maxMessagePublishBufferSize && + !isMessagePublishBufferThreshold) { + isMessagePublishBufferThreshold = true; + messagePublishBufferThrottleTimes++; + } + return isMessagePublishBufferThreshold; + } + + @VisibleForTesting + boolean decreasePublishBufferSizeAndCheckResumeRead(int msgSize) { + if (maxMessagePublishBufferSize < 0) { + return false; + } + if (currentMessagePublishBufferSize.addAndGet(-msgSize) < resumeProducerReadMessagePublishBufferSize && + isMessagePublishBufferThreshold) { + isMessagePublishBufferThreshold = false; + messagePublishBufferResumeTimes++; + enableTopicsAutoRead(); + return true; + } + return false; + } + + @VisibleForTesting + AtomicLong getCurrentMessagePublishBufferSize() { + return currentMessagePublishBufferSize; + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index b6c2b496c522c..a7c97da3414b7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -147,7 +147,7 @@ public void publishMessage(long producerId, long lowestSequenceId, long highestS cnx.ctx().channel().eventLoop().execute(() -> { cnx.ctx().writeAndFlush(Commands.newSendError(producerId, highestSequenceId, ServerError.MetadataError, "Invalid lowest or highest sequence id")); - cnx.completedSendOperation(isNonPersistentTopic); + cnx.completedSendOperation(isNonPersistentTopic, headersAndPayload.readableBytes()); }); return; } @@ -160,7 +160,7 @@ public void beforePublish(long producerId, long sequenceId, ByteBuf headersAndPa cnx.ctx().channel().eventLoop().execute(() -> { cnx.ctx().writeAndFlush(Commands.newSendError(producerId, sequenceId, ServerError.PersistenceError, "Producer is closed")); - cnx.completedSendOperation(isNonPersistentTopic); + cnx.completedSendOperation(isNonPersistentTopic, headersAndPayload.readableBytes()); }); return; @@ -170,7 +170,7 @@ public void beforePublish(long producerId, long sequenceId, ByteBuf headersAndPa cnx.ctx().channel().eventLoop().execute(() -> { cnx.ctx().writeAndFlush( Commands.newSendError(producerId, sequenceId, ServerError.ChecksumError, "Checksum failed on the broker")); - cnx.completedSendOperation(isNonPersistentTopic); + cnx.completedSendOperation(isNonPersistentTopic, headersAndPayload.readableBytes()); }); return; } @@ -187,7 +187,7 @@ public void beforePublish(long producerId, long sequenceId, ByteBuf headersAndPa cnx.ctx().channel().eventLoop().execute(() -> { cnx.ctx().writeAndFlush(Commands.newSendError(producerId, sequenceId, ServerError.MetadataError, "Messages must be encrypted")); - cnx.completedSendOperation(isNonPersistentTopic); + cnx.completedSendOperation(isNonPersistentTopic, headersAndPayload.readableBytes()); }); return; } @@ -353,7 +353,7 @@ public void completed(Exception exception, long ledgerId, long entryId) { producer.cnx.ctx().writeAndFlush(Commands.newSendError(producer.producerId, callBackSequenceId, serverError, exception.getMessage())); } - producer.cnx.completedSendOperation(producer.isNonPersistentTopic); + producer.cnx.completedSendOperation(producer.isNonPersistentTopic, msgSize); producer.publishOperationCompleted(); recycle(); }); @@ -385,7 +385,7 @@ public void run() { producer.cnx.ctx().writeAndFlush( Commands.newSendReceipt(producer.producerId, sequenceId, highestSequenceId, ledgerId, entryId), producer.cnx.ctx().voidPromise()); - producer.cnx.completedSendOperation(producer.isNonPersistentTopic); + producer.cnx.completedSendOperation(producer.isNonPersistentTopic, msgSize); producer.publishOperationCompleted(); recycle(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index d79f41bacbaa3..0820a0bdf7767 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1154,7 +1154,7 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { } } - startSendOperation(producer); + startSendOperation(producer, headersAndPayload.readableBytes()); // Persist the message if (send.hasHighestSequenceId() && send.getSequenceId() <= send.getHighestSequenceId()) { @@ -1677,9 +1677,9 @@ public boolean isWritable() { return ctx.channel().isWritable(); } - public void startSendOperation(Producer producer) { + private void startSendOperation(Producer producer, int msgSize) { boolean isPublishRateExceeded = producer.getTopic().isPublishRateExceeded(); - if (++pendingSendRequest == MaxPendingSendRequests || isPublishRateExceeded) { + if (++pendingSendRequest == MaxPendingSendRequests | isPublishRateExceeded | service.increasePublishBufferSizeAndCheckStopRead(msgSize)) { // When the quota of pending send requests is reached, stop reading from socket to cause backpressure on // client connection, possibly shared between multiple producers ctx.channel().config().setAutoRead(false); @@ -1687,8 +1687,8 @@ public void startSendOperation(Producer producer) { } } - public void completedSendOperation(boolean isNonPersistentTopic) { - if (--pendingSendRequest == ResumeReadsThreshold) { + void completedSendOperation(boolean isNonPersistentTopic, int msgSize) { + if (--pendingSendRequest == ResumeReadsThreshold | service.decreasePublishBufferSizeAndCheckResumeRead(msgSize)) { // Resume reading from socket ctx.channel().config().setAutoRead(true); // triggers channel read if autoRead couldn't trigger it @@ -1699,7 +1699,7 @@ public void completedSendOperation(boolean isNonPersistentTopic) { } } - public void enableCnxAutoRead() { + void enableCnxAutoRead() { // we can add check (&& pendingSendRequest < MaxPendingSendRequests) here but then it requires // pendingSendRequest to be volatile and it can be expensive while writing. also this will be called on if // throttling is enable on the topic. so, avoid pendingSendRequest check will be fine. @@ -1724,7 +1724,7 @@ private ServerError getErrorCode(CompletableFuture future) { return error; } - private final void disableTcpNoDelayIfNeeded(String topic, String producerName) { + private void disableTcpNoDelayIfNeeded(String topic, String producerName) { if (producerName != null && producerName.startsWith(replicatorPrefix)) { // Re-enable nagle algorithm on connections used for replication purposes try { 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 26af1c1c5c8bf..3f78f1c94b18e 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 @@ -201,4 +201,6 @@ void updateRates(NamespaceStats nsStats, NamespaceBundleStats currentBundleStats default Optional getDispatchRateLimiter() { return Optional.empty(); } + + void enableProducerRead(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java new file mode 100644 index 0000000000000..3307063791753 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java @@ -0,0 +1,97 @@ +/** + * 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 org.apache.pulsar.client.api.MessageId; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.common.util.FutureUtil; +import org.testng.Assert; +import org.testng.annotations.Test; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; + +/** + */ +public class MessagePublishBufferThrottleTest extends BrokerTestBase { + + @Override + protected void setup() throws Exception { + //No-op + } + + @Override + protected void cleanup() throws Exception { + //No-op + } + + @Test + public void testMessagePublishBufferThrottleDisabled() throws Exception { + conf.setMaxMessagePublishBufferSizeInMB(-1); + super.baseSetup(); + Assert.assertFalse(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(1)); + Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1)); + } + + @Test + public void testMessagePublishBufferThrottle() throws Exception { + conf.setMaxMessagePublishBufferSizeInMB(1); + super.baseSetup(); + Assert.assertFalse(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(512 * 1024)); + Assert.assertTrue(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(512 * 1024)); + Assert.assertTrue(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead( 1024)); + Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1024)); + Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(512 * 1024)); + Assert.assertTrue(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1024)); + Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1024)); + Assert.assertTrue(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(514 * 1024)); + } + + @Test + public void testCurrentPublishBufferShouldBeZeroWhenComplete() throws Exception { + conf.setMaxMessagePublishBufferSizeInMB(1); + super.baseSetup(); + + List> futures = new ArrayList<>(); + final int messages = 200; + final int producers = 10; + + List> producerList = new ArrayList<>(); + for (int i = 0; i < producers; i++) { + Producer producer = pulsarClient.newProducer() + .topic("persistent://prop/ns-abc/testCurrentPublishBufferShouldBeZeroWhenComplete") + .enableBatching(false) + .create(); + producerList.add(producer); + } + + for (Producer producer : producerList) { + for (int j = 0; j < messages; j++) { + futures.add(producer.sendAsync(new byte[1024])); + } + } + + FutureUtil.waitForAll(futures).get(); + Assert.assertEquals(futures.size(), messages * producers); + Assert.assertEquals(pulsar.getBrokerService().getCurrentMessagePublishBufferSize().get(), 0); + Assert.assertTrue(pulsar.getBrokerService().messagePublishBufferThrottleTimes > 0 + && pulsar.getBrokerService().messagePublishBufferThrottleTimes == pulsar.getBrokerService().messagePublishBufferResumeTimes); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/common/naming/ServiceConfigurationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/common/naming/ServiceConfigurationTest.java index 24e7152856718..258c1234411c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/common/naming/ServiceConfigurationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/common/naming/ServiceConfigurationTest.java @@ -63,6 +63,7 @@ public void testInit() throws Exception { assertEquals(config.getBrokerDeleteInactiveTopicsMode(), InactiveTopicDeleteMode.delete_when_subscriptions_caught_up); assertEquals(config.getDefaultNamespaceBundleSplitAlgorithm(), "topic_count_equally_divide"); assertEquals(config.getSupportedNamespaceBundleSplitAlgorithms().size(), 1); + assertEquals(config.getMaxMessagePublishBufferSizeInMB(), -1); } @Test diff --git a/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf b/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf index bd966ad76f5b7..cffb006d5b350 100644 --- a/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf +++ b/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf @@ -88,3 +88,4 @@ replicatorPrefix=pulsar.repl brokerDeleteInactiveTopicsMode=delete_when_subscriptions_caught_up supportedNamespaceBundleSplitAlgorithms=[range_equally_divide] defaultNamespaceBundleSplitAlgorithm=topic_count_equally_divide +maxMessagePublishBufferSizeInMB=-1 From c3bbbb02cf5704c82041c2d8398506dff1d5c9ca Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Tue, 11 Feb 2020 23:32:25 +0800 Subject: [PATCH 2/6] Apply comments --- conf/broker.conf | 4 + .../pulsar/broker/ServiceConfiguration.java | 6 + .../pulsar/broker/service/AbstractTopic.java | 23 ++- .../pulsar/broker/service/BrokerService.java | 69 +++++---- .../pulsar/broker/service/ServerCnx.java | 43 +++++- .../apache/pulsar/broker/service/Topic.java | 2 - .../MessagePublishBufferThrottleTest.java | 135 +++++++++++++----- 7 files changed, 204 insertions(+), 78 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 9e445428b0985..340d13694211d 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -332,6 +332,10 @@ replicatedSubscriptionsSnapshotMaxCachedPerSubscription=10 # Use -1 to disable the memory limitation. Default is 1/5 of direct memory. maxMessagePublishBufferSizeInMB= +# Interval between checks to see if message publish buffer size is exceed the max message publish buffer size +# Use 0 or negative number to disable the max publish buffer limiting. +messagePublishBufferCheckIntervalInMills=100 + ### --- Authentication --- ### # Role names that are treated as "proxy roles". If the broker sees a request with #role as proxyRoles - it will demand to see a valid original principal. 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 49fd2287e6b4f..fd32da1df7651 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 @@ -615,6 +615,12 @@ public class ServiceConfiguration implements PulsarConfiguration { private int maxMessagePublishBufferSizeInMB = Math.max(64, (int) (PlatformDependent.maxDirectMemory() / 5 / (1024 * 1024))); + @FieldContext( + category = CATEGORY_SERVER, + doc = "Interval between checks to see if message publish buffer size is exceed the max message publish buffer size" + ) + private int messagePublishBufferCheckIntervalInMills = 100; + /**** --- Messaging Protocols --- ****/ @FieldContext( 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 1275a59586f0c..df6c1d146621f 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 @@ -277,7 +277,7 @@ public void incrementPublishCount(int numOfMessages, long msgSizeInBytes) { public void resetTopicPublishCountAndEnableReadIfRequired() { // broker rate not exceeded. and completed topic limiter reset. if (!getBrokerPublishRateLimiter().isPublishRateExceeded() && topicPublishRateLimiter.resetPublishCount()) { - enableProducerRead(); + enableProducerReadForPublishRateLimiting(); } } @@ -285,17 +285,28 @@ public void resetTopicPublishCountAndEnableReadIfRequired() { public void resetBrokerPublishCountAndEnableReadIfRequired(boolean doneBrokerReset) { // topic rate not exceeded, and completed broker limiter reset. if (!topicPublishRateLimiter.isPublishRateExceeded() && doneBrokerReset) { - enableProducerRead(); + enableProducerReadForPublishRateLimiting(); } } /** * it sets cnx auto-readable if producer's cnx is disabled due to publish-throttling */ - @Override - public void enableProducerRead() { + public void enableProducerReadForPublishRateLimiting() { + if (producers != null) { + producers.values().forEach(producer -> { + producer.getCnx().cancelPublishRateLimiting(); + producer.getCnx().enableCnxAutoRead(); + }); + } + } + + public void enableProducerReadForPublishBufferLimiting() { if (producers != null) { - producers.values().forEach(producer -> producer.getCnx().enableCnxAutoRead()); + producers.values().forEach(producer -> { + producer.getCnx().cancelPublishBufferLimiting(); + producer.getCnx().enableCnxAutoRead(); + }); } } @@ -390,7 +401,7 @@ private void updatePublishDispatcher(Policies policies) { } else { log.info("Disabling publish throttling for {}", this.topic); this.topicPublishRateLimiter = PublishRateLimiter.DISABLED_RATE_LIMITER; - enableProducerRead(); + enableProducerReadForPublishRateLimiting(); } } 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 d2a13e3c957a2..b78b7ae101e36 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 @@ -188,6 +188,7 @@ public class BrokerService implements Closeable, ZooKeeperCacheListener 0 ? pulsar.getConfiguration().getMaxMessagePublishBufferSizeInMB() * 1024 * 1024 : -1; this.resumeProducerReadMessagePublishBufferSize = this.maxMessagePublishBufferSize / 2; - this.currentMessagePublishBufferSize = new AtomicLong(0); + this.currentMessagePublishBufferSize = 0; this.managedLedgerFactory = pulsar.getManagedLedgerFactory(); this.topics = new ConcurrentOpenHashMap<>(); this.replicationClients = new ConcurrentOpenHashMap<>(); @@ -273,6 +274,8 @@ public BrokerService(PulsarService pulsar) throws Exception { .newSingleThreadScheduledExecutor(new DefaultThreadFactory("pulsar-msg-expiry-monitor")); this.compactionMonitor = Executors.newSingleThreadScheduledExecutor(new DefaultThreadFactory("pulsar-compaction-monitor")); + this.messagePublishBufferMonitor = + Executors.newSingleThreadScheduledExecutor(new DefaultThreadFactory("pulsar-publish-buffer-monitor")); this.backlogQuotaManager = new BacklogQuotaManager(pulsar); this.backlogQuotaChecker = Executors @@ -402,6 +405,7 @@ public void start() throws Exception { this.startInactivityMonitor(); this.startMessageExpiryMonitor(); this.startCompactionMonitor(); + this.startMessagePublishBufferMonitor(); this.startBacklogQuotaChecker(); this.updateBrokerPublisherThrottlingMaxRate(); // register listener to capture zk-latency @@ -454,6 +458,14 @@ protected void startCompactionMonitor() { } } + protected void startMessagePublishBufferMonitor() { + int interval = pulsar().getConfiguration().getMessagePublishBufferCheckIntervalInMills(); + if (interval > 0 && maxMessagePublishBufferSize > 0) { + messagePublishBufferMonitor.scheduleAtFixedRate(safeRun(this::checkMessagePublishBuffer), + interval, interval, TimeUnit.MILLISECONDS); + } + } + protected void startBacklogQuotaChecker() { if (pulsar().getConfiguration().isBacklogQuotaCheckEnabled()) { final int interval = pulsar().getConfiguration().getBacklogQuotaCheckIntervalInSeconds(); @@ -2028,48 +2040,35 @@ public Optional getListenPortTls() { } } - private void enableTopicsAutoRead() { - topics.values().forEach(future -> { - if (future.isDone() && !future.isCompletedExceptionally()) { - try { - future.get().ifPresent(Topic::enableProducerRead); - } catch (InterruptedException | ExecutionException e) { - // no-op - } - } - }); - } - - @VisibleForTesting - boolean increasePublishBufferSizeAndCheckStopRead(int msgSize) { - if (maxMessagePublishBufferSize < 0) { - return false; - } - if (currentMessagePublishBufferSize.addAndGet(msgSize) >= maxMessagePublishBufferSize && - !isMessagePublishBufferThreshold) { + private void checkMessagePublishBuffer() { + currentMessagePublishBufferSize = 0; + foreachProducer(producer -> currentMessagePublishBufferSize += producer.getCnx().getMessagePublishBufferSize()); + if (currentMessagePublishBufferSize >= maxMessagePublishBufferSize + && !isMessagePublishBufferThreshold) { isMessagePublishBufferThreshold = true; messagePublishBufferThrottleTimes++; } - return isMessagePublishBufferThreshold; - } - - @VisibleForTesting - boolean decreasePublishBufferSizeAndCheckResumeRead(int msgSize) { - if (maxMessagePublishBufferSize < 0) { - return false; - } - if (currentMessagePublishBufferSize.addAndGet(-msgSize) < resumeProducerReadMessagePublishBufferSize && - isMessagePublishBufferThreshold) { + if (currentMessagePublishBufferSize < resumeProducerReadMessagePublishBufferSize + && isMessagePublishBufferThreshold) { isMessagePublishBufferThreshold = false; messagePublishBufferResumeTimes++; - enableTopicsAutoRead(); - return true; + forEachTopic(topic -> ((AbstractTopic) topic).enableProducerReadForPublishBufferLimiting()); } - return false; + } + + private void foreachProducer(Consumer consumer) { + topics.forEach((n, t) -> { + Optional topic = extractTopic(t); + topic.ifPresent(value -> value.getProducers().values().forEach(consumer)); + }); + } + + public boolean isMessagePublishBufferThreshold() { + return isMessagePublishBufferThreshold; } @VisibleForTesting - AtomicLong getCurrentMessagePublishBufferSize() { + long getCurrentMessagePublishBufferSize() { return currentMessagePublishBufferSize; } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 0820a0bdf7767..aa7139631053e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -25,6 +25,7 @@ import static org.apache.pulsar.common.protocol.Commands.newLookupErrorResponse; import static org.apache.pulsar.common.api.proto.PulsarApi.ProtocolVersion.v5; +import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Strings; import io.netty.buffer.ByteBuf; @@ -151,6 +152,9 @@ public class ServerCnx extends PulsarHandler { // Flag to manage throttling-rate by atomically enable/disable read-channel. private volatile boolean autoReadDisabledRateLimiting = false; private FeatureFlags features; + // Flag to manage throttling-publish-buffer by atomically enable/disable read-channel. + private volatile boolean autoReadDisabledPublishBufferLimiting = false; + private volatile long messagePublishBufferSize = 0; enum State { Start, Connected, Failed, Connecting @@ -1678,17 +1682,20 @@ public boolean isWritable() { } private void startSendOperation(Producer producer, int msgSize) { + messagePublishBufferSize += msgSize; boolean isPublishRateExceeded = producer.getTopic().isPublishRateExceeded(); - if (++pendingSendRequest == MaxPendingSendRequests | isPublishRateExceeded | service.increasePublishBufferSizeAndCheckStopRead(msgSize)) { + if (++pendingSendRequest == MaxPendingSendRequests | isPublishRateExceeded | getBrokerService().isMessagePublishBufferThreshold()) { // When the quota of pending send requests is reached, stop reading from socket to cause backpressure on // client connection, possibly shared between multiple producers ctx.channel().config().setAutoRead(false); autoReadDisabledRateLimiting = isPublishRateExceeded; + autoReadDisabledPublishBufferLimiting = true; } } void completedSendOperation(boolean isNonPersistentTopic, int msgSize) { - if (--pendingSendRequest == ResumeReadsThreshold | service.decreasePublishBufferSizeAndCheckResumeRead(msgSize)) { + messagePublishBufferSize -= msgSize; + if (--pendingSendRequest == ResumeReadsThreshold) { // Resume reading from socket ctx.channel().config().setAutoRead(true); // triggers channel read if autoRead couldn't trigger it @@ -1703,15 +1710,31 @@ void enableCnxAutoRead() { // we can add check (&& pendingSendRequest < MaxPendingSendRequests) here but then it requires // pendingSendRequest to be volatile and it can be expensive while writing. also this will be called on if // throttling is enable on the topic. so, avoid pendingSendRequest check will be fine. - if (!ctx.channel().config().isAutoRead() && autoReadDisabledRateLimiting) { + if (!ctx.channel().config().isAutoRead() && !autoReadDisabledRateLimiting && !autoReadDisabledPublishBufferLimiting) { // Resume reading from socket if pending-request is not reached to threshold ctx.channel().config().setAutoRead(true); // triggers channel read ctx.read(); + } + } + + @VisibleForTesting + void cancelPublishRateLimiting() { + if (autoReadDisabledRateLimiting) { autoReadDisabledRateLimiting = false; } } +<<<<<<< HEAD +======= + @VisibleForTesting + void cancelPublishBufferLimiting() { + if (autoReadDisabledPublishBufferLimiting) { + autoReadDisabledPublishBufferLimiting = false; + } + } + +>>>>>>> Apply comments private ServerError getErrorCode(CompletableFuture future) { ServerError error = ServerError.UnknownError; try { @@ -1799,4 +1822,18 @@ boolean supportsAuthenticationRefresh() { public String getClientVersion() { return clientVersion; } + + public long getMessagePublishBufferSize() { + return this.messagePublishBufferSize; + } + + @VisibleForTesting + void setMessagePublishBufferSize(long bufferSize) { + this.messagePublishBufferSize = bufferSize; + } + + @VisibleForTesting + void setAutoReadDisabledRateLimiting(boolean isLimiting) { + this.autoReadDisabledRateLimiting = isLimiting; + } } \ No newline at end of file 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 3f78f1c94b18e..26af1c1c5c8bf 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 @@ -201,6 +201,4 @@ void updateRates(NamespaceStats nsStats, NamespaceBundleStats currentBundleStats default Optional getDispatchRateLimiter() { return Optional.empty(); } - - void enableProducerRead(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java index 3307063791753..0c181de787a8b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java @@ -27,6 +27,8 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; /** */ @@ -45,53 +47,122 @@ protected void cleanup() throws Exception { @Test public void testMessagePublishBufferThrottleDisabled() throws Exception { conf.setMaxMessagePublishBufferSizeInMB(-1); + conf.setMessagePublishBufferCheckIntervalInMills(10); super.baseSetup(); - Assert.assertFalse(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(1)); - Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1)); + final String topic = "persistent://prop/ns-abc/testMessagePublishBufferThrottleDisabled"; + Producer producer = pulsarClient.newProducer() + .topic(topic) + .producerName("producer-name") + .create(); + Topic topicRef = pulsar.getBrokerService().getTopicReference(topic).get(); + Assert.assertNotNull(topicRef); + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); + Thread.sleep(20); + Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + List> futures = new ArrayList<>(); + // Make sure the producer can publish succeed. + for (int i = 0; i < 10; i++) { + futures.add(producer.sendAsync(new byte[1024 * 1024])); + } + FutureUtil.waitForAll(futures).get(); + for (CompletableFuture future : futures) { + Assert.assertNotNull(future.get()); + } + Thread.sleep(4); + Assert.assertEquals(pulsar.getBrokerService().getCurrentMessagePublishBufferSize(), 0L); + super.internalCleanup(); } @Test - public void testMessagePublishBufferThrottle() throws Exception { + public void testMessagePublishBufferThrottleEnable() throws Exception { conf.setMaxMessagePublishBufferSizeInMB(1); + conf.setMessagePublishBufferCheckIntervalInMills(2); super.baseSetup(); - Assert.assertFalse(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(512 * 1024)); - Assert.assertTrue(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(512 * 1024)); - Assert.assertTrue(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead( 1024)); - Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1024)); - Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(512 * 1024)); - Assert.assertTrue(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1024)); - Assert.assertFalse(pulsar.getBrokerService().decreasePublishBufferSizeAndCheckResumeRead(1024)); - Assert.assertTrue(pulsar.getBrokerService().increasePublishBufferSizeAndCheckStopRead(514 * 1024)); + Thread.sleep(4); + Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + final String topic = "persistent://prop/ns-abc/testMessagePublishBufferThrottleEnable"; + Producer producer = pulsarClient.newProducer() + .topic(topic) + .producerName("producer-name") + .create(); + Topic topicRef = pulsar.getBrokerService().getTopicReference(topic).get(); + Assert.assertNotNull(topicRef); + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); + Thread.sleep(4); + Assert.assertTrue(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + // The first message can publish success, but the second message should be blocked + producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); + MessageId messageId = null; + try { + messageId = producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); + Assert.fail("should failed, because producer blocked by publish buffer limiting"); + } catch (TimeoutException e) { + // No-op + } + Assert.assertNull(messageId); + + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(0L); + Thread.sleep(4); + + List> futures = new ArrayList<>(); + // Make sure the producer can publish succeed. + for (int i = 0; i < 10; i++) { + futures.add(producer.sendAsync(new byte[1024 * 1024])); + } + FutureUtil.waitForAll(futures).get(); + for (CompletableFuture future : futures) { + Assert.assertNotNull(future.get()); + } + Thread.sleep(4); + Assert.assertEquals(pulsar.getBrokerService().getCurrentMessagePublishBufferSize(), 0L); + super.internalCleanup(); } @Test - public void testCurrentPublishBufferShouldBeZeroWhenComplete() throws Exception { + public void testBlockByPublishRateLimiting() throws Exception { conf.setMaxMessagePublishBufferSizeInMB(1); + conf.setMessagePublishBufferCheckIntervalInMills(2); super.baseSetup(); + Thread.sleep(4); + Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + final String topic = "persistent://prop/ns-abc/testMessagePublishBufferThrottleEnable"; + Producer producer = pulsarClient.newProducer() + .topic(topic) + .producerName("producer-name") + .create(); + Topic topicRef = pulsar.getBrokerService().getTopicReference(topic).get(); + Assert.assertNotNull(topicRef); + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); + producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); - List> futures = new ArrayList<>(); - final int messages = 200; - final int producers = 10; - - List> producerList = new ArrayList<>(); - for (int i = 0; i < producers; i++) { - Producer producer = pulsarClient.newProducer() - .topic("persistent://prop/ns-abc/testCurrentPublishBufferShouldBeZeroWhenComplete") - .enableBatching(false) - .create(); - producerList.add(producer); + Thread.sleep(4); + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setAutoReadDisabledRateLimiting(true); + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(0); + Thread.sleep(4); + Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + MessageId messageId = null; + try { + messageId = producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); + Assert.fail("should failed, because producer blocked by publish buffer limiting"); + } catch (TimeoutException e) { + // No-op } + Assert.assertNull(messageId); - for (Producer producer : producerList) { - for (int j = 0; j < messages; j++) { - futures.add(producer.sendAsync(new byte[1024])); - } - } + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setAutoReadDisabledRateLimiting(false); + ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().enableCnxAutoRead(); + List> futures = new ArrayList<>(); + // Make sure the producer can publish succeed. + for (int i = 0; i < 10; i++) { + futures.add(producer.sendAsync(new byte[1024 * 1024])); + } FutureUtil.waitForAll(futures).get(); - Assert.assertEquals(futures.size(), messages * producers); - Assert.assertEquals(pulsar.getBrokerService().getCurrentMessagePublishBufferSize().get(), 0); - Assert.assertTrue(pulsar.getBrokerService().messagePublishBufferThrottleTimes > 0 - && pulsar.getBrokerService().messagePublishBufferThrottleTimes == pulsar.getBrokerService().messagePublishBufferResumeTimes); + for (CompletableFuture future : futures) { + Assert.assertNotNull(future.get()); + } + Thread.sleep(4); + Assert.assertEquals(pulsar.getBrokerService().getCurrentMessagePublishBufferSize(), 0L); + super.internalCleanup(); } } From 7bbc7a283a68f8e88d9ccfb273c1e87a18a21da9 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Tue, 11 Feb 2020 23:40:30 +0800 Subject: [PATCH 3/6] Clean code --- .../org/apache/pulsar/broker/service/AbstractTopic.java | 4 ++-- .../org/apache/pulsar/broker/service/BrokerService.java | 8 -------- .../java/org/apache/pulsar/broker/service/ServerCnx.java | 5 +---- 3 files changed, 3 insertions(+), 14 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index df6c1d146621f..b1a166ea26f14 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 @@ -292,7 +292,7 @@ public void resetBrokerPublishCountAndEnableReadIfRequired(boolean doneBrokerRes /** * it sets cnx auto-readable if producer's cnx is disabled due to publish-throttling */ - public void enableProducerReadForPublishRateLimiting() { + protected void enableProducerReadForPublishRateLimiting() { if (producers != null) { producers.values().forEach(producer -> { producer.getCnx().cancelPublishRateLimiting(); @@ -301,7 +301,7 @@ public void enableProducerReadForPublishRateLimiting() { } } - public void enableProducerReadForPublishBufferLimiting() { + protected void enableProducerReadForPublishBufferLimiting() { if (producers != null) { producers.values().forEach(producer -> { producer.getCnx().cancelPublishBufferLimiting(); 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 b78b7ae101e36..7f1cac65e4978 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 @@ -54,13 +54,11 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.ExecutionException; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.LongAdder; import java.util.concurrent.locks.ReadWriteLock; @@ -224,10 +222,6 @@ public class BrokerService implements Closeable, ZooKeeperCacheListener= maxMessagePublishBufferSize && !isMessagePublishBufferThreshold) { isMessagePublishBufferThreshold = true; - messagePublishBufferThrottleTimes++; } if (currentMessagePublishBufferSize < resumeProducerReadMessagePublishBufferSize && isMessagePublishBufferThreshold) { isMessagePublishBufferThreshold = false; - messagePublishBufferResumeTimes++; forEachTopic(topic -> ((AbstractTopic) topic).enableProducerReadForPublishBufferLimiting()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index aa7139631053e..bb7613e3e5286 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1684,7 +1684,7 @@ public boolean isWritable() { private void startSendOperation(Producer producer, int msgSize) { messagePublishBufferSize += msgSize; boolean isPublishRateExceeded = producer.getTopic().isPublishRateExceeded(); - if (++pendingSendRequest == MaxPendingSendRequests | isPublishRateExceeded | getBrokerService().isMessagePublishBufferThreshold()) { + if (++pendingSendRequest == MaxPendingSendRequests || isPublishRateExceeded || getBrokerService().isMessagePublishBufferThreshold()) { // When the quota of pending send requests is reached, stop reading from socket to cause backpressure on // client connection, possibly shared between multiple producers ctx.channel().config().setAutoRead(false); @@ -1725,8 +1725,6 @@ void cancelPublishRateLimiting() { } } -<<<<<<< HEAD -======= @VisibleForTesting void cancelPublishBufferLimiting() { if (autoReadDisabledPublishBufferLimiting) { @@ -1734,7 +1732,6 @@ void cancelPublishBufferLimiting() { } } ->>>>>>> Apply comments private ServerError getErrorCode(CompletableFuture future) { ServerError error = ServerError.UnknownError; try { From 917f29b45d5392d71270ab290266df232201c09b Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Wed, 12 Feb 2020 10:46:20 +0800 Subject: [PATCH 4/6] Fix unit test --- .../java/org/apache/pulsar/broker/service/ServerCnx.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index bb7613e3e5286..81696eadd8b42 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1684,11 +1684,15 @@ public boolean isWritable() { private void startSendOperation(Producer producer, int msgSize) { messagePublishBufferSize += msgSize; boolean isPublishRateExceeded = producer.getTopic().isPublishRateExceeded(); - if (++pendingSendRequest == MaxPendingSendRequests || isPublishRateExceeded || getBrokerService().isMessagePublishBufferThreshold()) { + if (++pendingSendRequest == MaxPendingSendRequests || isPublishRateExceeded) { // When the quota of pending send requests is reached, stop reading from socket to cause backpressure on // client connection, possibly shared between multiple producers ctx.channel().config().setAutoRead(false); autoReadDisabledRateLimiting = isPublishRateExceeded; + + } + if (getBrokerService().isMessagePublishBufferThreshold()) { + ctx.channel().config().setAutoRead(false); autoReadDisabledPublishBufferLimiting = true; } } From 47923d17fecde87a9017924828577be66c1e6b43 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Sun, 16 Feb 2020 10:57:07 +0800 Subject: [PATCH 5/6] Apply comments --- conf/broker.conf | 4 +- .../pulsar/broker/ServiceConfiguration.java | 6 +-- .../pulsar/broker/service/BrokerService.java | 43 ++++++++++--------- .../pulsar/broker/service/ServerCnx.java | 2 +- .../MessagePublishBufferThrottleTest.java | 16 +++---- 5 files changed, 36 insertions(+), 35 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 340d13694211d..55838f43f849d 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -329,12 +329,12 @@ replicatedSubscriptionsSnapshotMaxCachedPerSubscription=10 # from the connection. The processing messages means messages are sends to broker # but broker have not send response to client, usually waiting to write to bookies. # It's shared across all the topics running in the same broker. -# Use -1 to disable the memory limitation. Default is 1/5 of direct memory. +# Use -1 to disable the memory limitation. Default is 1/2 of direct memory. maxMessagePublishBufferSizeInMB= # Interval between checks to see if message publish buffer size is exceed the max message publish buffer size # Use 0 or negative number to disable the max publish buffer limiting. -messagePublishBufferCheckIntervalInMills=100 +messagePublishBufferCheckIntervalInMillis=100 ### --- Authentication --- ### # Role names that are treated as "proxy roles". If the broker sees a request with 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 fd32da1df7651..f015a92b66952 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 @@ -611,15 +611,15 @@ public class ServiceConfiguration implements PulsarConfiguration { + " from the connection. The processing messages means messages are sends to broker" + " but broker have not send response to client, usually waiting to write to bookies.\n\n" + " It's shared across all the topics running in the same broker.\n\n" - + " Use -1 to disable the memory limitation. Default is 1/5 of direct memory.\n\n") + + " Use -1 to disable the memory limitation. Default is 1/2 of direct memory.\n\n") private int maxMessagePublishBufferSizeInMB = Math.max(64, - (int) (PlatformDependent.maxDirectMemory() / 5 / (1024 * 1024))); + (int) (PlatformDependent.maxDirectMemory() / 2 / (1024 * 1024))); @FieldContext( category = CATEGORY_SERVER, doc = "Interval between checks to see if message publish buffer size is exceed the max message publish buffer size" ) - private int messagePublishBufferCheckIntervalInMills = 100; + private int messagePublishBufferCheckIntervalInMillis = 100; /**** --- Messaging Protocols --- ****/ 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 7f1cac65e4978..db055b69b9ddc 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 @@ -59,6 +59,7 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.LongAdder; import java.util.concurrent.locks.ReadWriteLock; @@ -218,17 +219,15 @@ public class BrokerService implements Closeable, ZooKeeperCacheListener 0 ? + this.maxMessagePublishBufferBytes = pulsar.getConfiguration().getMaxMessagePublishBufferSizeInMB() > 0 ? pulsar.getConfiguration().getMaxMessagePublishBufferSizeInMB() * 1024 * 1024 : -1; - this.resumeProducerReadMessagePublishBufferSize = this.maxMessagePublishBufferSize / 2; - this.currentMessagePublishBufferSize = 0; + this.resumeProducerReadMessagePublishBufferBytes = this.maxMessagePublishBufferBytes / 2; this.managedLedgerFactory = pulsar.getManagedLedgerFactory(); this.topics = new ConcurrentOpenHashMap<>(); this.replicationClients = new ConcurrentOpenHashMap<>(); @@ -453,8 +452,8 @@ protected void startCompactionMonitor() { } protected void startMessagePublishBufferMonitor() { - int interval = pulsar().getConfiguration().getMessagePublishBufferCheckIntervalInMills(); - if (interval > 0 && maxMessagePublishBufferSize > 0) { + int interval = pulsar().getConfiguration().getMessagePublishBufferCheckIntervalInMillis(); + if (interval > 0 && maxMessagePublishBufferBytes > 0) { messagePublishBufferMonitor.scheduleAtFixedRate(safeRun(this::checkMessagePublishBuffer), interval, interval, TimeUnit.MILLISECONDS); } @@ -2035,15 +2034,15 @@ public Optional getListenPortTls() { } private void checkMessagePublishBuffer() { - currentMessagePublishBufferSize = 0; - foreachProducer(producer -> currentMessagePublishBufferSize += producer.getCnx().getMessagePublishBufferSize()); - if (currentMessagePublishBufferSize >= maxMessagePublishBufferSize - && !isMessagePublishBufferThreshold) { - isMessagePublishBufferThreshold = true; - } - if (currentMessagePublishBufferSize < resumeProducerReadMessagePublishBufferSize - && isMessagePublishBufferThreshold) { - isMessagePublishBufferThreshold = false; + AtomicLong currentMessagePublishBufferBytes = new AtomicLong(); + foreachProducer(producer -> currentMessagePublishBufferBytes.addAndGet(producer.getCnx().getMessagePublishBufferSize())); + if (currentMessagePublishBufferBytes.get() >= maxMessagePublishBufferBytes + && !reachMessagePublishBufferThreshold) { + reachMessagePublishBufferThreshold = true; + } + if (currentMessagePublishBufferBytes.get() < resumeProducerReadMessagePublishBufferBytes + && reachMessagePublishBufferThreshold) { + reachMessagePublishBufferThreshold = false; forEachTopic(topic -> ((AbstractTopic) topic).enableProducerReadForPublishBufferLimiting()); } } @@ -2055,12 +2054,14 @@ private void foreachProducer(Consumer consumer) { }); } - public boolean isMessagePublishBufferThreshold() { - return isMessagePublishBufferThreshold; + public boolean isReachMessagePublishBufferThreshold() { + return reachMessagePublishBufferThreshold; } @VisibleForTesting long getCurrentMessagePublishBufferSize() { - return currentMessagePublishBufferSize; + AtomicLong currentMessagePublishBufferBytes = new AtomicLong(); + foreachProducer(producer -> currentMessagePublishBufferBytes.addAndGet(producer.getCnx().getMessagePublishBufferSize())); + return currentMessagePublishBufferBytes.get(); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 81696eadd8b42..50a5107c05fd9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1691,7 +1691,7 @@ private void startSendOperation(Producer producer, int msgSize) { autoReadDisabledRateLimiting = isPublishRateExceeded; } - if (getBrokerService().isMessagePublishBufferThreshold()) { + if (getBrokerService().isReachMessagePublishBufferThreshold()) { ctx.channel().config().setAutoRead(false); autoReadDisabledPublishBufferLimiting = true; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java index 0c181de787a8b..363e78a6486fe 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java @@ -47,7 +47,7 @@ protected void cleanup() throws Exception { @Test public void testMessagePublishBufferThrottleDisabled() throws Exception { conf.setMaxMessagePublishBufferSizeInMB(-1); - conf.setMessagePublishBufferCheckIntervalInMills(10); + conf.setMessagePublishBufferCheckIntervalInMillis(10); super.baseSetup(); final String topic = "persistent://prop/ns-abc/testMessagePublishBufferThrottleDisabled"; Producer producer = pulsarClient.newProducer() @@ -58,7 +58,7 @@ public void testMessagePublishBufferThrottleDisabled() throws Exception { Assert.assertNotNull(topicRef); ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); Thread.sleep(20); - Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); List> futures = new ArrayList<>(); // Make sure the producer can publish succeed. for (int i = 0; i < 10; i++) { @@ -76,10 +76,10 @@ public void testMessagePublishBufferThrottleDisabled() throws Exception { @Test public void testMessagePublishBufferThrottleEnable() throws Exception { conf.setMaxMessagePublishBufferSizeInMB(1); - conf.setMessagePublishBufferCheckIntervalInMills(2); + conf.setMessagePublishBufferCheckIntervalInMillis(2); super.baseSetup(); Thread.sleep(4); - Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); final String topic = "persistent://prop/ns-abc/testMessagePublishBufferThrottleEnable"; Producer producer = pulsarClient.newProducer() .topic(topic) @@ -89,7 +89,7 @@ public void testMessagePublishBufferThrottleEnable() throws Exception { Assert.assertNotNull(topicRef); ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); Thread.sleep(4); - Assert.assertTrue(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + Assert.assertTrue(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); // The first message can publish success, but the second message should be blocked producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); MessageId messageId = null; @@ -121,10 +121,10 @@ public void testMessagePublishBufferThrottleEnable() throws Exception { @Test public void testBlockByPublishRateLimiting() throws Exception { conf.setMaxMessagePublishBufferSizeInMB(1); - conf.setMessagePublishBufferCheckIntervalInMills(2); + conf.setMessagePublishBufferCheckIntervalInMillis(2); super.baseSetup(); Thread.sleep(4); - Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); final String topic = "persistent://prop/ns-abc/testMessagePublishBufferThrottleEnable"; Producer producer = pulsarClient.newProducer() .topic(topic) @@ -139,7 +139,7 @@ public void testBlockByPublishRateLimiting() throws Exception { ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setAutoReadDisabledRateLimiting(true); ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(0); Thread.sleep(4); - Assert.assertFalse(pulsar.getBrokerService().isMessagePublishBufferThreshold()); + Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); MessageId messageId = null; try { messageId = producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); From f752a6ce54cbd58dfc4cb873ccca3a6fdcdc7890 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Sun, 16 Feb 2020 21:48:27 +0800 Subject: [PATCH 6/6] Fix unit tests. --- .../broker/service/MessagePublishBufferThrottleTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java index 363e78a6486fe..5397725a99c67 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java @@ -68,8 +68,8 @@ public void testMessagePublishBufferThrottleDisabled() throws Exception { for (CompletableFuture future : futures) { Assert.assertNotNull(future.get()); } - Thread.sleep(4); - Assert.assertEquals(pulsar.getBrokerService().getCurrentMessagePublishBufferSize(), 0L); + Thread.sleep(20); + Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); super.internalCleanup(); }