diff --git a/conf/broker.conf b/conf/broker.conf index 4d217a3e1981f..1c0242f583dbc 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -73,8 +73,8 @@ backlogQuotaCheckEnabled=true # How often to check for topics that have reached the quota 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 60d5a9cfa3920..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 @@ -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/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; } 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); }