From 6ee175b830571a68341984e8b24ca2d1eb6d9787 Mon Sep 17 00:00:00 2001 From: Masahiro Sakamoto Date: Thu, 30 Jan 2020 15:18:48 +0900 Subject: [PATCH] Restore clusterDispatchRate policy for compatibility --- .../broker/admin/impl/NamespacesBase.java | 6 ++ .../persistent/DispatchRateLimiter.java | 8 ++ .../api/MessageDispatchThrottlingTest.java | 83 ++++++++++++++++++- .../pulsar/common/policies/data/Policies.java | 6 +- 4 files changed, 98 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index 33592e6f99612..f28205130ec17 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -855,6 +855,7 @@ protected PublishRate internalGetPublishRate() { } } + @SuppressWarnings("deprecation") protected void internalSetTopicDispatchRate(DispatchRate dispatchRate) { log.info("[{}] Set namespace dispatch-rate {}/{}", clientAppId(), namespaceName, dispatchRate); validateSuperUserAccess(); @@ -867,6 +868,7 @@ protected void internalSetTopicDispatchRate(DispatchRate dispatchRate) { policiesNode = policiesCache().getWithStat(path).orElseThrow( () -> new RestException(Status.NOT_FOUND, "Namespace " + namespaceName + " does not exist")); policiesNode.getKey().topicDispatchRate.put(pulsar().getConfiguration().getClusterName(), dispatchRate); + policiesNode.getKey().clusterDispatchRate.put(pulsar().getConfiguration().getClusterName(), dispatchRate); // Write back the new policies into zookeeper globalZk().setData(path, jsonMapper().writeValueAsBytes(policiesNode.getKey()), @@ -892,11 +894,15 @@ protected void internalSetTopicDispatchRate(DispatchRate dispatchRate) { } } + @SuppressWarnings("deprecation") protected DispatchRate internalGetTopicDispatchRate() { validateAdminAccessForTenant(namespaceName.getTenant()); Policies policies = getNamespacePolicies(namespaceName); DispatchRate dispatchRate = policies.topicDispatchRate.get(pulsar().getConfiguration().getClusterName()); + if (dispatchRate == null) { + dispatchRate = policies.clusterDispatchRate.get(pulsar().getConfiguration().getClusterName()); + } if (dispatchRate != null) { return dispatchRate; } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java index 3f437ddc81566..29eed7305f0d2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/DispatchRateLimiter.java @@ -186,6 +186,7 @@ public static boolean isDispatchRateNeeded(final ServiceConfiguration serviceCon return true; } + @SuppressWarnings("deprecation") public void onPoliciesUpdate(Policies data) { String cluster = brokerService.pulsar().getConfiguration().getClusterName(); @@ -194,6 +195,9 @@ public void onPoliciesUpdate(Policies data) { switch (type) { case TOPIC: dispatchRate = data.topicDispatchRate.get(cluster); + if (dispatchRate == null) { + dispatchRate = data.clusterDispatchRate.get(cluster); + } break; case SUBSCRIPTION: dispatchRate = data.subscriptionDispatchRate.get(cluster); @@ -219,6 +223,7 @@ public void onPoliciesUpdate(Policies data) { } } + @SuppressWarnings("deprecation") public static DispatchRate getPoliciesDispatchRate(final String cluster, Optional policies, Type type) { // return policy-dispatch rate only if it's enabled in policies return policies.map(p -> { @@ -226,6 +231,9 @@ public static DispatchRate getPoliciesDispatchRate(final String cluster, Optiona switch (type) { case TOPIC: dispatchRate = p.topicDispatchRate.get(cluster); + if (dispatchRate == null) { + dispatchRate = p.clusterDispatchRate.get(cluster); + } break; case SUBSCRIPTION: dispatchRate = p.subscriptionDispatchRate.get(cluster); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MessageDispatchThrottlingTest.java index ff93d3aed297d..d75fcfd0e581b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/MessageDispatchThrottlingTest.java @@ -18,14 +18,17 @@ */ package org.apache.pulsar.client.api; -import com.google.common.collect.Sets; - import static org.testng.Assert.assertNotNull; +import com.google.common.collect.Maps; +import com.google.common.collect.Sets; + import java.lang.reflect.Field; import java.util.Arrays; import java.util.LinkedList; import java.util.List; +import java.util.Map; +import java.util.Optional; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -37,6 +40,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.DispatchRate; +import org.apache.pulsar.common.policies.data.Policies; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -101,6 +105,7 @@ enum DispatchRateType { * * @throws Exception */ + @SuppressWarnings("deprecation") @Test public void testMessageRateDynamicallyChange() throws Exception { @@ -116,7 +121,7 @@ public void testMessageRateDynamicallyChange() throws Exception { // (1) verify message-rate is -1 initially Assert.assertFalse(topic.getDispatchRateLimiter().isPresent()); - // (1) change to 100 + // (2) change to 100 int messageRate = 100; DispatchRate dispatchRate = new DispatchRate(messageRate, -1, 360); admin.namespaces().setDispatchRate(namespace, dispatchRate); @@ -134,8 +139,13 @@ public void testMessageRateDynamicallyChange() throws Exception { } Assert.assertTrue(isDispatchRateUpdate); Assert.assertEquals(admin.namespaces().getDispatchRate(namespace), dispatchRate); + Policies policies = admin.namespaces().getPolicies(namespace); + Map dispatchRateMap = Maps.newHashMap(); + dispatchRateMap.put("test", dispatchRate); + Assert.assertEquals(policies.clusterDispatchRate, dispatchRateMap); + Assert.assertEquals(policies.topicDispatchRate, dispatchRateMap); - // (1) change to 500 + // (3) change to 500 messageRate = 500; dispatchRate = new DispatchRate(-1, messageRate, 360); admin.namespaces().setDispatchRate(namespace, dispatchRate); @@ -152,6 +162,10 @@ public void testMessageRateDynamicallyChange() throws Exception { } Assert.assertTrue(isDispatchRateUpdate); Assert.assertEquals(admin.namespaces().getDispatchRate(namespace), dispatchRate); + policies = admin.namespaces().getPolicies(namespace); + dispatchRateMap.put("test", dispatchRate); + Assert.assertEquals(policies.clusterDispatchRate, dispatchRateMap); + Assert.assertEquals(policies.topicDispatchRate, dispatchRateMap); producer.close(); } @@ -896,6 +910,67 @@ public void testClosingRateLimiter(SubscriptionType subscription) throws Excepti log.info("-- Exiting {} test --", methodName); } + @SuppressWarnings("deprecation") + @Test + public void testDispatchRateCompatibility1() throws Exception { + final String cluster = "test"; + + Optional policies = Optional.of(new Policies()); + DispatchRate clusterDispatchRate = new DispatchRate(100, 512, 1); + DispatchRate topicDispatchRate = new DispatchRate(200, 1024, 1); + + // (1) If both clusterDispatchRate and topicDispatchRate are empty, dispatch throttling is disabled + DispatchRate dispatchRate = DispatchRateLimiter.getPoliciesDispatchRate(cluster, policies, + DispatchRateLimiter.Type.TOPIC); + Assert.assertNull(dispatchRate); + + // (2) If topicDispatchRate is empty, clusterDispatchRate is effective + policies.get().clusterDispatchRate.put(cluster, clusterDispatchRate); + dispatchRate = DispatchRateLimiter.getPoliciesDispatchRate(cluster, policies, DispatchRateLimiter.Type.TOPIC); + Assert.assertEquals(dispatchRate, clusterDispatchRate); + + // (3) If topicDispatchRate is not empty, topicDispatchRate is effective + policies.get().topicDispatchRate.put(cluster, topicDispatchRate); + dispatchRate = DispatchRateLimiter.getPoliciesDispatchRate(cluster, policies, DispatchRateLimiter.Type.TOPIC); + Assert.assertEquals(dispatchRate, topicDispatchRate); + } + + @SuppressWarnings("deprecation") + @Test + public void testDispatchRateCompatibility2() throws Exception { + final String namespace = "my-property/dispatch-rate-compatibility"; + final String topicName = "persistent://" + namespace + "/t1"; + final String cluster = "test"; + admin.namespaces().createNamespace(namespace, Sets.newHashSet(cluster)); + Producer producer = pulsarClient.newProducer().topic(topicName).create(); + PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); + DispatchRateLimiter dispatchRateLimiter = new DispatchRateLimiter(topic, DispatchRateLimiter.Type.TOPIC); + + Policies policies = new Policies(); + DispatchRate clusterDispatchRate = new DispatchRate(100, 512, 1); + DispatchRate topicDispatchRate = new DispatchRate(200, 1024, 1); + + // (1) If both clusterDispatchRate and topicDispatchRate are empty, dispatch throttling is disabled + dispatchRateLimiter.onPoliciesUpdate(policies); + Assert.assertEquals(dispatchRateLimiter.getDispatchRateOnMsg(), -1); + Assert.assertEquals(dispatchRateLimiter.getDispatchRateOnByte(), -1); + + // (2) If topicDispatchRate is empty, clusterDispatchRate is effective + policies.clusterDispatchRate.put(cluster, clusterDispatchRate); + dispatchRateLimiter.onPoliciesUpdate(policies); + Assert.assertEquals(dispatchRateLimiter.getDispatchRateOnMsg(), 100); + Assert.assertEquals(dispatchRateLimiter.getDispatchRateOnByte(), 512); + + // (3) If topicDispatchRate is not empty, topicDispatchRate is effective + policies.topicDispatchRate.put(cluster, topicDispatchRate); + dispatchRateLimiter.onPoliciesUpdate(policies); + Assert.assertEquals(dispatchRateLimiter.getDispatchRateOnMsg(), 200); + Assert.assertEquals(dispatchRateLimiter.getDispatchRateOnByte(), 1024); + + producer.close(); + topic.close().get(); + } + protected void deactiveCursors(ManagedLedgerImpl ledger) throws Exception { Field statsUpdaterField = BrokerService.class.getDeclaredField("statsUpdater"); statsUpdaterField.setAccessible(true); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/Policies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/Policies.java index 70ad896b63dd0..175503adab466 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/Policies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/Policies.java @@ -39,6 +39,8 @@ public class Policies { public BundlesData bundles; @SuppressWarnings("checkstyle:MemberName") public Map backlog_quota_map = Maps.newHashMap(); + @Deprecated + public Map clusterDispatchRate = Maps.newHashMap(); public Map topicDispatchRate = Maps.newHashMap(); public Map subscriptionDispatchRate = Maps.newHashMap(); public Map replicatorDispatchRate = Maps.newHashMap(); @@ -103,7 +105,7 @@ public class Policies { @Override public int hashCode() { return Objects.hash(auth_policies, replication_clusters, - backlog_quota_map, publishMaxMessageRate, + backlog_quota_map, publishMaxMessageRate, clusterDispatchRate, topicDispatchRate, subscriptionDispatchRate, replicatorDispatchRate, clusterSubscribeRate, deduplicationEnabled, persistence, bundles, latency_stats_sample_rate, @@ -128,6 +130,7 @@ public boolean equals(Object obj) { return Objects.equals(auth_policies, other.auth_policies) && Objects.equals(replication_clusters, other.replication_clusters) && Objects.equals(backlog_quota_map, other.backlog_quota_map) + && Objects.equals(clusterDispatchRate, other.clusterDispatchRate) && Objects.equals(topicDispatchRate, other.topicDispatchRate) && Objects.equals(subscriptionDispatchRate, other.subscriptionDispatchRate) && Objects.equals(replicatorDispatchRate, other.replicatorDispatchRate) @@ -182,6 +185,7 @@ public String toString() { .add("replication_clusters", replication_clusters).add("bundles", bundles) .add("backlog_quota_map", backlog_quota_map).add("persistence", persistence) .add("deduplicationEnabled", deduplicationEnabled) + .add("clusterDispatchRate", clusterDispatchRate) .add("topicDispatchRate", topicDispatchRate) .add("subscriptionDispatchRate", subscriptionDispatchRate) .add("replicatorDispatchRate", replicatorDispatchRate)