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..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,6 +168,11 @@ interface Reader { } static boolean isSystemTopic(TopicName topicName) { + if (topicName.isPartitioned()) { + return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME + .equals(TopicName.get(topicName.getPartitionedTopicName()).getLocalName()); + } + return EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME.equals(topicName.getLocalName()); } 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",