From 3b8fb93be8f94b1f2fac9546d87fbd4f222c468d Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Sat, 19 Aug 2023 12:18:51 +0800 Subject: [PATCH 1/3] Ensure that the cache is updated before updateTopicPoliciesAsync returns. --- .../SystemTopicBasedTopicPoliciesService.java | 24 ++++++++++++++- ...temTopicBasedTopicPoliciesServiceTest.java | 30 +++++++++++++++++++ 2 files changed, 53 insertions(+), 1 deletion(-) 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 09f8de818db0a..ba03d394725c0 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 @@ -24,9 +24,11 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.BlockingDeque; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.LinkedBlockingDeque; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import javax.annotation.Nonnull; @@ -84,6 +86,8 @@ public class SystemTopicBasedTopicPoliciesService implements TopicPoliciesServic @VisibleForTesting final Map>> listeners = new ConcurrentHashMap<>(); + final Map> results = new ConcurrentHashMap<>(); + public SystemTopicBasedTopicPoliciesService(PulsarService pulsarService) { this.pulsarService = pulsarService; this.clusterName = pulsarService.getConfiguration().getClusterName(); @@ -97,7 +101,10 @@ public CompletableFuture deleteTopicPoliciesAsync(TopicName topicName) { @Override public CompletableFuture updateTopicPoliciesAsync(TopicName topicName, TopicPolicies policies) { - return sendTopicPolicyEvent(topicName, ActionType.UPDATE, policies); + CompletableFuture result = new CompletableFuture<>(); + results.computeIfAbsent(topicName, k -> new LinkedBlockingDeque<>()).add(result); + return sendTopicPolicyEvent(topicName, ActionType.UPDATE, policies) + .thenCombine(result, (__, ___) -> null); } private CompletableFuture sendTopicPolicyEvent(TopicName topicName, ActionType actionType, @@ -420,11 +427,26 @@ private void cleanCacheAndCloseReader(@Nonnull NamespaceName namespace, boolean }); } + private void completeResults(TopicName topicName) { + if (results.get(topicName) != null) { + BlockingDeque futures = results.get(topicName); + CompletableFuture future; + while (true) { + future = futures.poll(); + if (future == null) { + break; + } + future.complete(null); + } + } + } + private void readMorePolicies(SystemTopicClient.Reader reader) { reader.readNextAsync() .thenAccept(msg -> { refreshTopicPoliciesCache(msg); notifyListener(msg); + completeResults(TopicName.get(TopicName.get(msg.getKey()).getPartitionedTopicName())); }) .whenComplete((__, ex) -> { if (ex == null) { 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 31b5bcb23cd98..99fe3b2498efc 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 @@ -239,6 +239,36 @@ public void testGetPolicy() throws ExecutionException, InterruptedException, Top Assert.assertEquals(policies1, policiesGet1); } + @Test + public void testUpdateAndGetPolicy() throws ExecutionException, InterruptedException, TopicPoliciesCacheNotInitException { + // Init topic policies + TopicPolicies initPolicy = TopicPolicies.builder() + .maxConsumerPerTopic(10) + .build(); + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy).get(); + + // Wait for all topic policies updated. + Awaitility.await().untilAsserted(() -> + Assert.assertTrue(systemTopicBasedTopicPoliciesService + .getPoliciesCacheInit(TOPIC1.getNamespaceObject()))); + + // Assert broker is cache all topic policies + Awaitility.await().untilAsserted(() -> + Assert.assertEquals(systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1) + .getMaxConsumerPerTopic().intValue(), 10)); + + // Update policy for TOPIC1 + TopicPolicies policies1 = TopicPolicies.builder() + .maxProducerPerTopic(1) + .build(); + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, policies1).get(); + + // Ensure that the cache is updated before updateTopicPoliciesAsync returns. + TopicPolicies policiesGet1 = systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1); + Assert.assertEquals(policiesGet1, policies1); + } + + @Test public void testCacheCleanup() throws Exception { final String topic = "persistent://" + NAMESPACE1 + "/test" + UUID.randomUUID(); From 7a91103f28afdae6f70011ace2fab691b2ff3f43 Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 21 Aug 2023 10:40:00 +0800 Subject: [PATCH 2/3] remove the result blocking queue if topic is removed. --- .../SystemTopicBasedTopicPoliciesService.java | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) 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 ba03d394725c0..a465d7485888b 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 @@ -441,6 +441,17 @@ private void completeResults(TopicName topicName) { } } + private void completeResults(BlockingDeque futures) { + CompletableFuture future; + while (true) { + future = futures.poll(); + if (future == null) { + break; + } + future.complete(null); + } + } + private void readMorePolicies(SystemTopicClient.Reader reader) { reader.readNextAsync() .thenAccept(msg -> { @@ -649,6 +660,11 @@ public void unregisterListener(TopicName topicName, TopicPolicyListener futures = results.remove(topicName); + if (futures != null) { + completeResults(futures); + } } } return topicListeners; From 73393a6aa0a3a2b92ec8abcbb603334c60c9426f Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 21 Aug 2023 16:22:16 +0800 Subject: [PATCH 3/3] add test code. --- ...temTopicBasedTopicPoliciesServiceTest.java | 38 +++++++++++++++++++ 1 file changed, 38 insertions(+) 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 99fe3b2498efc..09d72235f8424 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 @@ -268,6 +268,44 @@ public void testUpdateAndGetPolicy() throws ExecutionException, InterruptedExcep Assert.assertEquals(policiesGet1, policies1); } + @Test + public void testGetPolicyAndRemoveTopic() throws ExecutionException, InterruptedException { + class TopicPolicyListenerImpl implements TopicPolicyListener { + + @Override + public void onUpdate(TopicPolicies data) { + //no op. + } + } + TopicPolicyListener listener = new TopicPolicyListenerImpl(); + systemTopicBasedTopicPoliciesService.registerListener(TOPIC1, listener); + Assert.assertTrue(systemTopicBasedTopicPoliciesService.listeners.get(TOPIC1).size() >= 1); + + // Init topic policies + TopicPolicies initPolicy = TopicPolicies.builder() + .maxConsumerPerTopic(10) + .build(); + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy).get(); + + // Wait for all topic policies updated. + Awaitility.await().untilAsserted(() -> + Assert.assertTrue(systemTopicBasedTopicPoliciesService + .getPoliciesCacheInit(TOPIC1.getNamespaceObject()))); + + // Assert broker cache all topic policies. + Awaitility.await().untilAsserted(() -> + Assert.assertEquals(systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1) + .getMaxConsumerPerTopic().intValue(), 10)); + + // Assert policy setting result queue is created. + Assert.assertTrue(systemTopicBasedTopicPoliciesService.results.get(TOPIC1) != null); + + // unregister listener and remove the policy setting result queue. + systemTopicBasedTopicPoliciesService.unregisterListener(TOPIC1, listener); + Assert.assertFalse(systemTopicBasedTopicPoliciesService.listeners.containsKey(TOPIC1)); + Assert.assertFalse(systemTopicBasedTopicPoliciesService.results.containsKey(TOPIC1)); + } + @Test public void testCacheCleanup() throws Exception {