From 2a2e59e97efc6169e9fcdc9243baa9d1eaeb2662 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 21 May 2019 14:49:59 +0800 Subject: [PATCH 1/5] [broker-conf] Change broker level backlog quota retention policy to consumer_backlog_eviction. --- conf/broker.conf | 4 ++-- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 4d217a3e1981f..4db8429dc8435 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -76,11 +76,11 @@ backlogQuotaCheckIntervalInSeconds=60 # Default per-topic backlog quota limit backlogQuotaDefaultLimitGB=10 -# Default backlog quota retention policy. Default is producer_request_hold +# Default backlog quota retention policy. Default is consumer_backlog_eviction # 'producer_request_hold' Policy which holds producer's send request until the resource becomes available (or holding times out) # 'producer_exception' Policy which throws javax.jms.ResourceAllocationException to the producer # 'consumer_backlog_eviction' Policy which evicts the oldest message from the slowest consumer's backlog -backlogQuotaDefaultRetentionPolicy=producer_request_hold +backlogQuotaDefaultRetentionPolicy=consumer_backlog_eviction # Default ttl for namespaces if ttl is not already configured at namespace policies. (disable default-ttl with value 0) ttlDurationDefaultInSeconds=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 60d5a9cfa3920..4d5e92271a83e 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 @@ -215,13 +215,13 @@ public class ServiceConfiguration implements PulsarConfiguration { private long backlogQuotaDefaultLimitGB = 50; @FieldContext( category = CATEGORY_POLICIES, - doc = "Default backlog quota retention policy. Default is producer_request_hold\n\n" + doc = "Default backlog quota retention policy. Default is consumer_backlog_eviction\n\n" + "'producer_request_hold' Policy which holds producer's send request until the" + "resource becomes available (or holding times out)\n" + "'producer_exception' Policy which throws javax.jms.ResourceAllocationException to the producer\n" + "'consumer_backlog_eviction' Policy which evicts the oldest message from the slowest consumer's backlog" ) - private BacklogQuota.RetentionPolicy backlogQuotaDefaultRetentionPolicy = BacklogQuota.RetentionPolicy.producer_request_hold; + private BacklogQuota.RetentionPolicy backlogQuotaDefaultRetentionPolicy = BacklogQuota.RetentionPolicy.consumer_backlog_eviction; @FieldContext( category = CATEGORY_POLICIES, doc = "Default ttl for namespaces if ttl is not already configured at namespace policies. " From 4cce583c3c122baf3c4cbd5f22cb113eb22cd5cc Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 28 May 2019 10:04:04 +0800 Subject: [PATCH 2/5] disable backlog quota check by default. --- conf/broker.conf | 10 +++++----- .../org/apache/pulsar/broker/ServiceConfiguration.java | 6 +++--- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 4db8429dc8435..f6ca07de4bef4 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -68,19 +68,19 @@ zooKeeperOperationTimeoutSeconds=30 brokerShutdownTimeoutMs=60000 # Enable backlog quota check. Enforces action on topic when the quota is reached -backlogQuotaCheckEnabled=true +backlogQuotaCheckEnabled=false # How often to check for topics that have reached the quota -backlogQuotaCheckIntervalInSeconds=60 +# backlogQuotaCheckIntervalInSeconds=60 # Default per-topic backlog quota limit -backlogQuotaDefaultLimitGB=10 +# backlogQuotaDefaultLimitGB=10 -# Default backlog quota retention policy. Default is consumer_backlog_eviction +# Default backlog quota retention policy. Default is producer_request_hold # 'producer_request_hold' Policy which holds producer's send request until the resource becomes available (or holding times out) # 'producer_exception' Policy which throws javax.jms.ResourceAllocationException to the producer # 'consumer_backlog_eviction' Policy which evicts the oldest message from the slowest consumer's backlog -backlogQuotaDefaultRetentionPolicy=consumer_backlog_eviction +backlogQuotaDefaultRetentionPolicy=producer_request_hold # Default ttl for namespaces if ttl is not already configured at namespace policies. (disable default-ttl with value 0) ttlDurationDefaultInSeconds=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 4d5e92271a83e..4533b6a44bfbc 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 @@ -201,7 +201,7 @@ public class ServiceConfiguration implements PulsarConfiguration { category = CATEGORY_POLICIES, doc = "Enable backlog quota check. Enforces actions on topic when the quota is reached" ) - private boolean backlogQuotaCheckEnabled = true; + private boolean backlogQuotaCheckEnabled = false; @FieldContext( category = CATEGORY_POLICIES, doc = "How often to check for topics that have reached the quota." @@ -215,13 +215,13 @@ public class ServiceConfiguration implements PulsarConfiguration { private long backlogQuotaDefaultLimitGB = 50; @FieldContext( category = CATEGORY_POLICIES, - doc = "Default backlog quota retention policy. Default is consumer_backlog_eviction\n\n" + doc = "Default backlog quota retention policy. Default is producer_request_hold\n\n" + "'producer_request_hold' Policy which holds producer's send request until the" + "resource becomes available (or holding times out)\n" + "'producer_exception' Policy which throws javax.jms.ResourceAllocationException to the producer\n" + "'consumer_backlog_eviction' Policy which evicts the oldest message from the slowest consumer's backlog" ) - private BacklogQuota.RetentionPolicy backlogQuotaDefaultRetentionPolicy = BacklogQuota.RetentionPolicy.consumer_backlog_eviction; + private BacklogQuota.RetentionPolicy backlogQuotaDefaultRetentionPolicy = BacklogQuota.RetentionPolicy.producer_request_hold; @FieldContext( category = CATEGORY_POLICIES, doc = "Default ttl for namespaces if ttl is not already configured at namespace policies. " From a25dc44cb9a7dab2f9b68acc21c2e9a2ddc3a9c3 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 28 May 2019 15:30:46 +0800 Subject: [PATCH 3/5] Fix unit tests. --- .../test/java/org/apache/pulsar/PulsarBrokerStarterTest.java | 2 +- .../apache/pulsar/broker/service/BacklogQuotaManagerTest.java | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java index 5ccea59f3301e..1ff53d6b3cd26 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java @@ -115,7 +115,7 @@ public void testLoadConfig() throws SecurityException, NoSuchMethodException, IO Assert.assertEquals(serviceConfig.getManagedLedgerCursorMaxEntriesPerLedger(), 50); Assert.assertTrue(serviceConfig.isClientLibraryVersionCheckEnabled()); Assert.assertEquals(serviceConfig.getManagedLedgerMinLedgerRolloverTimeMinutes(), 34); - Assert.assertEquals(serviceConfig.isBacklogQuotaCheckEnabled(), true); + Assert.assertEquals(serviceConfig.isBacklogQuotaCheckEnabled(), false); Assert.assertEquals(serviceConfig.getManagedLedgerDefaultMarkDeleteRateLimit(), 5.0); Assert.assertEquals(serviceConfig.getReplicationProducerQueueSize(), 50); Assert.assertEquals(serviceConfig.isReplicationMetricsEnabled(), false); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java index fdbc0fc851982..2885fd2f31e20 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java @@ -85,6 +85,7 @@ void setup() throws Exception { config.setBrokerServicePort(Optional.ofNullable(BROKER_SERVICE_PORT)); config.setAuthorizationEnabled(false); config.setAuthenticationEnabled(false); + config.setBacklogQuotaCheckEnabled(true); config.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); config.setManagedLedgerMaxEntriesPerLedger(5); config.setManagedLedgerMinLedgerRolloverTimeMinutes(0); From fea3e8118632fb132c6da20834ec1ccf1ff7be92 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Wed, 29 May 2019 12:04:54 +0800 Subject: [PATCH 4/5] SET default backlogQuotaDefaultLimitGB=-1 --- conf/broker.conf | 8 ++++---- .../org/apache/pulsar/broker/ServiceConfiguration.java | 7 ++++--- .../org/apache/pulsar/broker/service/BrokerService.java | 3 +++ .../java/org/apache/pulsar/PulsarBrokerStarterTest.java | 2 +- .../pulsar/broker/service/BacklogQuotaManagerTest.java | 1 - 5 files changed, 12 insertions(+), 9 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index f6ca07de4bef4..1c0242f583dbc 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -68,13 +68,13 @@ zooKeeperOperationTimeoutSeconds=30 brokerShutdownTimeoutMs=60000 # Enable backlog quota check. Enforces action on topic when the quota is reached -backlogQuotaCheckEnabled=false +backlogQuotaCheckEnabled=true # How often to check for topics that have reached the quota -# backlogQuotaCheckIntervalInSeconds=60 +backlogQuotaCheckIntervalInSeconds=60 -# Default per-topic backlog quota limit -# backlogQuotaDefaultLimitGB=10 +# Default per-topic backlog quota limit, less than 0 means no limitation. default is -1. +backlogQuotaDefaultLimitGB=-1 # Default backlog quota retention policy. Default is producer_request_hold # 'producer_request_hold' Policy which holds producer's send request until the resource becomes available (or holding times out) 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 4533b6a44bfbc..873fa66e2dd18 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 @@ -201,7 +201,7 @@ public class ServiceConfiguration implements PulsarConfiguration { category = CATEGORY_POLICIES, doc = "Enable backlog quota check. Enforces actions on topic when the quota is reached" ) - private boolean backlogQuotaCheckEnabled = false; + private boolean backlogQuotaCheckEnabled = true; @FieldContext( category = CATEGORY_POLICIES, doc = "How often to check for topics that have reached the quota." @@ -210,9 +210,10 @@ public class ServiceConfiguration implements PulsarConfiguration { private int backlogQuotaCheckIntervalInSeconds = 60; @FieldContext( category = CATEGORY_POLICIES, - doc = "Default per-topic backlog quota limit. Increase it if you want to allow larger msg backlog" + doc = "Default per-topic backlog quota limit, less than 0 means no limitation. default is -1." + + " Increase it if you want to allow larger msg backlog" ) - private long backlogQuotaDefaultLimitGB = 50; + private long backlogQuotaDefaultLimitGB = -1; @FieldContext( category = CATEGORY_POLICIES, doc = "Default backlog quota retention policy. Default is producer_request_hold\n\n" 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 b26a6c1861290..9618bc1a04c29 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 @@ -916,6 +916,9 @@ public BacklogQuotaManager getBacklogQuotaManager() { public boolean isBacklogExceeded(PersistentTopic topic) { TopicName topicName = TopicName.get(topic.getName()); long backlogQuotaLimitInBytes = getBacklogQuotaManager().getBacklogQuotaLimit(topicName.getNamespace()); + if (backlogQuotaLimitInBytes < 0) { + return false; + } if (log.isDebugEnabled()) { log.debug("[{}] - backlog quota limit = [{}]", topic.getName(), backlogQuotaLimitInBytes); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java index 1ff53d6b3cd26..5ccea59f3301e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/PulsarBrokerStarterTest.java @@ -115,7 +115,7 @@ public void testLoadConfig() throws SecurityException, NoSuchMethodException, IO Assert.assertEquals(serviceConfig.getManagedLedgerCursorMaxEntriesPerLedger(), 50); Assert.assertTrue(serviceConfig.isClientLibraryVersionCheckEnabled()); Assert.assertEquals(serviceConfig.getManagedLedgerMinLedgerRolloverTimeMinutes(), 34); - Assert.assertEquals(serviceConfig.isBacklogQuotaCheckEnabled(), false); + Assert.assertEquals(serviceConfig.isBacklogQuotaCheckEnabled(), true); Assert.assertEquals(serviceConfig.getManagedLedgerDefaultMarkDeleteRateLimit(), 5.0); Assert.assertEquals(serviceConfig.getReplicationProducerQueueSize(), 50); Assert.assertEquals(serviceConfig.isReplicationMetricsEnabled(), false); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java index 2885fd2f31e20..fdbc0fc851982 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java @@ -85,7 +85,6 @@ void setup() throws Exception { config.setBrokerServicePort(Optional.ofNullable(BROKER_SERVICE_PORT)); config.setAuthorizationEnabled(false); config.setAuthenticationEnabled(false); - config.setBacklogQuotaCheckEnabled(true); config.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); config.setManagedLedgerMaxEntriesPerLedger(5); config.setManagedLedgerMinLedgerRolloverTimeMinutes(0); From 7b3373616b441239f4a5739ab2d2dd6367fc5edd Mon Sep 17 00:00:00 2001 From: lipenghui Date: Wed, 29 May 2019 12:10:44 +0800 Subject: [PATCH 5/5] ADD disable backlog quota check for checkQuotas(). --- .../org/apache/pulsar/broker/admin/impl/NamespacesBase.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index eab5cd9d9c3eb..7a6589d287e5f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -1416,6 +1416,9 @@ private boolean checkQuotas(Policies policies, RetentionPolicies retention) { if (quota == null) { quota = pulsar().getBrokerService().getBacklogQuotaManager().getDefaultQuota(); } + if (quota.getLimit() < 0 && (retention.getRetentionSizeInMB() > 0 || retention.getRetentionTimeInMinutes() > 0)) { + return false; + } if (quota.getLimit() >= ((long) retention.getRetentionSizeInMB() * 1024 * 1024)) { return false; }