From 886c5953c767e9dae8e6d261dc76d579b00eec51 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Tue, 3 Jan 2017 11:40:53 -0800 Subject: [PATCH] Fix: Fair message consumption from all partitions in partitioned-consumer --- .../api/PartitionedProducerConsumerTest.java | 64 +++++++++++++++++++ .../client/impl/PartitionedConsumerImpl.java | 6 +- 2 files changed, 68 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/com/yahoo/pulsar/client/api/PartitionedProducerConsumerTest.java b/pulsar-broker/src/test/java/com/yahoo/pulsar/client/api/PartitionedProducerConsumerTest.java index a25503abe5b83..cc001a601f66a 100644 --- a/pulsar-broker/src/test/java/com/yahoo/pulsar/client/api/PartitionedProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/com/yahoo/pulsar/client/api/PartitionedProducerConsumerTest.java @@ -431,6 +431,70 @@ public void testAsyncPartitionedProducerConsumerQueueSizeOne() throws Exception log.info("-- Exiting {} test --", methodName); } + /** + * It verifies that consumer consumes from all the partitions fairly. + * + * @throws Exception + */ + @Test + public void testFairDistributionForPartitionConsumers() throws Exception { + log.info("-- Starting {} test --", methodName); + + final int numPartitions = 2; + final String topicName = "persistent://my-property/use/my-ns/my-topic"; + final String producer1Msg = "producer1"; + final String producer2Msg = "producer2"; + final int queueSize = 10; + + ConsumerConfiguration conf = new ConsumerConfiguration(); + conf.setReceiverQueueSize(queueSize); + + admin.persistentTopics().createPartitionedTopic(topicName, numPartitions); + + ProducerConfiguration producerConf = new ProducerConfiguration(); + producerConf.setMessageRoutingMode(MessageRoutingMode.RoundRobinPartition); + Producer producer1 = pulsarClient.createProducer(topicName + "-partition-0", producerConf); + Producer producer2 = pulsarClient.createProducer(topicName + "-partition-1", producerConf); + + Consumer consumer = pulsarClient.subscribe(topicName, "my-partitioned-subscriber", conf); + + int partition2Msgs = 0; + + // produce messages on Partition-1: which will makes partitioned-consumer's queue full + for (int i = 0; i < queueSize - 1; i++) { + producer1.send((producer1Msg + "-" + i).getBytes()); + } + + Thread.sleep(1000); + + // now queue is full : so, partition-2 consumer will be pushed to paused-consumer list + for (int i = 0; i < 5; i++) { + producer2.send((producer2Msg + "-" + i).getBytes()); + } + + // now, Queue should take both partition's messages + // also: we will keep producing messages to partition-1 + int produceMsgInPartition1AfterNumberOfConsumeMessages = 2; + for (int i = 0; i < 3 * queueSize; i++) { + Message msg = consumer.receive(); + partition2Msgs += (new String(msg.getData())).startsWith(producer2Msg) ? 1 : 0; + if (i >= produceMsgInPartition1AfterNumberOfConsumeMessages) { + producer1.send(producer1Msg.getBytes()); + Thread.sleep(100); + } + + } + + assertTrue(partition2Msgs >= 4); + producer1.close(); + producer2.close(); + consumer.unsubscribe(); + consumer.close(); + admin.persistentTopics().deletePartitionedTopic(topicName); + + log.info("-- Exiting {} test --", methodName); + } + private void receiveAsync(Consumer consumer, int totalMessage, int currentMessage, CountDownLatch latch, final Set consumeMsg, ExecutorService executor) throws PulsarClientException { if (currentMessage < totalMessage) { diff --git a/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/PartitionedConsumerImpl.java b/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/PartitionedConsumerImpl.java index 702cb323577eb..d8c2057c8b88c 100644 --- a/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/PartitionedConsumerImpl.java +++ b/pulsar-client/src/main/java/com/yahoo/pulsar/client/impl/PartitionedConsumerImpl.java @@ -126,8 +126,10 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { // Process the message, add to the queue and trigger listener or async callback messageReceived(message); - if (incomingMessages.size() >= maxReceiverQueueSize) { - // No more space left in shared queue, mark this consumer to be resumed later + if (incomingMessages.size() >= maxReceiverQueueSize + || (incomingMessages.size() > sharedQueueResumeThreshold && !pausedConsumers.isEmpty())) { + // mark this consumer to be resumed later: if No more space left in shared queue, + // or if any consumer is already paused (to create fair chance for already paused consumers) pausedConsumers.add(consumer); } else { // Schedule next receiveAsync() if the incoming queue is not full. Use a different thread to avoid