From 1b72bbfec275da3bacae9080737a75fb22afa3e7 Mon Sep 17 00:00:00 2001 From: Michael Marshall Date: Thu, 21 Oct 2021 16:44:27 -0500 Subject: [PATCH 1/5] [Java Client] Remove data race in MultiTopicsConsumerImpl to ensure correct message order --- .../pulsar/client/impl/ConsumerImpl.java | 3 +- .../client/impl/MultiTopicsConsumerImpl.java | 30 ++++++++++--------- 2 files changed, 17 insertions(+), 16 deletions(-) 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 b29d6c73deff6..d4c6871fb2481 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 @@ -427,8 +427,7 @@ protected CompletableFuture> internalReceiveAsync() { if (message == null) { pendingReceives.add(result); cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); - } - if (message != null) { + } else { messageProcessed(message); result.complete(beforeConsume(message)); } 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 520e7f3ab099c..8e7b6112f5375 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 @@ -245,7 +245,7 @@ private void startReceivingMessages(List> newConsumers) { } private void receiveMessageFromConsumer(ConsumerImpl consumer) { - consumer.receiveAsync().thenAccept(message -> { + consumer.receiveAsync().thenAcceptAsync(message -> { if (log.isDebugEnabled()) { log.debug("[{}] [{}] Receive message from sub consumer:{}", topic, subscription, consumer.getTopic()); @@ -260,7 +260,7 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { // or if any consumer is already paused (to create fair chance for already paused consumers) pausedConsumers.add(consumer); - // Since we din't get a mutex, the condition on the incoming queue might have changed after + // Since we didn't get a mutex, the condition on the incoming queue might have changed after // we have paused the current consumer. We need to re-check in order to avoid this consumer // from getting stalled. resumeReceivingFromPausedConsumersIfNeeded(); @@ -269,7 +269,7 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { // recursion and stack overflow internalPinnedExecutor.execute(() -> receiveMessageFromConsumer(consumer)); } - }).exceptionally(ex -> { + }, internalPinnedExecutor).exceptionally(ex -> { if (ex instanceof PulsarClientException.AlreadyClosedException || ex.getCause() instanceof PulsarClientException.AlreadyClosedException) { // ignore the exception that happens when the consumer is closed @@ -281,6 +281,7 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { }); } + // Must be called from the internalPinnedExecutor thread private void messageReceived(ConsumerImpl consumer, Message message) { checkArgument(message instanceof MessageImpl); TopicMessageImpl topicMessage = new TopicMessageImpl<>(consumer.getTopic(), @@ -409,17 +410,18 @@ protected CompletableFuture> internalBatchReceiveAsync() { protected CompletableFuture> internalReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - Message message = incomingMessages.poll(); - if (message == null) { - pendingReceives.add(result); - cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); - } else { - decreaseIncomingMessageSize(message); - checkState(message instanceof TopicMessageImpl); - unAckedMessageTracker.add(message.getMessageId()); - resumeReceivingFromPausedConsumersIfNeeded(); - result.complete(message); - } + internalPinnedExecutor.execute(() -> { + Message message = incomingMessages.poll(); + if (message == null) { + pendingReceives.add(result); + cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); + } else { + decreaseIncomingMessageSize(message); + unAckedMessageTracker.add(message.getMessageId()); + resumeReceivingFromPausedConsumersIfNeeded(); + result.complete(message); + } + }); return result; } From 8810b58b08bf2d732210b8a7e14fd5c604e2ca00 Mon Sep 17 00:00:00 2001 From: Michael Marshall Date: Fri, 22 Oct 2021 09:29:47 -0500 Subject: [PATCH 2/5] Fix test --- .../apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java index 6af8914d6943d..faa621cae2193 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImplTest.java @@ -36,6 +36,7 @@ import org.apache.pulsar.client.util.ExecutorProvider; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.util.netty.EventLoopUtil; +import org.awaitility.Awaitility; import org.junit.After; import org.junit.Before; import org.testng.annotations.AfterMethod; @@ -165,7 +166,7 @@ public void testReceiveAsyncCanBeCancelled() { // given MultiTopicsConsumerImpl consumer = createMultiTopicsConsumer(); CompletableFuture> future = consumer.receiveAsync(); - assertTrue(consumer.hasNextPendingReceive()); + Awaitility.await().untilAsserted(() -> assertTrue(consumer.hasNextPendingReceive())); // when future.cancel(true); // then From e13139082ed64144b897f49493276929b18af661 Mon Sep 17 00:00:00 2001 From: Michael Marshall Date: Fri, 22 Oct 2021 09:38:57 -0500 Subject: [PATCH 3/5] Return the checkState method call to keep original behavior --- .../org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java | 1 + 1 file changed, 1 insertion(+) 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 8e7b6112f5375..15ab78137a7a0 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 @@ -417,6 +417,7 @@ protected CompletableFuture> internalReceiveAsync() { cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); } else { decreaseIncomingMessageSize(message); + checkState(message instanceof TopicMessageImpl); unAckedMessageTracker.add(message.getMessageId()); resumeReceivingFromPausedConsumersIfNeeded(); result.complete(message); From f74aaf50c7457da90b75870acf65ae68457c765b Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sat, 23 Oct 2021 20:22:56 +0300 Subject: [PATCH 4/5] Reproduce out-of-order delivery issue in PR 12456 --- .../client/api/MultiTopicsConsumerTest.java | 75 +++++++++++++++++++ 1 file changed, 75 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MultiTopicsConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MultiTopicsConsumerTest.java index 715f3adb7aaf6..d8c8bd657f8cc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MultiTopicsConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MultiTopicsConsumerTest.java @@ -16,6 +16,7 @@ * specific language governing permissions and limitations * under the License. */ + package org.apache.pulsar.client.api; import static org.mockito.ArgumentMatchers.any; @@ -24,9 +25,16 @@ import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import com.google.common.collect.Lists; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; import lombok.Cleanup; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.impl.ClientBuilderImpl; import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; @@ -34,6 +42,7 @@ import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -70,6 +79,7 @@ protected PulsarClient createNewPulsarClient(ClientBuilder clientBuilder) throws // method calls on the interface. Mockito.withSettings().defaultAnswer(AdditionalAnswers.delegatesTo(internalExecutorService))); } + @Override public ExecutorService getInternalExecutorService() { return internalExecutorServiceDelegate; @@ -119,4 +129,69 @@ public void testMultiTopicsConsumerCloses() throws Exception { verify(internalExecutorServiceDelegate, times(0)) .schedule(any(Runnable.class), anyLong(), any()); } + + // test that reproduces the issue that PR https://github.com/apache/pulsar/pull/12456 fixes + // where MultiTopicsConsumerImpl has a data race that causes out-of-order delivery of messages + @Test + public void testShouldMaintainOrderForIndividualTopicInMultiTopicsConsumer() + throws PulsarAdminException, PulsarClientException, ExecutionException, InterruptedException, + TimeoutException { + String topicName = newTopicName(); + int numPartitions = 2; + int numMessages = 100000; + admin.topics().createPartitionedTopic(topicName, numPartitions); + + Producer[] producers = new Producer[numPartitions]; + + for (int i = 0; i < numPartitions; i++) { + producers[i] = pulsarClient.newProducer(Schema.INT64) + // produce to each partition directly so that order can be maintained in sending + .topic(topicName + "-partition-" + i) + .enableBatching(true) + .maxPendingMessages(30000) + .maxPendingMessagesAcrossPartitions(60000) + .batchingMaxMessages(10000) + .batchingMaxPublishDelay(5, TimeUnit.SECONDS) + .batchingMaxBytes(4 * 1024 * 1024) + .blockIfQueueFull(true) + .create(); + } + + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.INT64) + // consume on the partitioned topic + .topic(topicName) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) + .receiverQueueSize(numMessages) + .subscriptionName(methodName) + .subscribe(); + + // produce sequence numbers to each partition topic + long sequenceNumber = 1L; + for (int i = 0; i < numMessages; i++) { + for (Producer producer : producers) { + producer.newMessage() + .value(sequenceNumber) + .sendAsync(); + } + sequenceNumber++; + } + for (Producer producer : producers) { + producer.close(); + } + + // receive and validate sequences in the partitioned topic + Map receivedSequences = new HashMap<>(); + int receivedCount = 0; + while (receivedCount < numPartitions * numMessages) { + Message message = consumer.receiveAsync().get(5, TimeUnit.SECONDS); + consumer.acknowledge(message); + receivedCount++; + AtomicLong receivedSequenceCounter = + receivedSequences.computeIfAbsent(message.getTopicName(), k -> new AtomicLong(1L)); + Assert.assertEquals(message.getValue().longValue(), receivedSequenceCounter.getAndIncrement()); + } + Assert.assertEquals(numPartitions * numMessages, receivedCount); + } } From 082b46b3b6484aa8dc9a15625a35851c56c9880c Mon Sep 17 00:00:00 2001 From: Michael Marshall Date: Mon, 25 Oct 2021 11:22:31 -0500 Subject: [PATCH 5/5] Remove unnecessary scheduling of receiveMessageFromConsumer --- .../apache/pulsar/client/impl/MultiTopicsConsumerImpl.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 15ab78137a7a0..21ae2d7ecb0cd 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 @@ -265,9 +265,9 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { // from getting stalled. resumeReceivingFromPausedConsumersIfNeeded(); } else { - // Schedule next receiveAsync() if the incoming queue is not full. Use a different thread to avoid - // recursion and stack overflow - internalPinnedExecutor.execute(() -> receiveMessageFromConsumer(consumer)); + // Call receiveAsync() if the incoming queue is not full. Because this block is run with + // thenAcceptAsync, there is no chance for recursion that would lead to stack overflow. + receiveMessageFromConsumer(consumer); } }, internalPinnedExecutor).exceptionally(ex -> { if (ex instanceof PulsarClientException.AlreadyClosedException