From 63bcc969c049a43914115490a705c730eb6323a1 Mon Sep 17 00:00:00 2001 From: Baodi Shi Date: Tue, 12 Nov 2024 22:04:16 +0800 Subject: [PATCH] [fix][client] The partitionedProducer maxPendingMessages always is 0 --- .../service/PersistentTopicE2ETest.java | 44 +++++++++++++++++++ .../client/impl/PartitionedProducerImpl.java | 15 +++++-- .../impl/PartitionedProducerImplTest.java | 38 ++++++++++++++++ 3 files changed, 94 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java index 4d79e7ccdf0d1..8e3b920e002b6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java @@ -75,6 +75,7 @@ import org.apache.pulsar.client.impl.ConsumerImpl; import org.apache.pulsar.client.impl.LookupService; import org.apache.pulsar.client.impl.MessageIdImpl; +import org.apache.pulsar.client.impl.PartitionedProducerImpl; import org.apache.pulsar.client.impl.ProducerImpl; import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.TypedMessageBuilderImpl; @@ -1340,6 +1341,49 @@ public void testProducerQueueFullBlocking() throws Exception { setup(); } + @Test + public void testProducerQueueFullBlockingWithPartitionedTopic() throws Exception { + final String topicName = "persistent://prop/ns-abc/topic-xyzx2"; + admin.topics().createPartitionedTopic(topicName, 2); + + @Cleanup + PulsarClient client = PulsarClient.builder().serviceUrl(brokerUrl.toString()).build(); + + // 1. Producer connect + PartitionedProducerImpl producer = (PartitionedProducerImpl) client.newProducer() + .topic(topicName) + .maxPendingMessages(1) + .blockIfQueueFull(true) + .sendTimeout(1, TimeUnit.SECONDS) + .enableBatching(false) + .messageRoutingMode(MessageRoutingMode.SinglePartition) + .create(); + + // 2. Stop broker + cleanup(); + + // 2. producer publish messages + long startTime = System.nanoTime(); + producer.sendAsync("msg".getBytes()); + + // Verify thread was not blocked + long delayNs = System.nanoTime() - startTime; + assertTrue(delayNs < TimeUnit.SECONDS.toNanos(1)); + + // Next send operation must block, until all the messages in the queue expire + startTime = System.nanoTime(); + producer.sendAsync("msg".getBytes()); + delayNs = System.nanoTime() - startTime; + assertTrue(delayNs > TimeUnit.MILLISECONDS.toNanos(500)); + assertTrue(delayNs < TimeUnit.MILLISECONDS.toNanos(1500)); + + // 4. producer disconnect + producer.close(); + + // 5. Restart broker + setup(); + } + @Test public void testProducerQueueFullNonBlocking() throws Exception { final String topicName = "persistent://prop/ns-abc/topic-xyzx"; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PartitionedProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PartitionedProducerImpl.java index 2dc826d9e3af3..903a1beaaeee3 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PartitionedProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PartitionedProducerImpl.java @@ -19,6 +19,8 @@ package org.apache.pulsar.client.impl; import static com.google.common.base.Preconditions.checkArgument; +import static org.apache.pulsar.client.impl.conf.ProducerConfigurationData.DEFAULT_MAX_PENDING_MESSAGES; +import static org.apache.pulsar.client.impl.conf.ProducerConfigurationData.DEFAULT_MAX_PENDING_MESSAGES_ACROSS_PARTITIONS; import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.ImmutableList; import io.netty.util.Timeout; @@ -84,9 +86,16 @@ public PartitionedProducerImpl(PulsarClientImpl client, String topic, ProducerCo : null; // MaxPendingMessagesAcrossPartitions doesn't support partial partition such as SinglePartition correctly - int maxPendingMessages = Math.min(conf.getMaxPendingMessages(), - conf.getMaxPendingMessagesAcrossPartitions() / numPartitions); - conf.setMaxPendingMessages(maxPendingMessages); + int maxPendingMessages = conf.getMaxPendingMessages(); + int maxPendingMessagesAcrossPartitions = conf.getMaxPendingMessagesAcrossPartitions(); + if (maxPendingMessagesAcrossPartitions != DEFAULT_MAX_PENDING_MESSAGES_ACROSS_PARTITIONS) { + int maxPendingMsgsForOnePartition = maxPendingMessagesAcrossPartitions / numPartitions; + maxPendingMessages = (maxPendingMessages == DEFAULT_MAX_PENDING_MESSAGES) + ? maxPendingMsgsForOnePartition + : Math.min(maxPendingMessages, maxPendingMsgsForOnePartition); + conf.setMaxPendingMessages(maxPendingMessages); + } + final List indexList; if (conf.isLazyStartPartitionedProducers() diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/PartitionedProducerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/PartitionedProducerImplTest.java index f96d2e2e0b0e9..0ce95dc1264a7 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/PartitionedProducerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/PartitionedProducerImplTest.java @@ -272,6 +272,44 @@ public void testGetNumOfPartitions() throws Exception { assertEquals(producerImpl.getNumOfPartitions(), 0); } + @Test + public void testMaxPendingQueueSize() throws Exception { + String topicName = "test-max-pending-queue-size"; + ClientConfigurationData conf = new ClientConfigurationData(); + conf.setServiceUrl("pulsar://localhost:6650"); + conf.setStatsIntervalSeconds(100); + + ThreadFactory threadFactory = new DefaultThreadFactory("client-test-stats", Thread.currentThread().isDaemon()); + @Cleanup("shutdownGracefully") + EventLoopGroup eventLoopGroup = EventLoopUtil.newEventLoopGroup(conf.getNumIoThreads(), false, threadFactory); + + @Cleanup + PulsarClientImpl clientImpl = new PulsarClientImpl(conf, eventLoopGroup); + + // Test set maxPendingMessage to 10 + ProducerConfigurationData producerConfData = new ProducerConfigurationData(); + producerConfData.setMessageRoutingMode(MessageRoutingMode.CustomPartition); + producerConfData.setCustomMessageRouter(new CustomMessageRouter()); + producerConfData.setMaxPendingMessages(10); + PartitionedProducerImpl partitionedProducerImpl = new PartitionedProducerImpl( + clientImpl, topicName, producerConfData, 1, null, null, null); + assertEquals(partitionedProducerImpl.getConfiguration().getMaxPendingMessages(), 10); + + // Test set MaxPendingMessagesAcrossPartitions=5 + producerConfData.setMaxPendingMessages(ProducerConfigurationData.DEFAULT_MAX_PENDING_MESSAGES); + producerConfData.setMaxPendingMessagesAcrossPartitions(5); + partitionedProducerImpl = new PartitionedProducerImpl( + clientImpl, topicName, producerConfData, 1, null, null, null); + assertEquals(partitionedProducerImpl.getConfiguration().getMaxPendingMessages(), 5); + + // Test set maxPendingMessage=10 and MaxPendingMessagesAcrossPartitions=10 with 2 partitions + producerConfData.setMaxPendingMessages(10); + producerConfData.setMaxPendingMessagesAcrossPartitions(10); + partitionedProducerImpl = new PartitionedProducerImpl( + clientImpl, topicName, producerConfData, 2, null, null, null); + assertEquals(partitionedProducerImpl.getConfiguration().getMaxPendingMessages(), 5); + } + @Test public void testOnTopicsExtended() throws Exception {