From ffee00b41272cd38313d02c31939af3be0a58150 Mon Sep 17 00:00:00 2001 From: chenhang Date: Mon, 10 May 2021 20:12:15 +0800 Subject: [PATCH 1/2] fix partitioned system topic check bug --- .../pulsar/broker/systopic/SystemTopicClient.java | 5 ++++- .../NamespaceEventsSystemTopicServiceTest.java | 14 ++++++++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java index 4422d2f3a6baa..27e20e8b66058 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java @@ -168,7 +168,10 @@ interface Reader { } static boolean isSystemTopic(TopicName topicName) { - return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME.equals(topicName.getLocalName()); + String localName = topicName.getLocalName(); + int partitionIndex = localName.indexOf(TopicName.PARTITIONED_TOPIC_SUFFIX); + String topicNameWithoutSuffix = partitionIndex == -1 ? localName : localName.substring(0, partitionIndex); + return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME.equals(topicNameWithoutSuffix); } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/NamespaceEventsSystemTopicServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/NamespaceEventsSystemTopicServiceTest.java index abbab205df670..eb08913fca727 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/NamespaceEventsSystemTopicServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/NamespaceEventsSystemTopicServiceTest.java @@ -24,9 +24,11 @@ import org.apache.pulsar.client.api.Message; import org.apache.pulsar.common.events.ActionType; import org.apache.pulsar.common.events.EventType; +import org.apache.pulsar.common.events.EventsTopicNames; import org.apache.pulsar.common.events.PulsarEvent; import org.apache.pulsar.common.events.TopicPoliciesEvent; import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.policies.data.TopicPolicies; @@ -108,6 +110,18 @@ public void testSendAndReceiveNamespaceEvents() throws Exception { Assert.assertEquals(systemTopicClientForNamespace1.getReaders().size(), 0); } + @Test(timeOut = 30000) + public void checkSystemTopic() throws PulsarAdminException { + final String systemTopic = "persistent://" + NAMESPACE1 + "/" + EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME; + final String normalTopic = "persistent://" + NAMESPACE1 + "/normal_topic"; + admin.topics().createPartitionedTopic(normalTopic, 3); + TopicName systemTopicName = TopicName.get(systemTopic); + TopicName normalTopicName = TopicName.get(normalTopic); + + Assert.assertEquals(SystemTopicClient.isSystemTopic(systemTopicName), true); + Assert.assertEquals(SystemTopicClient.isSystemTopic(normalTopicName), false); + } + private void prepareData() throws PulsarAdminException { admin.clusters().createCluster("test", new ClusterData(pulsar.getBrokerServiceUrl())); admin.tenants().createTenant("system-topic", From 3f70d8b7b9f18dbf60a2c77d49d5745ed7dc8b6f Mon Sep 17 00:00:00 2001 From: chenhang Date: Tue, 11 May 2021 10:18:26 +0800 Subject: [PATCH 2/2] update code --- .../pulsar/broker/systopic/SystemTopicClient.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java index 27e20e8b66058..3f5a0a94cd977 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java @@ -168,10 +168,12 @@ interface Reader { } static boolean isSystemTopic(TopicName topicName) { - String localName = topicName.getLocalName(); - int partitionIndex = localName.indexOf(TopicName.PARTITIONED_TOPIC_SUFFIX); - String topicNameWithoutSuffix = partitionIndex == -1 ? localName : localName.substring(0, partitionIndex); - return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME.equals(topicNameWithoutSuffix); + if (topicName.isPartitioned()) { + return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME + .equals(TopicName.get(topicName.getPartitionedTopicName()).getLocalName()); + } + + return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME.equals(topicName.getLocalName()); } }