diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java index 09a9a003f3ec5..d1c569565bcdd 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/PatternTopicsConsumerImplTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.client.impl; import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.spy; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertSame; import static org.testng.Assert.assertTrue; @@ -33,6 +34,7 @@ import java.util.regex.Pattern; import java.util.stream.IntStream; +import io.netty.util.Timeout; import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; @@ -809,7 +811,9 @@ public void testAutoUnsubscribePatternConsumer() throws Exception { // 7. call recheckTopics to unsubscribe topic 1,3, verify topics number: 2=6-1-3 log.debug("recheck topics change"); PatternMultiTopicsConsumerImpl consumer1 = ((PatternMultiTopicsConsumerImpl) consumer); - consumer1.run(consumer1.getRecheckPatternTimeout()); + Timeout recheckPatternTimeout = spy(consumer1.getRecheckPatternTimeout()); + doReturn(false).when(recheckPatternTimeout).isCancelled(); + consumer1.run(recheckPatternTimeout); Thread.sleep(100); assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitions().size(), 2); assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getConsumers().size(), 2); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicsConsumerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicsConsumerImplTest.java index 73fe97996424c..51b32c2b44ecf 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicsConsumerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicsConsumerImplTest.java @@ -1327,6 +1327,7 @@ public void testPartitionsUpdatesForMultipleTopics() throws Exception { assertEquals(admin.topics().getPartitionedTopicMetadata(topicName1).partitions, 3); consumer.getRecheckPatternTimeout().task().run(consumer.getRecheckPatternTimeout()); + Assert.assertTrue(consumer.getRecheckPatternTimeout().isCancelled()); Awaitility.await().untilAsserted(() -> { Assert.assertEquals(consumer.getPartitionsOfTheTopicMap(), 8); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java index 12c7e4e4ba3c7..c6ea6216cc1f4 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/PatternMultiTopicsConsumerImpl.java @@ -83,10 +83,12 @@ public PatternMultiTopicsConsumerImpl(Pattern topicsPattern, long watcherId = client.newTopicListWatcherId(); new TopicListWatcher(topicsChangeListener, client, topicsPattern, watcherId, namespaceName, topicsHash, watcherFuture); - watcherFuture.exceptionally(ex -> { - log.debug("Unable to create topic list watcher. Falling back to only polling for new topics", ex); - return null; - }); + watcherFuture + .thenAccept(__ -> recheckPatternTimeout.cancel()) + .exceptionally(ex -> { + log.warn("Unable to create topic list watcher. Falling back to only polling for new topics", ex); + return null; + }); } else { log.debug("Not creating topic list watcher for subscription mode {}", subscriptionMode); watcherFuture.complete(null);