From 3aafb0977009045685a8a43c1d6afd08f487db44 Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Thu, 14 Jul 2022 15:55:36 +0800 Subject: [PATCH 1/2] [fix][flaky-test] org.apache.pulsar.client.impl.PatternTopicsConsumerImplTest#testAutoSubscribePatternConsumer --- .../impl/PatternTopicsConsumerImplTest.java | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) 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 65c05341b7726..e7743e07e255c 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 @@ -439,7 +439,7 @@ public void testStartEmptyPatternConsumer() throws Exception { // 2. Create consumer, this should success, but with empty sub-consumser internal Consumer consumer = pulsarClient.newConsumer() .topicsPattern(pattern) - .patternAutoDiscoveryPeriod(2) + .patternAutoDiscoveryPeriod(2, TimeUnit.SECONDS) .subscriptionName(subscriptionName) .subscriptionType(SubscriptionType.Shared) .ackTimeout(ackTimeOutMillis, TimeUnit.MILLISECONDS) @@ -473,13 +473,14 @@ public void testStartEmptyPatternConsumer() throws Exception { log.debug("recheck topics change"); PatternMultiTopicsConsumerImpl consumer1 = ((PatternMultiTopicsConsumerImpl) consumer); consumer1.run(consumer1.getRecheckPatternTimeout()); - Thread.sleep(100); // 6. verify consumer get methods, to get number of partitions and topics, value 6=1+2+3. - assertSame(pattern, ((PatternMultiTopicsConsumerImpl) consumer).getPattern()); - assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitions().size(), 6); - assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getConsumers().size(), 6); - assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitionedTopics().size(), 2); + Awaitility.await().untilAsserted(() -> { + assertSame(pattern, ((PatternMultiTopicsConsumerImpl) consumer).getPattern()); + assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitions().size(), 6); + assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getConsumers().size(), 6); + assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitionedTopics().size(), 2); + }); // 7. produce data @@ -543,7 +544,7 @@ public void testAutoSubscribePatternConsumer() throws Exception { Consumer consumer = pulsarClient.newConsumer() .topicsPattern(pattern) - .patternAutoDiscoveryPeriod(2) + .patternAutoDiscoveryPeriod(2, TimeUnit.SECONDS) .subscriptionName(subscriptionName) .subscriptionType(SubscriptionType.Shared) .ackTimeout(ackTimeOutMillis, TimeUnit.MILLISECONDS) From 9d8a5d75b97efb64cf377e1357ee06a46ce71f94 Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Thu, 14 Jul 2022 18:05:55 +0800 Subject: [PATCH 2/2] fix error --- .../client/impl/PatternTopicsConsumerImplTest.java | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) 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 e7743e07e255c..9b3ff40113d78 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 @@ -439,7 +439,7 @@ public void testStartEmptyPatternConsumer() throws Exception { // 2. Create consumer, this should success, but with empty sub-consumser internal Consumer consumer = pulsarClient.newConsumer() .topicsPattern(pattern) - .patternAutoDiscoveryPeriod(2, TimeUnit.SECONDS) + .patternAutoDiscoveryPeriod(2) .subscriptionName(subscriptionName) .subscriptionType(SubscriptionType.Shared) .ackTimeout(ackTimeOutMillis, TimeUnit.MILLISECONDS) @@ -469,6 +469,11 @@ public void testStartEmptyPatternConsumer() throws Exception { .messageRoutingMode(org.apache.pulsar.client.api.MessageRoutingMode.RoundRobinPartition) .create(); + List topicNames = Lists.newArrayList(topicName1, topicName2, topicName3); + NamespaceService nss = pulsar.getNamespaceService(); + doReturn(CompletableFuture.completedFuture(topicNames)).when(nss) + .getListOfPersistentTopics(NamespaceName.get("my-property/my-ns")); + // 5. call recheckTopics to subscribe each added topics above log.debug("recheck topics change"); PatternMultiTopicsConsumerImpl consumer1 = ((PatternMultiTopicsConsumerImpl) consumer); @@ -544,7 +549,7 @@ public void testAutoSubscribePatternConsumer() throws Exception { Consumer consumer = pulsarClient.newConsumer() .topicsPattern(pattern) - .patternAutoDiscoveryPeriod(2, TimeUnit.SECONDS) + .patternAutoDiscoveryPeriod(2) .subscriptionName(subscriptionName) .subscriptionType(SubscriptionType.Shared) .ackTimeout(ackTimeOutMillis, TimeUnit.MILLISECONDS) @@ -585,6 +590,11 @@ public void testAutoSubscribePatternConsumer() throws Exception { .messageRoutingMode(org.apache.pulsar.client.api.MessageRoutingMode.RoundRobinPartition) .create(); + List topicNames = Lists.newArrayList(topicName1, topicName2, topicName3, topicName4); + NamespaceService nss = pulsar.getNamespaceService(); + doReturn(CompletableFuture.completedFuture(topicNames)).when(nss) + .getListOfPersistentTopics(NamespaceName.get("my-property/my-ns")); + // 7. call recheckTopics to subscribe each added topics above, verify topics number: 10=1+2+3+4 log.debug("recheck topics change"); PatternMultiTopicsConsumerImpl consumer1 = ((PatternMultiTopicsConsumerImpl) consumer);