From 4cfef6410dce639d3b11f8f127ea8eec74e20733 Mon Sep 17 00:00:00 2001 From: JiangHaiting Date: Sun, 16 Jan 2022 12:04:31 +0800 Subject: [PATCH 1/2] optmize deduplicationSnapshotIntervalSeconds --- .../pulsar/broker/service/AbstractTopic.java | 6 ++++++ .../persistent/MessageDeduplication.java | 21 +------------------ .../persistent/TopicDuplicationTest.java | 4 ++-- .../policies/data/HierarchyTopicPolicies.java | 2 ++ 4 files changed, 11 insertions(+), 22 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 95f661abed84e..2c4ec834e3114 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -163,6 +163,8 @@ protected void updateTopicPolicy(TopicPolicies data) { topicPolicies.getMaxConsumersPerSubscription().updateTopicValue(data.getMaxConsumersPerSubscription()); topicPolicies.getInactiveTopicPolicies().updateTopicValue(data.getInactiveTopicPolicies()); topicPolicies.getDeduplicationEnabled().updateTopicValue(data.getDeduplicationEnabled()); + topicPolicies.getDeduplicationSnapshotIntervalSeconds().updateTopicValue( + data.getDeduplicationSnapshotIntervalSeconds()); topicPolicies.getSubscriptionTypesEnabled().updateTopicValue( CollectionUtils.isEmpty(data.getSubscriptionTypesEnabled()) ? null : EnumSet.copyOf(data.getSubscriptionTypesEnabled())); @@ -194,6 +196,8 @@ protected void updateTopicPolicyByNamespacePolicy(Policies namespacePolicies) { .updateNamespaceValue(namespacePolicies.max_consumers_per_subscription); topicPolicies.getInactiveTopicPolicies().updateNamespaceValue(namespacePolicies.inactive_topic_policies); topicPolicies.getDeduplicationEnabled().updateNamespaceValue(namespacePolicies.deduplicationEnabled); + topicPolicies.getDeduplicationSnapshotIntervalSeconds().updateNamespaceValue( + namespacePolicies.deduplicationSnapshotIntervalSeconds); topicPolicies.getDelayedDeliveryEnabled().updateNamespaceValue( Optional.ofNullable(namespacePolicies.delayed_delivery_policies) .map(DelayedDeliveryPolicies::isActive).orElse(null)); @@ -220,6 +224,8 @@ private void updateTopicPolicyByBrokerConfig() { topicPolicies.getMaxConsumerPerTopic().updateBrokerValue(config.getMaxConsumersPerTopic()); topicPolicies.getMaxConsumersPerSubscription().updateBrokerValue(config.getMaxConsumersPerSubscription()); topicPolicies.getDeduplicationEnabled().updateBrokerValue(config.isBrokerDeduplicationEnabled()); + topicPolicies.getDeduplicationSnapshotIntervalSeconds().updateBrokerValue( + config.getBrokerDeduplicationSnapshotIntervalSeconds()); //init backlogQuota topicPolicies.getBackLogQuotaMap() .get(BacklogQuota.BacklogQuotaType.destination_storage) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index af9907d5602c0..11aaef324608e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -472,26 +472,7 @@ public long getLastPublishedSequenceId(String producerName) { } public void takeSnapshot() { - // try to get topic-level policies - Integer interval = topic.getTopicPolicies() - .map(TopicPolicies::getDeduplicationSnapshotIntervalSeconds) - .orElse(null); - try { - //if topic-level policies not exists, try to get namespace-level policies - if (interval == null) { - final Optional policies = pulsar.getPulsarResources().getNamespaceResources() - .getPolicies(TopicName.get(topic.getName()).getNamespaceObject()); - if (policies.isPresent()) { - interval = policies.get().deduplicationSnapshotIntervalSeconds; - } - } - } catch (Exception e) { - log.error("Failed to get namespace policies", e); - } - //There is no other level of policies, use the broker-level by default - if (interval == null) { - interval = pulsar.getConfiguration().getBrokerDeduplicationSnapshotIntervalSeconds(); - } + Integer interval = topic.getHierarchyTopicPolicies().getDeduplicationSnapshotIntervalSeconds().get(); long currentTimeStamp = System.currentTimeMillis(); if (interval == null || interval <= 0 || currentTimeStamp - lastSnapshotTimestamp < TimeUnit.SECONDS.toMillis(interval)) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/TopicDuplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/TopicDuplicationTest.java index 7b63ea5a391dc..c7cf44cb07acc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/TopicDuplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/TopicDuplicationTest.java @@ -240,7 +240,7 @@ public void testTopicPolicyTakeSnapshot() throws Exception { Producer producer = pulsarClient .newProducer(Schema.STRING).topic(topicName).enableBatching(false).producerName(producerName).create(); waitCacheInit(topicName); - admin.topics().setDeduplicationSnapshotInterval(topicName, 3); + admin.topicPolicies().setDeduplicationSnapshotInterval(topicName, 3); admin.namespaces().setDeduplicationSnapshotInterval(myNamespace, 5); int msgNum = 10; @@ -264,7 +264,7 @@ public void testTopicPolicyTakeSnapshot() throws Exception { assertEquals(position, markDeletedPosition); //remove topic-level policies, namespace-level should be used, interval becomes 5 seconds - admin.topics().removeDeduplicationSnapshotInterval(topicName); + admin.topicPolicies().removeDeduplicationSnapshotInterval(topicName); producer.newMessage().value("msg").send(); //zk update time + 5 second interval time Awaitility.await() diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/HierarchyTopicPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/HierarchyTopicPolicies.java index 820320b34ab4d..a61e90531bc30 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/HierarchyTopicPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/HierarchyTopicPolicies.java @@ -34,6 +34,7 @@ public class HierarchyTopicPolicies { final PolicyHierarchyValue> replicationClusters; final PolicyHierarchyValue deduplicationEnabled; + final PolicyHierarchyValue deduplicationSnapshotIntervalSeconds; final PolicyHierarchyValue inactiveTopicPolicies; final PolicyHierarchyValue> subscriptionTypesEnabled; final PolicyHierarchyValue maxSubscriptionsPerTopic; @@ -51,6 +52,7 @@ public class HierarchyTopicPolicies { public HierarchyTopicPolicies() { replicationClusters = new PolicyHierarchyValue<>(); deduplicationEnabled = new PolicyHierarchyValue<>(); + deduplicationSnapshotIntervalSeconds = new PolicyHierarchyValue<>(); inactiveTopicPolicies = new PolicyHierarchyValue<>(); subscriptionTypesEnabled = new PolicyHierarchyValue<>(); maxSubscriptionsPerTopic = new PolicyHierarchyValue<>(); From d12494ecb7275c79ecb7ab10ab8ea1b864e3bb78 Mon Sep 17 00:00:00 2001 From: JiangHaiting Date: Tue, 18 Jan 2022 09:26:17 +0800 Subject: [PATCH 2/2] remove unused import --- .../broker/service/persistent/MessageDeduplication.java | 4 ---- 1 file changed, 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 11aaef324608e..04c8ca71855da 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -24,7 +24,6 @@ import java.util.Iterator; import java.util.List; import java.util.Map; -import java.util.Optional; import java.util.TreeMap; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; @@ -40,9 +39,6 @@ import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.service.Topic.PublishContext; import org.apache.pulsar.common.api.proto.MessageMetadata; -import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.policies.data.Policies; -import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; import org.slf4j.Logger;