From e14a5c515884c98fb6a814bfb2aff1864234c3e0 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Tue, 11 Jul 2023 16:50:05 +0800 Subject: [PATCH 1/2] Disable polling pattern topics when TopicListWatcher is enabled. --- .../pulsar/client/impl/TopicsConsumerImplTest.java | 1 + .../client/impl/PatternMultiTopicsConsumerImpl.java | 10 ++++++---- 2 files changed, 7 insertions(+), 4 deletions(-) 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); From b7b350dff52ce8c6ef78cb09fb898d2ccadfd2a3 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Fri, 14 Jul 2023 14:59:01 +0800 Subject: [PATCH 2/2] fix test. --- .../pulsar/client/impl/PatternTopicsConsumerImplTest.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) 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);