From a754347e7085cc170adf095bc925183de6cb195b Mon Sep 17 00:00:00 2001 From: Qiang Huang Date: Thu, 16 Dec 2021 00:03:39 +0800 Subject: [PATCH 1/2] Enable client memory limit controller by default --- .../java/org/apache/pulsar/client/api/ProducerBuilder.java | 2 +- .../pulsar/client/impl/conf/ClientConfigurationData.java | 4 ++-- .../pulsar/client/impl/conf/ProducerConfigurationData.java | 2 +- 3 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java index fd91101ee4986..0b93ce27bfc81 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java @@ -172,7 +172,7 @@ public interface ProducerBuilder extends Cloneable { * the client application. Until, the producer gets a successful acknowledgment back from the broker, * it will keep in memory (direct memory pool) all the messages in the pending queue. * - *

Default is 1000. + *

Default is 0, disable the pending messages check. * * @param maxPendingMessages * the max size of the pending messages queue for the producer diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java index e4f86d383aaca..b2c62ca10b40d 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java @@ -270,9 +270,9 @@ public class ClientConfigurationData implements Serializable, Cloneable { @ApiModelProperty( name = "memoryLimitBytes", - value = "Limit of client memory usage (in byte)." + value = "Limit of client memory usage (in byte). The 64M default can guarantee a high producer throughput." ) - private long memoryLimitBytes = 0; + private long memoryLimitBytes = 64 * 1024 * 1024; @ApiModelProperty( name = "proxyServiceUrl", diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java index 18a70a6e8ccc7..f94bf5b057a30 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java @@ -53,7 +53,7 @@ public class ProducerConfigurationData implements Serializable, Cloneable { private static final long serialVersionUID = 1L; public static final int DEFAULT_BATCHING_MAX_MESSAGES = 1000; - public static final int DEFAULT_MAX_PENDING_MESSAGES = 1000; + public static final int DEFAULT_MAX_PENDING_MESSAGES = 0; public static final int DEFAULT_MAX_PENDING_MESSAGES_ACROSS_PARTITIONS = 50000; private String topicName = null; From f4e6d2d19df1da2168253ec39ffb853d6dac655c Mon Sep 17 00:00:00 2001 From: Qiang Huang Date: Thu, 16 Dec 2021 11:49:31 +0800 Subject: [PATCH 2/2] Deprecate maxPendingMessagesAcrossPartitions --- .../java/org/apache/pulsar/client/api/ProducerBuilder.java | 3 ++- .../org/apache/pulsar/client/impl/ProducerBuilderImpl.java | 1 + .../pulsar/client/impl/conf/ProducerConfigurationData.java | 2 +- .../apache/pulsar/client/impl/ProducerBuilderImplTest.java | 4 ++-- 4 files changed, 6 insertions(+), 4 deletions(-) diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java index 0b93ce27bfc81..53a6f26170c78 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java @@ -188,7 +188,7 @@ public interface ProducerBuilder extends Cloneable { * The purpose of this setting is to have an upper-limit on the number * of pending messages when publishing on a partitioned topic. * - *

Default is 50000. + *

Default is 0, disable the pending messages across partitions check. * *

If publishing at high rate over a topic with many partitions (especially when publishing messages without a * partitioning key), it might be beneficial to increase this parameter to allow for more pipelining within the @@ -198,6 +198,7 @@ public interface ProducerBuilder extends Cloneable { * max pending messages across all the partitions * @return the producer builder instance */ + @Deprecated ProducerBuilder maxPendingMessagesAcrossPartitions(int maxPendingMessagesAcrossPartitions); /** diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java index 497b82a78c7f8..90a2fc8666a1e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java @@ -145,6 +145,7 @@ public ProducerBuilder maxPendingMessages(int maxPendingMessages) { return this; } + @Deprecated @Override public ProducerBuilder maxPendingMessagesAcrossPartitions(int maxPendingMessagesAcrossPartitions) { conf.setMaxPendingMessagesAcrossPartitions(maxPendingMessagesAcrossPartitions); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java index f94bf5b057a30..08fc096e66029 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java @@ -54,7 +54,7 @@ public class ProducerConfigurationData implements Serializable, Cloneable { public static final int DEFAULT_BATCHING_MAX_MESSAGES = 1000; public static final int DEFAULT_MAX_PENDING_MESSAGES = 0; - public static final int DEFAULT_MAX_PENDING_MESSAGES_ACROSS_PARTITIONS = 50000; + public static final int DEFAULT_MAX_PENDING_MESSAGES_ACROSS_PARTITIONS = 0; private String topicName = null; private String producerName = null; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ProducerBuilderImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ProducerBuilderImplTest.java index 21708c25f7537..9a53bc6146037 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ProducerBuilderImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ProducerBuilderImplTest.java @@ -346,12 +346,12 @@ public void testProducerBuilderImplWhenSendTimeoutPropertyIsNegative() { @Test(expectedExceptions = IllegalArgumentException.class) public void testProducerBuilderImplWhenMaxPendingMessagesAcrossPartitionsPropertyIsInvalid() { - producerBuilderImpl.maxPendingMessagesAcrossPartitions(999); + producerBuilderImpl.maxPendingMessagesAcrossPartitions(-1); } @Test(expectedExceptions = IllegalArgumentException.class, expectedExceptionsMessageRegExp = "maxPendingMessagesAcrossPartitions needs to be >= maxPendingMessages") public void testProducerBuilderImplWhenMaxPendingMessagesAcrossPartitionsPropertyIsInvalidErrorMessages() { - producerBuilderImpl.maxPendingMessagesAcrossPartitions(999); + producerBuilderImpl.maxPendingMessagesAcrossPartitions(-1); } @Test