From a5d8a16ee69ebd4a1ef8a13ef7f56a8ac13d8ec4 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Sun, 25 Apr 2021 00:38:22 +0800 Subject: [PATCH 01/12] lock-free --- .../client/impl/MultiTopicsReaderTest.java | 7 +++ .../pulsar/client/impl/ConsumerImpl.java | 44 +++++++++---------- .../pulsar/client/impl/PulsarClientImpl.java | 4 ++ 3 files changed, 32 insertions(+), 23 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java index c4adb9712a091..e6d14bd95ff41 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java @@ -254,6 +254,13 @@ public void testMultiReaderSeek() throws Exception { publishMessages(topic,100,false); } + @Test + public void test() throws Exception { + for (int i = 0; i < 100; i++) { + testMultiTopic(); + } + } + @Test(timeOut = 20000) public void testMultiTopic() throws Exception { final String topic = "persistent://my-property/my-ns/topic" + UUID.randomUUID(); 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 13b6a96a482ac..1075a8b8ae24d 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 @@ -412,22 +412,17 @@ protected Message internalReceive() throws PulsarClientException { protected CompletableFuture> internalReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - Message message = null; - try { - message = incomingMessages.poll(0, TimeUnit.MILLISECONDS); + client.getInternalExecutorService(consumerId).execute(() -> { + Message message = incomingMessages.poll(); if (message == null) { pendingReceives.add(result); cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); } - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - result.completeExceptionally(e); - } - - if (message != null) { - messageProcessed(message); - result.complete(beforeConsume(message)); - } + if (message != null) { + messageProcessed(message); + result.complete(beforeConsume(message)); + } + }); return result; } @@ -1066,11 +1061,13 @@ uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redelive possibleSendToDeadLetterTopicMessages.put((MessageIdImpl)message.getMessageId(), Collections.singletonList(message)); } - if (peekPendingReceive() != null) { - notifyPendingReceivedCallback(message, null); - } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { - notifyPendingBatchReceivedCallBack(); - } + client.getInternalExecutorService(consumerId).execute(() -> { + if (peekPendingReceive() != null) { + notifyPendingReceivedCallback(message, null); + } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { + notifyPendingBatchReceivedCallBack(); + } + }); } else { // handle batch message enqueuing; uncompressed payload has all messages in batch receiveIndividualMessagesFromBatch(msgMetadata, redeliveryCount, ackSet, uncompressedPayload, messageId, cnx); @@ -1280,12 +1277,13 @@ void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata, int redeliv if (possibleToDeadLetter != null) { possibleToDeadLetter.add(message); } - - if (peekPendingReceive() != null) { - notifyPendingReceivedCallback(message, null); - } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { - notifyPendingBatchReceivedCallBack(); - } + client.getInternalExecutorService(consumerId).execute(() -> { + if (peekPendingReceive() != null) { + notifyPendingReceivedCallback(message, null); + } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { + notifyPendingBatchReceivedCallBack(); + } + }); singleMessagePayload.release(); } if (ackBitSet != null) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java index df121095dd5fd..d58ea4b221a37 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java @@ -1029,6 +1029,10 @@ protected CompletableFuture> preProcessSchemaBeforeSubscribe(Pulsa public ExecutorService getInternalExecutorService() { return internalExecutorService.getExecutor(); } + + public ExecutorService getInternalExecutorService(Object obj) { + return internalExecutorService.getExecutor(obj); + } // // Transaction related API // From 335522d891b6963b35e94a1c074cf68a21799ec6 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Sun, 25 Apr 2021 19:51:35 +0800 Subject: [PATCH 02/12] Change each locked place to use thread pool --- .../pulsar/client/impl/ConsumerImpl.java | 31 +++++++++++-------- 1 file changed, 18 insertions(+), 13 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 1075a8b8ae24d..d70a5ca3e4987 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 @@ -47,6 +47,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; @@ -188,6 +189,7 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final boolean poolMessages; private final AtomicReference clientCnxUsedForConsumerRegistration = new AtomicReference<>(); + private final ExecutorService internalPinnedExecutor; static ConsumerImpl newConsumerImpl(PulsarClientImpl client, String topic, @@ -255,6 +257,7 @@ protected ConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurat this.expireTimeOfIncompleteChunkedMessageMillis = conf.getExpireTimeOfIncompleteChunkedMessageMillis(); this.autoAckOldestChunkedMessageOnQueueFull = conf.isAutoAckOldestChunkedMessageOnQueueFull(); this.poolMessages = conf.isPoolMessages(); + this.internalPinnedExecutor = client.getInternalExecutorService(consumerId); if (client.getConfiguration().getStatsIntervalSeconds() > 0) { stats = new ConsumerStatsRecorderImpl(client, conf, this); @@ -412,7 +415,7 @@ protected Message internalReceive() throws PulsarClientException { protected CompletableFuture> internalReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - client.getInternalExecutorService(consumerId).execute(() -> { + internalPinnedExecutor.execute(() -> { Message message = incomingMessages.poll(); if (message == null) { pendingReceives.add(result); @@ -467,8 +470,7 @@ protected Messages internalBatchReceive() throws PulsarClientException { protected CompletableFuture> internalBatchReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - try { - lock.writeLock().lock(); + internalPinnedExecutor.execute(() -> { if (pendingBatchReceives == null) { pendingBatchReceives = Queues.newConcurrentLinkedQueue(); } @@ -490,10 +492,7 @@ protected CompletableFuture> internalBatchReceiveAsync() { pendingBatchReceives.add(opBatchReceive); cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); } - } finally { - lock.writeLock().unlock(); - } - + }); return result; } @@ -952,10 +951,12 @@ private void closeConsumerTasks() { } private void failPendingReceive() { - if (pinnedExecutor != null && !pinnedExecutor.isShutdown()) { - failPendingReceives(this.pendingReceives); - failPendingBatchReceives(this.pendingBatchReceives); - } + internalPinnedExecutor.execute(() -> { + if (pinnedExecutor != null && !pinnedExecutor.isShutdown()) { + failPendingReceives(this.pendingReceives); + failPendingBatchReceives(this.pendingBatchReceives); + } + }); } void activeConsumerChanged(boolean isActive) { @@ -1061,20 +1062,24 @@ uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redelive possibleSendToDeadLetterTopicMessages.put((MessageIdImpl)message.getMessageId(), Collections.singletonList(message)); } - client.getInternalExecutorService(consumerId).execute(() -> { + internalPinnedExecutor.execute(() -> { if (peekPendingReceive() != null) { notifyPendingReceivedCallback(message, null); } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { notifyPendingBatchReceivedCallBack(); } + tryTriggerListener(); }); } else { // handle batch message enqueuing; uncompressed payload has all messages in batch receiveIndividualMessagesFromBatch(msgMetadata, redeliveryCount, ackSet, uncompressedPayload, messageId, cnx); uncompressedPayload.release(); + tryTriggerListener(); } + } + private void tryTriggerListener() { if (listener != null) { triggerListener(); } @@ -1277,7 +1282,7 @@ void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata, int redeliv if (possibleToDeadLetter != null) { possibleToDeadLetter.add(message); } - client.getInternalExecutorService(consumerId).execute(() -> { + internalPinnedExecutor.execute(() -> { if (peekPendingReceive() != null) { notifyPendingReceivedCallback(message, null); } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { From d9c7e8ffef6bb21f5b03406174af2cb5d8465106 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Sun, 25 Apr 2021 21:48:58 +0800 Subject: [PATCH 03/12] Fix unit test --- pulsar-client/pom.xml | 7 +++++++ .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 2 +- .../org/apache/pulsar/client/impl/PulsarClientImpl.java | 4 ---- .../org/apache/pulsar/client/impl/ClientTestFixtures.java | 2 +- .../org/apache/pulsar/client/impl/ConsumerImplTest.java | 3 ++- 5 files changed, 11 insertions(+), 7 deletions(-) diff --git a/pulsar-client/pom.xml b/pulsar-client/pom.xml index 0fb444909c54a..142671a8591a0 100644 --- a/pulsar-client/pom.xml +++ b/pulsar-client/pom.xml @@ -169,6 +169,13 @@ ${skyscreamer.version} test + + + org.awaitility + awaitility + test + + 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 d70a5ca3e4987..5f40cbc8aa280 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 @@ -257,7 +257,7 @@ protected ConsumerImpl(PulsarClientImpl client, String topic, ConsumerConfigurat this.expireTimeOfIncompleteChunkedMessageMillis = conf.getExpireTimeOfIncompleteChunkedMessageMillis(); this.autoAckOldestChunkedMessageOnQueueFull = conf.isAutoAckOldestChunkedMessageOnQueueFull(); this.poolMessages = conf.isPoolMessages(); - this.internalPinnedExecutor = client.getInternalExecutorService(consumerId); + this.internalPinnedExecutor = client.getInternalExecutorService(); if (client.getConfiguration().getStatsIntervalSeconds() > 0) { stats = new ConsumerStatsRecorderImpl(client, conf, this); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java index d58ea4b221a37..df121095dd5fd 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java @@ -1029,10 +1029,6 @@ protected CompletableFuture> preProcessSchemaBeforeSubscribe(Pulsa public ExecutorService getInternalExecutorService() { return internalExecutorService.getExecutor(); } - - public ExecutorService getInternalExecutorService(Object obj) { - return internalExecutorService.getExecutor(obj); - } // // Transaction related API // diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java index 8bb7bbc41a1c7..f493f5ea925d8 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java @@ -21,7 +21,6 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.channel.EventLoop; import io.netty.util.Timer; -import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.client.util.ExecutorProvider; import org.mockito.Mockito; @@ -48,6 +47,7 @@ static PulsarClientImpl createPulsarClientMock() { when(clientMock.timer()).thenReturn(mock(Timer.class)); when(clientMock.externalExecutorProvider()).thenReturn(mock(ExecutorProvider.class)); + when(clientMock.getInternalExecutorService()).thenReturn(Executors.newSingleThreadExecutor()); when(clientMock.eventLoopGroup().next()).thenReturn(mock(EventLoop.class)); return clientMock; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java index f3ccfd86e80b9..de6f1c41b3fce 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java @@ -38,6 +38,7 @@ import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.apache.pulsar.client.util.ExecutorProvider; +import org.awaitility.Awaitility; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -173,7 +174,7 @@ public void testReceiveAsyncCanBeCancelled() { public void testBatchReceiveAsyncCanBeCancelled() { // given CompletableFuture> future = consumer.batchReceiveAsync(); - Assert.assertTrue(consumer.hasPendingBatchReceive()); + Awaitility.await().untilAsserted(() -> Assert.assertTrue(consumer.hasPendingBatchReceive())); // when future.cancel(true); // then From efa0fd4da54d6676639811cdceb03fd42f3f0fcf Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Sun, 25 Apr 2021 21:50:09 +0800 Subject: [PATCH 04/12] remove test --- .../apache/pulsar/client/impl/MultiTopicsReaderTest.java | 7 ------- 1 file changed, 7 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java index e6d14bd95ff41..c4adb9712a091 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MultiTopicsReaderTest.java @@ -254,13 +254,6 @@ public void testMultiReaderSeek() throws Exception { publishMessages(topic,100,false); } - @Test - public void test() throws Exception { - for (int i = 0; i < 100; i++) { - testMultiTopic(); - } - } - @Test(timeOut = 20000) public void testMultiTopic() throws Exception { final String topic = "persistent://my-property/my-ns/topic" + UUID.randomUUID(); From cb8fd205b90c9324db459924659173b72964c014 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Sun, 25 Apr 2021 23:59:55 +0800 Subject: [PATCH 05/12] fix unit test --- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 5 +++-- .../java/org/apache/pulsar/client/impl/ConsumerImplTest.java | 2 +- 2 files changed, 4 insertions(+), 3 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 5f40cbc8aa280..e196f57c6bb83 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 @@ -1075,8 +1075,8 @@ uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redelive receiveIndividualMessagesFromBatch(msgMetadata, redeliveryCount, ackSet, uncompressedPayload, messageId, cnx); uncompressedPayload.release(); - tryTriggerListener(); } + } private void tryTriggerListener() { @@ -1288,8 +1288,9 @@ void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata, int redeliv } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { notifyPendingBatchReceivedCallBack(); } + singleMessagePayload.release(); + tryTriggerListener(); }); - singleMessagePayload.release(); } if (ackBitSet != null) { ackBitSet.recycle(); diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java index de6f1c41b3fce..816ddb60c9fc1 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java @@ -163,7 +163,7 @@ public void testNotifyPendingReceivedCallback_WorkNormally() { public void testReceiveAsyncCanBeCancelled() { // given CompletableFuture> future = consumer.receiveAsync(); - Assert.assertEquals(consumer.peekPendingReceive(), future); + Awaitility.await().untilAsserted(() -> Assert.assertEquals(consumer.peekPendingReceive(), future)); // when future.cancel(true); // then From 9e0719d28791b47ec18e3da4dcf22d0c1278b949 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Mon, 26 Apr 2021 13:13:04 +0800 Subject: [PATCH 06/12] fix unit test --- .../java/org/apache/pulsar/client/impl/ReaderImplTest.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java index 6e587f720e6b4..6f22cddfecc90 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java @@ -22,6 +22,7 @@ import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.impl.conf.ReaderConfigurationData; +import org.awaitility.Awaitility; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -47,7 +48,9 @@ void setupReader() { void shouldSupportCancellingReadNextAsync() { // given CompletableFuture> future = reader.readNextAsync(); - assertNotNull(reader.getConsumer().peekPendingReceive()); + Awaitility.await().untilAsserted(() -> { + assertNotNull(reader.getConsumer().peekPendingReceive()); + }); // when future.cancel(false); From 87355ee712ff6d2f9eb3f925e607f60e35a63b71 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 27 Apr 2021 00:14:30 +0800 Subject: [PATCH 07/12] Trigger CI to see if there are occasional problems --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 1 + 1 file changed, 1 insertion(+) 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 e196f57c6bb83..d6235fad911e4 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 @@ -421,6 +421,7 @@ protected CompletableFuture> internalReceiveAsync() { pendingReceives.add(result); cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); } + if (message != null) { messageProcessed(message); result.complete(beforeConsume(message)); From 9a1324f63362111281fdf128c8648597cada58ca Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 27 Apr 2021 11:31:19 +0800 Subject: [PATCH 08/12] Trigger CI to see if there are occasional problems --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 1 - 1 file changed, 1 deletion(-) 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 d6235fad911e4..e196f57c6bb83 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 @@ -421,7 +421,6 @@ protected CompletableFuture> internalReceiveAsync() { pendingReceives.add(result); cancellationHandler.setCancelAction(() -> pendingReceives.remove(result)); } - if (message != null) { messageProcessed(message); result.complete(beforeConsume(message)); From e797a03b3ecaa2767d89dc947b1b4bcd79606f0a Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 27 Apr 2021 15:08:37 +0800 Subject: [PATCH 09/12] Shutdown executor --- .../pulsar/client/impl/ClientTestFixtures.java | 1 - .../pulsar/client/impl/ConsumerImplTest.java | 11 ++++++++++- .../apache/pulsar/client/impl/ReaderImplTest.java | 15 +++++++++++++++ 3 files changed, 25 insertions(+), 2 deletions(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java index f493f5ea925d8..0adb165f066b1 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ClientTestFixtures.java @@ -47,7 +47,6 @@ static PulsarClientImpl createPulsarClientMock() { when(clientMock.timer()).thenReturn(mock(Timer.class)); when(clientMock.externalExecutorProvider()).thenReturn(mock(ExecutorProvider.class)); - when(clientMock.getInternalExecutorService()).thenReturn(Executors.newSingleThreadExecutor()); when(clientMock.eventLoopGroup().next()).thenReturn(mock(EventLoop.class)); return clientMock; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java index 816ddb60c9fc1..2702439cee345 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ConsumerImplTest.java @@ -25,12 +25,14 @@ import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import io.netty.util.concurrent.DefaultThreadFactory; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.Messages; @@ -50,12 +52,15 @@ public class ConsumerImplTest { private ExecutorProvider executorProvider; private ConsumerImpl consumer; private ConsumerConfigurationData consumerConf; + private ExecutorService executorService; @BeforeMethod(alwaysRun = true) public void setUp() { executorProvider = new ExecutorProvider(1, "ConsumerImplTest"); consumerConf = new ConsumerConfigurationData<>(); PulsarClientImpl client = ClientTestFixtures.createPulsarClientMock(); + executorService = Executors.newSingleThreadExecutor(); + when(client.getInternalExecutorService()).thenReturn(executorService); ClientConfigurationData clientConf = client.getConfiguration(); clientConf.setOperationTimeoutMs(100); clientConf.setStatsIntervalSeconds(0); @@ -75,6 +80,10 @@ public void cleanup() { executorProvider.shutdownNow(); executorProvider = null; } + if (executorService != null) { + executorService.shutdownNow(); + executorService = null; + } } @Test(invocationTimeOut = 1000) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java index 6f22cddfecc90..d0c4023e09e75 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ReaderImplTest.java @@ -23,27 +23,42 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.impl.conf.ReaderConfigurationData; import org.awaitility.Awaitility; +import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; public class ReaderImplTest { ReaderImpl reader; + private ExecutorService executorService; @BeforeMethod void setupReader() { PulsarClientImpl mockedClient = ClientTestFixtures.createPulsarClientMockWithMockedClientCnx(); ReaderConfigurationData readerConfiguration = new ReaderConfigurationData<>(); readerConfiguration.setTopicName("topicName"); + executorService = Executors.newSingleThreadExecutor(); + when(mockedClient.getInternalExecutorService()).thenReturn(executorService); CompletableFuture> consumerFuture = new CompletableFuture<>(); reader = new ReaderImpl<>(mockedClient, readerConfiguration, ClientTestFixtures.createMockedExecutorProvider(), consumerFuture, Schema.BYTES); } + @AfterMethod + public void clean() { + if (executorService != null) { + executorService.shutdownNow(); + executorService = null; + } + } + @Test void shouldSupportCancellingReadNextAsync() { // given From 2b8bb8e9c30f5d4894fb28ab1f8e4edd5bf54d5e Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Wed, 28 Apr 2021 20:04:17 +0800 Subject: [PATCH 10/12] Trigger CI to see if there are occasional problems --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 9170a156b894b..5e987369b1e01 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 @@ -1060,7 +1060,7 @@ uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redelive internalPinnedExecutor.execute(() -> { if (deadLetterPolicy != null && possibleSendToDeadLetterTopicMessages != null && redeliveryCount >= deadLetterPolicy.getMaxRedeliverCount()) { - possibleSendToDeadLetterTopicMessages.put((MessageIdImpl)message.getMessageId(), + possibleSendToDeadLetterTopicMessages.put((MessageIdImpl) message.getMessageId(), Collections.singletonList(message)); } if (peekPendingReceive() != null) { From c1517ffbe21f767a1147633f71e0bc0153b5a7db Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 29 Apr 2021 21:12:06 +0800 Subject: [PATCH 11/12] Apply suggestion --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 2 -- 1 file changed, 2 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 5e987369b1e01..b5e790df2fa8b 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 @@ -130,8 +130,6 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final int receiverQueueRefillThreshold; - private final ReadWriteLock lock = new ReentrantReadWriteLock(); - private final UnAckedMessageTracker unAckedMessageTracker; private final AcknowledgmentsGroupingTracker acknowledgmentsGroupingTracker; private final NegativeAcksTracker negativeAcksTracker; From 9773d7f1e5634f1340c702b03488e94565fd7664 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 4 May 2021 18:23:18 +0800 Subject: [PATCH 12/12] move triggerListener --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 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 b5e790df2fa8b..90d037cc10984 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 @@ -1066,7 +1066,6 @@ uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redelive } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { notifyPendingBatchReceivedCallBack(); } - tryTriggerListener(); }); } else { // handle batch message enqueuing; uncompressed payload has all messages in batch @@ -1074,6 +1073,8 @@ uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redelive uncompressedPayload.release(); } + internalPinnedExecutor.execute(() + -> tryTriggerListener()); } @@ -1287,7 +1288,6 @@ void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata, int redeliv notifyPendingBatchReceivedCallBack(); } singleMessagePayload.release(); - tryTriggerListener(); }); } if (ackBitSet != null) {