diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index 3dc5a1014cb67..0541c7a498534 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -144,6 +144,10 @@ public boolean cacheIsInitialized(TopicName topicName) { @Override public TopicPolicies getTopicPolicies(TopicName topicName) throws TopicPoliciesCacheNotInitException { + if (!policyCacheInitMap.containsKey(topicName.getNamespaceObject())) { + NamespaceName namespace = topicName.getNamespaceObject(); + prepareInitPoliciesCache(namespace, new CompletableFuture<>()); + } if (policyCacheInitMap.containsKey(topicName.getNamespaceObject()) && !policyCacheInitMap.get(topicName.getNamespaceObject())) { throw new TopicPoliciesCacheNotInitException(); @@ -151,6 +155,25 @@ public TopicPolicies getTopicPolicies(TopicName topicName) throws TopicPoliciesC return policiesCache.get(TopicName.get(topicName.getPartitionedTopicName())); } + private void prepareInitPoliciesCache(NamespaceName namespace, CompletableFuture result) { + if (policyCacheInitMap.putIfAbsent(namespace, false) == null) { + CompletableFuture readerCompletableFuture = namespaceEventsSystemTopicFactory + .createTopicPoliciesSystemTopicClient(namespace).newReaderAsync(); + readerCaches.put(namespace, readerCompletableFuture); + readerCompletableFuture.whenComplete((reader, ex) -> { + if (ex != null) { + log.error("[{}] Failed to create reader on __change_events topic", namespace, ex); + result.completeExceptionally(ex); + readerCaches.remove(namespace); + reader.closeAsync(); + } else { + initPolicesCache(reader, result); + result.thenRun(() -> readMorePolicies(reader)); + } + }); + } + } + @Override public CompletableFuture getTopicPoliciesBypassCacheAsync(TopicName topicName) { CompletableFuture result = new CompletableFuture<>(); @@ -180,21 +203,8 @@ public CompletableFuture addOwnedNamespaceBundleAsync(NamespaceBundle name ownedBundlesCountPerNamespace.get(namespace).incrementAndGet(); result.complete(null); } else { - SystemTopicClient systemTopicClient = namespaceEventsSystemTopicFactory.createSystemTopic(namespace - , EventType.TOPIC_POLICY); ownedBundlesCountPerNamespace.putIfAbsent(namespace, new AtomicInteger(1)); - policyCacheInitMap.put(namespace, false); - CompletableFuture readerCompletableFuture = systemTopicClient.newReaderAsync(); - readerCaches.put(namespace, readerCompletableFuture); - readerCompletableFuture.whenComplete((reader, ex) -> { - if (ex != null) { - log.error("[{}] Failed to create reader on __change_events topic", namespace, ex); - result.completeExceptionally(ex); - } else { - initPolicesCache(reader, result); - result.thenRun(() -> readMorePolicies(reader)); - } - }); + prepareInitPoliciesCache(namespace, result); } } return result; @@ -249,14 +259,20 @@ private void initPolicesCache(SystemTopicClient.Reader reader, CompletableFuture reader.getSystemTopic().getTopicName(), ex); future.completeExceptionally(ex); readerCaches.remove(reader.getSystemTopic().getTopicName().getNamespaceObject()); + policyCacheInitMap.remove(reader.getSystemTopic().getTopicName().getNamespaceObject()); + reader.closeAsync(); + return; } if (hasMore) { reader.readNextAsync().whenComplete((msg, e) -> { if (e != null) { log.error("[{}] Failed to read event from the system topic.", - reader.getSystemTopic().getTopicName(), ex); + reader.getSystemTopic().getTopicName(), e); future.completeExceptionally(e); readerCaches.remove(reader.getSystemTopic().getTopicName().getNamespaceObject()); + policyCacheInitMap.remove(reader.getSystemTopic().getTopicName().getNamespaceObject()); + reader.closeAsync(); + return; } refreshTopicPoliciesCache(msg); if (log.isDebugEnabled()) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java index 9eb866980bfd1..c8068404e2868 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java @@ -22,6 +22,7 @@ import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest; import org.apache.pulsar.broker.service.BrokerServiceException.TopicPoliciesCacheNotInitException; import org.apache.pulsar.broker.systopic.NamespaceEventsSystemTopicFactory; +import org.apache.pulsar.broker.systopic.SystemTopicClient; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; @@ -33,8 +34,13 @@ import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; - +import java.lang.reflect.Field; +import java.util.Map; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; public class SystemTopicBasedTopicPoliciesServiceTest extends MockedPulsarServiceBaseTest { @@ -193,4 +199,33 @@ private void prepareData() throws PulsarAdminException { systemTopicFactory = new NamespaceEventsSystemTopicFactory(pulsarClient); systemTopicBasedTopicPoliciesService = (SystemTopicBasedTopicPoliciesService) pulsar.getTopicPoliciesService(); } + + @Test + public void testGetTopicPoliciesWithRetry() throws Exception { + Field initMapField = SystemTopicBasedTopicPoliciesService.class.getDeclaredField("policyCacheInitMap"); + initMapField.setAccessible(true); + Map initMap = (Map)initMapField.get(systemTopicBasedTopicPoliciesService); + initMap.remove(NamespaceName.get(NAMESPACE1)); + Field readerCaches = SystemTopicBasedTopicPoliciesService.class.getDeclaredField("readerCaches"); + readerCaches.setAccessible(true); + Map> readers = (Map)readerCaches.get(systemTopicBasedTopicPoliciesService); + readers.remove(NamespaceName.get(NAMESPACE1)); + TopicPolicies initPolicy = TopicPolicies.builder() + .maxConsumerPerTopic(10) + .build(); + ScheduledExecutorService executors = Executors.newScheduledThreadPool(1); + executors.schedule(() -> { + try { + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy).get(); + } catch (Exception ignore) {} + }, 2000, TimeUnit.MILLISECONDS); + Awaitility.await().untilAsserted(() -> { + try { + TopicPolicies topicPolicies = systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1); + Assert.assertNotNull(topicPolicies); + } catch (Exception ex) { + Assert.fail(); + } + }); + } }