diff --git a/pulsar-broker/src/main/java/com/yahoo/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/com/yahoo/pulsar/broker/service/persistent/PersistentReplicator.java index 160fe9d48b464..b8d8e7c6f78b9 100644 --- a/pulsar-broker/src/main/java/com/yahoo/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/com/yahoo/pulsar/broker/service/persistent/PersistentReplicator.java @@ -85,7 +85,7 @@ public class PersistentReplicator implements ReadEntriesCallback, DeleteCallback private final Rate msgExpired = new Rate(); private static final ProducerConfiguration producerConfiguration = new ProducerConfiguration().setSendTimeout(0, - TimeUnit.SECONDS); + TimeUnit.SECONDS).setBlockIfQueueFull(true); private final Backoff backOff = new Backoff(100, TimeUnit.MILLISECONDS, 1, TimeUnit.MINUTES); private int messageTTLInSeconds = 0; diff --git a/pulsar-client-cpp/lib/ProducerConfigurationImpl.h b/pulsar-client-cpp/lib/ProducerConfigurationImpl.h index 0d2448d22a28a..37c2056d0928f 100644 --- a/pulsar-client-cpp/lib/ProducerConfigurationImpl.h +++ b/pulsar-client-cpp/lib/ProducerConfigurationImpl.h @@ -38,7 +38,7 @@ struct ProducerConfigurationImpl { compressionType(CompressionNone), maxPendingMessages(30000), routingMode(ProducerConfiguration::UseSinglePartition), - blockIfQueueFull(true), + blockIfQueueFull(false), batchingEnabled(false), batchingMaxMessages(1000), batchingMaxAllowedSizeInBytes(128 * 1024), // 128 KB diff --git a/pulsar-client-cpp/perf/PerfProducer.cc b/pulsar-client-cpp/perf/PerfProducer.cc index bd65f4f185d8d..96fdd3bbfa00c 100644 --- a/pulsar-client-cpp/perf/PerfProducer.cc +++ b/pulsar-client-cpp/perf/PerfProducer.cc @@ -286,6 +286,9 @@ int main(int argc, char** argv) { producerConf.setBatchingMaxAllowedSizeInBytes(args.batchingMaxAllowedSizeInBytes); producerConf.setBatchingMaxPublishDelayMs(args.batchingMaxPublishDelayMs); } + + // Block if queue is full else we will start seeing errors in sendAsync + producerConf.setBlockIfQueueFull(true); pulsar::ClientConfiguration conf; conf.setUseTls(args.isUseTls); conf.setTlsAllowInsecureConnection(args.isTlsAllowInsecureConnection); diff --git a/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/Producer.java b/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/Producer.java index a3fa7b82fe622..1e6629e47cb71 100644 --- a/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/Producer.java +++ b/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/Producer.java @@ -50,8 +50,7 @@ public interface Producer extends Closeable { /** * Send a message asynchronously *

- * When the producer queue is full, by default this method will block until there will be space available in - * the queue. + * When the producer queue is full, by default this method will complete the future with an exception {@link PulsarClientException#ProducerQueueIsFullError} *

* See {@link ProducerConfiguration#setMaxPendingMessages} to configure the producer queue size and * {@link ProducerConfiguration#setBlockIfQueueFull(boolean)} to change the blocking behavior. @@ -90,8 +89,7 @@ public interface Producer extends Closeable { * }); * *

- * When the producer queue is full, by default this method will block until there will be space available in - * the queue. + * When the producer queue is full, by default this method will complete the future with an exception {@link PulsarClientException#ProducerQueueIsFullError} *

* See {@link ProducerConfiguration#setMaxPendingMessages} to configure the producer queue size and * {@link ProducerConfiguration#setBlockIfQueueFull(boolean)} to change the blocking behavior. diff --git a/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/ProducerConfiguration.java b/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/ProducerConfiguration.java index 20bf0c3ed6a8b..371ef3664151c 100644 --- a/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/ProducerConfiguration.java +++ b/pulsar-client/src/main/java/com/yahoo/pulsar/client/api/ProducerConfiguration.java @@ -36,8 +36,8 @@ public class ProducerConfiguration implements Serializable { */ private static final long serialVersionUID = 1L; private long sendTimeoutMs = 30000; + private boolean blockIfQueueFull = false; private int maxPendingMessages = 30000; - private boolean blockIfQueueFull = true; private MessageRoutingMode messageRouteMode = MessageRoutingMode.SinglePartition; private MessageRouter customMessageRouter = null; private long batchingMaxPublishDelayMs = 10; @@ -83,8 +83,8 @@ public int getMaxPendingMessages() { /** * Set the max size of the queue holding the messages pending to receive an acknowledgment from the broker. *

- * When the queue is full, by default, all calls to {@link Producer#send} and {@link Producer#sendAsync} and send - * will block the calling thread. Use {@link #setBlockIfQueueFull} to change the blocking behavior. + * When the queue is full, by default, all calls to {@link Producer#send} and {@link Producer#sendAsync} + * will fail unless blockIfQueueFull is set to true. Use {@link #setBlockIfQueueFull} to change the blocking behavior. * * @param maxPendingMessages * @return @@ -108,7 +108,7 @@ public boolean getBlockIfQueueFull() { * Set whether the {@link Producer#send} and {@link Producer#sendAsync} operations should block when the outgoing * message queue is full. *

- * Default is true. If set to false, send operations will immediately fail with + * Default is false. If set to false, send operations will immediately fail with * {@link ProducerQueueIsFullError} when there is no space left in pending queue. * * @param blockIfQueueFull diff --git a/pulsar-testclient/src/main/java/com/yahoo/pulsar/testclient/PerformanceProducer.java b/pulsar-testclient/src/main/java/com/yahoo/pulsar/testclient/PerformanceProducer.java index 3679ff2e0e981..2540b171d0a27 100644 --- a/pulsar-testclient/src/main/java/com/yahoo/pulsar/testclient/PerformanceProducer.java +++ b/pulsar-testclient/src/main/java/com/yahoo/pulsar/testclient/PerformanceProducer.java @@ -233,6 +233,9 @@ public static void main(String[] args) throws Exception { producerConf.setMaxPendingMessages(arguments.msgRate); } + // Block if queue is full else we will start seeing errors in sendAsync + producerConf.setBlockIfQueueFull(true); + for (int i = 0; i < arguments.numTopics; i++) { String topic = (arguments.numTopics == 1) ? prefixTopicName : String.format("%s-%d", prefixTopicName, i); log.info("Adding {} publishers on destination {}", arguments.numProducers, topic);