diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 7bdf1f3491c23..4cc8ae178c3e0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -3160,7 +3160,8 @@ protected void handleCommandWatchTopicList(CommandWatchTopicList commandWatchTop final NamespaceName namespaceName = NamespaceName.get(commandWatchTopicList.getNamespace()); Pattern topicsPattern = Pattern.compile(commandWatchTopicList.hasTopicsPattern() - ? commandWatchTopicList.getTopicsPattern() : TopicList.ALL_TOPICS_PATTERN); + ? TopicList.removeTopicDomainScheme(commandWatchTopicList.getTopicsPattern()) + : TopicList.ALL_TOPICS_PATTERN); String topicsHash = commandWatchTopicList.hasTopicsHash() ? commandWatchTopicList.getTopicsHash() : null; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicListService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicListService.java index e04d07460a2cb..818188fd1829e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicListService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TopicListService.java @@ -70,7 +70,8 @@ public List getMatchingTopics() { */ @Override public void accept(String topicName, NotificationType notificationType) { - if (topicsPattern.matcher(TopicName.get(topicName).getPartitionedTopicName()).matches()) { + String partitionedTopicName = TopicName.get(topicName).getPartitionedTopicName(); + if (topicsPattern.matcher(TopicList.removeTopicDomainScheme(partitionedTopicName)).matches()) { List newTopics; List deletedTopics; if (notificationType == NotificationType.Deleted) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicListWatcherTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicListWatcherTest.java index 641b1bd4e74b3..23bb884d25e47 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicListWatcherTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicListWatcherTest.java @@ -40,7 +40,7 @@ public class TopicListWatcherTest { ); private static final long ID = 7; - private static final Pattern PATTERN = Pattern.compile("persistent://tenant/ns/topic\\d+"); + private static final Pattern PATTERN = Pattern.compile("tenant/ns/topic\\d+"); private TopicListService topicListService; 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 4823426c8b83a..775ea4c3dcb9a 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 @@ -684,6 +684,40 @@ public void testAutoSubscribePatterConsumerFromBrokerWatcher(boolean delayWatchi } } + @Test(timeOut = testTimeout) + public void testSubscribePatterWithOutTopicDomain() throws Exception { + final String key = "testSubscribePatterWithOutTopicDomain"; + final String subscriptionName = "my-ex-subscription-" + key; + final Pattern pattern = Pattern.compile("my-property/my-ns/test-pattern.*"); + + Consumer consumer = pulsarClient.newConsumer() + .topicsPattern(pattern) + .subscriptionName(subscriptionName) + .subscriptionType(SubscriptionType.Shared) + .receiverQueueSize(4) + .subscribe(); + + // 0. Need make sure topic watcher started + waitForTopicListWatcherStarted(consumer); + + // 1. create partition topic + String topicName = "persistent://my-property/my-ns/test-pattern" + key; + admin.topics().createPartitionedTopic(topicName, 4); + + // 2. verify broker will push the changes to update(CommandWatchTopicUpdate). + assertSame(pattern.pattern(), ((PatternMultiTopicsConsumerImpl) consumer).getPattern().pattern()); + Awaitility.await().atMost(Duration.ofSeconds(5)).untilAsserted(() -> { + assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitions().size(), 4); + assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getConsumers().size(), 4); + assertEquals(((PatternMultiTopicsConsumerImpl) consumer).getPartitionedTopics().size(), 1); + }); + + // cleanup. + consumer.close(); + admin.topics().deletePartitionedTopic(topicName); + pulsarClient.close(); + } + @DataProvider(name= "regexpConsumerArgs") public Object[][] regexpConsumerArgs(){ return new Object[][]{ diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java index ed77652c82340..8eaf5ca969f67 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java @@ -128,6 +128,8 @@ public interface ConsumerBuilder extends Cloneable { /** * Specify a pattern for topics(not contains the partition suffix) that this consumer subscribes to. * + *

Will ignore the topic domain("persistent://" or "non-persistent://") when pattern matching. + * *

The pattern is applied to subscribe to all topics, within a single namespace, that match the * pattern. * @@ -143,7 +145,9 @@ public interface ConsumerBuilder extends Cloneable { * Specify a pattern for topics(not contains the partition suffix) that this consumer subscribes to. * *

It accepts a regular expression that is compiled into a pattern internally. E.g., - * "persistent://public/default/pattern-topic-.*" + * "persistent://public/default/pattern-topic-.*" or "public/default/pattern-topic-.*" + * + *

Will ignore the topic domain("persistent://" or "non-persistent://") when pattern matching. * *

The pattern is applied to subscribe to all topics, within a single namespace, that match the * pattern. diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/topics/TopicList.java b/pulsar-common/src/main/java/org/apache/pulsar/common/topics/TopicList.java index 9e24483df8239..98e4d0865030a 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/topics/TopicList.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/topics/TopicList.java @@ -18,7 +18,6 @@ */ package org.apache.pulsar.common.topics; -import com.google.common.annotations.VisibleForTesting; import com.google.common.hash.Hashing; import com.google.re2j.Pattern; import java.nio.charset.StandardCharsets; @@ -85,8 +84,7 @@ public static Set minus(Collection list1, Collection lis return s1; } - @VisibleForTesting - static String removeTopicDomainScheme(String originalRegexp) { + public static String removeTopicDomainScheme(String originalRegexp) { if (!originalRegexp.toString().contains(SCHEME_SEPARATOR)) { return originalRegexp; }