From 5dd74cf71fa30cfa3ce0478c7fb91ee25ef66b1e Mon Sep 17 00:00:00 2001 From: chenhang Date: Tue, 29 Jun 2021 13:09:17 +0800 Subject: [PATCH 1/3] add parameter to control publish message check policy frequency --- conf/broker.conf | 4 ++++ conf/standalone.conf | 4 ++++ .../terraform-ansible/templates/broker.conf | 4 ++++ .../pulsar/broker/ServiceConfiguration.java | 6 +++++ .../pulsar/broker/service/AbstractTopic.java | 24 ++++++++++++------- .../broker/admin/TopicPoliciesTest.java | 12 ++++++++-- 6 files changed, 44 insertions(+), 10 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 0efc6569abfb5..685b4ba15ac28 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -487,6 +487,10 @@ maxMessagePublishBufferSizeInMB= # Use 0 or negative number to disable the check retentionCheckIntervalInSeconds=120 +# Check between intervals to see if max message size in topic policies has been updated. +# Default is 60s +maxMessageSizeCheckIntervalInSeconds=60 + # Max number of partitions per partitioned topic # Use 0 or negative number to disable the check maxNumPartitionsPerPartitionedTopic=0 diff --git a/conf/standalone.conf b/conf/standalone.conf index 16a877dc62899..2ef53d6169aee 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -118,6 +118,10 @@ maxPendingPublishRequestsPerConnection=1000 # How frequently to proactively check and purge expired messages messageExpiryCheckIntervalInMinutes=5 +# Check between intervals to see if max message size in topic policies has been updated. +# Default is 60s +maxMessageSizeCheckIntervalInSeconds=60 + # How long to delay rewinding cursor and dispatching messages when active consumer is changed activeConsumerFailoverDelayTimeMillis=1000 diff --git a/deployment/terraform-ansible/templates/broker.conf b/deployment/terraform-ansible/templates/broker.conf index b17684dc38f2e..1e0e43fc943c1 100644 --- a/deployment/terraform-ansible/templates/broker.conf +++ b/deployment/terraform-ansible/templates/broker.conf @@ -412,6 +412,10 @@ messagePublishBufferCheckIntervalInMillis=100 # Use 0 or negative number to disable the check retentionCheckIntervalInSeconds=120 +# Check between intervals to see if max message size in topic policies has been updated. +# Default is 60s +maxMessageSizeCheckIntervalInSeconds=60 + # Max number of partitions per partitioned topic # Use 0 or negative number to disable the check maxNumPartitionsPerPartitionedTopic=0 diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index b41f76fc6a0d0..9b5f9399e16c8 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 @@ -916,6 +916,12 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private int retentionCheckIntervalInSeconds = 120; + @FieldContext( + category = CATEGORY_SERVER, + doc = "Check between intervals to see if max message size of topic policy has updated. default is 60s" + ) + private int maxMessageSizeCheckIntervalInSeconds = 60; + @FieldContext( category = CATEGORY_SERVER, doc = "The number of partitions per partitioned topic.\n" 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 e2133f6345fb4..33c973cbcc1e5 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 @@ -120,6 +120,10 @@ public abstract class AbstractTopic implements Topic { AtomicLongFieldUpdater.newUpdater(AbstractTopic.class, "usageCount"); private volatile long usageCount = 0; + private volatile int topicMaxMessageSize = 0; + private volatile long lastTopicMaxMessageSizeCheckTimeStamp = 0; + private final long topicMaxMessageSizeCheckIntervalMs; + public AbstractTopic(String topic, BrokerService brokerService) { this.topic = topic; this.brokerService = brokerService; @@ -132,6 +136,8 @@ public AbstractTopic(String topic, BrokerService brokerService) { .getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds()); this.inactiveTopicPolicies.setInactiveTopicDeleteMode(brokerService.pulsar().getConfiguration() .getBrokerDeleteInactiveTopicsMode()); + this.topicMaxMessageSizeCheckIntervalMs = TimeUnit.SECONDS.toMillis(brokerService.pulsar().getConfiguration() + .getMaxMessageSizeCheckIntervalInSeconds()); this.lastActive = System.nanoTime(); Policies policies = null; try { @@ -853,14 +859,16 @@ protected int getWaitingProducersCount() { } protected boolean isExceedMaximumMessageSize(int size) { - return getTopicPolicies() - .map(TopicPolicies::getMaxMessageSize) - .map(maxMessageSize -> { - if (maxMessageSize == 0) { - return false; - } - return size > maxMessageSize; - }).orElse(false); + if (lastTopicMaxMessageSizeCheckTimeStamp + topicMaxMessageSizeCheckIntervalMs < System.currentTimeMillis()) { + // refresh topicMaxMessageSize from topic policies + topicMaxMessageSize = getTopicPolicies().map(TopicPolicies::getMaxMessageSize).orElse(0); + lastTopicMaxMessageSizeCheckTimeStamp = System.currentTimeMillis(); + } + + if (topicMaxMessageSize == 0) { + return false; + } + return size > topicMaxMessageSize; } /** diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java index 3934d3b342c4e..b38631ccc1f0c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java @@ -56,6 +56,7 @@ import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.apache.pulsar.common.policies.data.TopicStats; import org.awaitility.Awaitility; +import static org.mockito.Mockito.doReturn; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -100,6 +101,7 @@ protected void setup() throws Exception { this.conf.setSystemTopicEnabled(true); this.conf.setTopicLevelPoliciesEnabled(true); this.conf.setDefaultNumberOfNamespaceBundles(1); + this.conf.setMaxMessageSizeCheckIntervalInSeconds(1); super.internalSetup(); admin.clusters().createCluster("test", ClusterData.builder().serviceUrl(pulsar.getWebServiceAddress()).build()); @@ -1781,8 +1783,14 @@ private void doTestTopicMaxMessageSize(boolean isPartitioned) throws Exception { assertEquals(e.getStatusCode(), 412); } - MessageId messageId = producer.send(new byte[1024]); - assertNotNull(messageId); + Awaitility.await().untilAsserted(() -> { + try { + MessageId messageId = producer.send(new byte[1024]); + assertNotNull(messageId); + } catch (PulsarClientException e) { + fail("failed to send message"); + } + }); producer.close(); } From 40300df08e4efc48843d5294a69ab4c827dcab0c Mon Sep 17 00:00:00 2001 From: chenhang Date: Tue, 29 Jun 2021 13:15:38 +0800 Subject: [PATCH 2/3] format code --- .../java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java index b38631ccc1f0c..fd664aeb164ae 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java @@ -56,7 +56,6 @@ import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.apache.pulsar.common.policies.data.TopicStats; import org.awaitility.Awaitility; -import static org.mockito.Mockito.doReturn; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; From caae321d6a54c54e3326b6396b3ac83566dfae95 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Thu, 1 Jul 2021 16:27:55 +0800 Subject: [PATCH 3/3] Update conf/broker.conf Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- conf/broker.conf | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 685b4ba15ac28..47795a56bca30 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -487,8 +487,8 @@ maxMessagePublishBufferSizeInMB= # Use 0 or negative number to disable the check retentionCheckIntervalInSeconds=120 -# Check between intervals to see if max message size in topic policies has been updated. -# Default is 60s +# Control the frequency of checking the max message size in the topic policy. +# The default interval is 60 seconds. maxMessageSizeCheckIntervalInSeconds=60 # Max number of partitions per partitioned topic