From 02049aaa9af9fbe3ef2190f30d276db2ab513a13 Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 12 Jan 2021 11:20:56 +0800 Subject: [PATCH 1/2] Fix incoming queue size issue that introduced in #9113 --- .../api/SimpleProducerConsumerTest.java | 85 +++++++++++++++++++ .../pulsar/client/impl/ConsumerBase.java | 20 +++-- .../pulsar/client/impl/ConsumerImpl.java | 4 +- .../client/impl/MultiTopicsConsumerImpl.java | 12 +-- 4 files changed, 108 insertions(+), 13 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 6354e605ce153..15e16439c5e9f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -77,6 +77,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.impl.ClientCnx; +import org.apache.pulsar.client.impl.ConsumerBase; import org.apache.pulsar.client.impl.ConsumerImpl; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.MessageImpl; @@ -3692,4 +3693,88 @@ public void testGetStatsForPartitionedTopic() throws Exception { consumer.close(); producer.close(); } + + @Test + public void testIncomingMessageSizeForNonPartitionedTopic() throws Exception { + final String topicName = "persistent://my-property/my-ns/testIncomingMessageSizeForNonPartitionedTopic-" + + UUID.randomUUID().toString(); + final String subName = "my-sub"; + + @Cleanup + Consumer consumer = pulsarClient.newConsumer() + .topic(topicName) + .subscriptionName(subName) + .subscribe(); + + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(topicName) + .create(); + + final int messages = 100; + List> messageIds = new ArrayList<>(messages); + for (int i = 0; i < messages; i++) { + messageIds.add(producer.newMessage().value(("Message-" + i).getBytes()).sendAsync()); + } + FutureUtil.waitForAll(messageIds).get(); + + Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { + long size = ((ConsumerBase) consumer).getIncomingMessageSize(); + log.info("Check the incoming message size should greater that 0, current size is {}", size); + Assert.assertTrue(size > 0); + }); + + for (int i = 0; i < messages; i++) { + consumer.acknowledge(consumer.receive()); + } + + Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { + long size = ((ConsumerBase) consumer).getIncomingMessageSize(); + log.info("Check the incoming message size should be 0, current size is {}", size); + Assert.assertEquals(size, 0); + }); + } + + @Test + public void testIncomingMessageSizeForPartitionedTopic() throws Exception { + final String topicName = "persistent://my-property/my-ns/testIncomingMessageSizeForPartitionedTopic-" + + UUID.randomUUID().toString(); + final String subName = "my-sub"; + + admin.topics().createPartitionedTopic(topicName, 3); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer() + .topic(topicName) + .subscriptionName(subName) + .subscribe(); + + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(topicName) + .create(); + + final int messages = 100; + List> messageIds = new ArrayList<>(messages); + for (int i = 0; i < messages; i++) { + messageIds.add(producer.newMessage().key(i + "").value(("Message-" + i).getBytes()).sendAsync()); + } + FutureUtil.waitForAll(messageIds).get(); + + Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { + long size = ((ConsumerBase) consumer).getIncomingMessageSize(); + log.info("Check the incoming message size should greater that 0, current size is {}", size); + Assert.assertTrue(size > 0); + }); + + for (int i = 0; i < messages; i++) { + consumer.acknowledge(consumer.receive()); + } + + Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { + long size = ((ConsumerBase) consumer).getIncomingMessageSize(); + log.info("Check the incoming message size should be 0, current size is {}", size); + Assert.assertEquals(size, 0); + }); + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java index e258704f2b0bc..2479c8b248b35 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java @@ -20,6 +20,7 @@ import static com.google.common.base.Preconditions.checkArgument; +import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Queues; import java.util.Collections; import java.util.List; @@ -664,8 +665,7 @@ protected boolean canEnqueueMessage(Message message) { protected boolean enqueueMessageAndCheckBatchReceive(Message message) { if (canEnqueueMessage(message) && incomingMessages.offer(message)) { - INCOMING_MESSAGES_SIZE_UPDATER.addAndGet( - this, message.getData() == null ? 0 : message.getData().length); + increaseIncomingMessageSize(message); } return hasEnoughMessagesForBatchReceive(); } @@ -675,7 +675,7 @@ protected boolean hasEnoughMessagesForBatchReceive() { return false; } return (batchReceivePolicy.getMaxNumMessages() > 0 && incomingMessages.size() >= batchReceivePolicy.getMaxNumMessages()) - || (batchReceivePolicy.getMaxNumBytes() > 0 && INCOMING_MESSAGES_SIZE_UPDATER.get(this) >= batchReceivePolicy.getMaxNumBytes()); + || (batchReceivePolicy.getMaxNumBytes() > 0 && getIncomingMessageSize() >= batchReceivePolicy.getMaxNumBytes()); } private void verifyConsumerState() throws PulsarClientException { @@ -847,13 +847,23 @@ protected boolean hasPendingBatchReceive() { return pendingBatchReceives != null && peekNextBatchReceive() != null; } + protected void increaseIncomingMessageSize(final Message message) { + INCOMING_MESSAGES_SIZE_UPDATER.addAndGet( + this, message.getData() == null ? 0 : message.getData().length); + } + protected void resetIncomingMessageSize() { INCOMING_MESSAGES_SIZE_UPDATER.set(this, 0); } - protected void updateIncomingMessageSize(final Message message) { + protected void decreaseIncomingMessageSize(final Message message) { INCOMING_MESSAGES_SIZE_UPDATER.addAndGet(this, - (message.getData() != null) ? message.getData().length : 0); + (message.getData() != null) ? -message.getData().length : 0); + } + + @VisibleForTesting + public long getIncomingMessageSize() { + return INCOMING_MESSAGES_SIZE_UPDATER.get(this); } protected abstract void completeOpBatchReceive(OpBatchReceive op); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 74c8eea9f5a34..e9d7ad2149fb9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -1507,7 +1507,7 @@ protected synchronized void messageProcessed(Message msg) { stats.updateNumMsgsReceived(msg); trackMessage(msg); - updateIncomingMessageSize(msg); + decreaseIncomingMessageSize(msg); } protected void trackMessage(Message msg) { @@ -2197,7 +2197,7 @@ private int removeExpiredMessagesFromQueue(Set messageIds) { // try not to remove elements that are added while we remove Message message = incomingMessages.poll(); while (message != null) { - updateIncomingMessageSize(message); + decreaseIncomingMessageSize(message); messagesFromQueue++; MessageIdImpl id = getMessageIdImpl(message); if (!messageIds.contains(id)) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java index ef5d1f26f2380..8910393477b3c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java @@ -312,7 +312,7 @@ private void messageReceived(ConsumerImpl consumer, Message message) { @Override protected synchronized void messageProcessed(Message msg) { unAckedMessageTracker.add(msg.getMessageId()); - updateIncomingMessageSize(msg); + decreaseIncomingMessageSize(msg); } private void resumeReceivingFromPausedConsumersIfNeeded() { @@ -335,7 +335,7 @@ protected Message internalReceive() throws PulsarClientException { Message message; try { message = incomingMessages.take(); - updateIncomingMessageSize(message); + decreaseIncomingMessageSize(message); checkState(message instanceof TopicMessageImpl); unAckedMessageTracker.add(message.getMessageId()); resumeReceivingFromPausedConsumersIfNeeded(); @@ -351,7 +351,7 @@ protected Message internalReceive(int timeout, TimeUnit unit) throws PulsarCl try { message = incomingMessages.poll(timeout, unit); if (message != null) { - updateIncomingMessageSize(message); + decreaseIncomingMessageSize(message); checkArgument(message instanceof TopicMessageImpl); unAckedMessageTracker.add(message.getMessageId()); } @@ -392,7 +392,7 @@ protected CompletableFuture> internalBatchReceiveAsync() { while (msgPeeked != null && messages.canAdd(msgPeeked)) { Message msg = incomingMessages.poll(); if (msg != null) { - updateIncomingMessageSize(msg); + decreaseIncomingMessageSize(msg); Message interceptMsg = beforeConsume(msg); messages.add(interceptMsg); } @@ -420,7 +420,7 @@ protected CompletableFuture> internalReceiveAsync() { pendingReceives.add(result); cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); } else { - updateIncomingMessageSize(message); + decreaseIncomingMessageSize(message); checkState(message instanceof TopicMessageImpl); unAckedMessageTracker.add(message.getMessageId()); resumeReceivingFromPausedConsumersIfNeeded(); @@ -785,7 +785,7 @@ private void removeExpiredMessagesFromQueue(Set messageIds) { Message message = incomingMessages.poll(); checkState(message instanceof TopicMessageImpl); while (message != null) { - updateIncomingMessageSize(message); + decreaseIncomingMessageSize(message); MessageId messageId = message.getMessageId(); if (!messageIds.contains(messageId)) { messageIds.add(messageId); From ee87882cd061930e487d142032bf85e59f881380 Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 12 Jan 2021 13:37:20 +0800 Subject: [PATCH 2/2] Apply comment --- .../api/SimpleProducerConsumerTest.java | 52 ++++--------------- .../pulsar/client/impl/ConsumerBase.java | 1 - 2 files changed, 9 insertions(+), 44 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 15e16439c5e9f..9fc14c03b0105 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -3694,54 +3694,20 @@ public void testGetStatsForPartitionedTopic() throws Exception { producer.close(); } - @Test - public void testIncomingMessageSizeForNonPartitionedTopic() throws Exception { - final String topicName = "persistent://my-property/my-ns/testIncomingMessageSizeForNonPartitionedTopic-" + - UUID.randomUUID().toString(); - final String subName = "my-sub"; - - @Cleanup - Consumer consumer = pulsarClient.newConsumer() - .topic(topicName) - .subscriptionName(subName) - .subscribe(); - - @Cleanup - Producer producer = pulsarClient.newProducer() - .topic(topicName) - .create(); - - final int messages = 100; - List> messageIds = new ArrayList<>(messages); - for (int i = 0; i < messages; i++) { - messageIds.add(producer.newMessage().value(("Message-" + i).getBytes()).sendAsync()); - } - FutureUtil.waitForAll(messageIds).get(); - - Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { - long size = ((ConsumerBase) consumer).getIncomingMessageSize(); - log.info("Check the incoming message size should greater that 0, current size is {}", size); - Assert.assertTrue(size > 0); - }); - - for (int i = 0; i < messages; i++) { - consumer.acknowledge(consumer.receive()); - } - - Awaitility.await().atMost(3, TimeUnit.SECONDS).untilAsserted(() -> { - long size = ((ConsumerBase) consumer).getIncomingMessageSize(); - log.info("Check the incoming message size should be 0, current size is {}", size); - Assert.assertEquals(size, 0); - }); + @DataProvider(name = "partitioned") + public static Object[] isPartitioned() { + return new Object[] {false, true}; } - @Test - public void testIncomingMessageSizeForPartitionedTopic() throws Exception { - final String topicName = "persistent://my-property/my-ns/testIncomingMessageSizeForPartitionedTopic-" + + @Test(dataProvider = "partitioned") + public void testIncomingMessageSize(boolean isPartitioned) throws Exception { + final String topicName = "persistent://my-property/my-ns/testIncomingMessageSize-" + UUID.randomUUID().toString(); final String subName = "my-sub"; - admin.topics().createPartitionedTopic(topicName, 3); + if (isPartitioned) { + admin.topics().createPartitionedTopic(topicName, 3); + } @Cleanup Consumer consumer = pulsarClient.newConsumer() diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java index 2479c8b248b35..b100a7a257b8d 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java @@ -861,7 +861,6 @@ protected void decreaseIncomingMessageSize(final Message message) { (message.getData() != null) ? -message.getData().length : 0); } - @VisibleForTesting public long getIncomingMessageSize() { return INCOMING_MESSAGES_SIZE_UPDATER.get(this); }