From 14d1dc01daf98f73191bf70f5681e559a25edfd4 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 16:57:43 +0800 Subject: [PATCH 1/9] [fix][broker] Fix topic policies deadlock --- .../SystemTopicBasedTopicPoliciesService.java | 57 +++++++++---------- 1 file changed, 28 insertions(+), 29 deletions(-) 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 cc3938491e637..b4fdac7d9b338 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 @@ -34,7 +34,11 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import javax.annotation.Nonnull; + +import it.unimi.dsi.fastutil.Pair; +import lombok.Data; import org.apache.commons.lang3.concurrent.ConcurrentInitializer; import org.apache.commons.lang3.concurrent.LazyInitializer; import org.apache.pulsar.broker.PulsarServerException; @@ -56,6 +60,7 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.SystemTopicNames; 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.util.FutureUtil; import org.slf4j.Logger; @@ -267,37 +272,31 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam return CompletableFuture.completedFuture(Optional.empty()); } final CompletableFuture preparedFuture = prepareInitPoliciesCacheAsync(topicName.getNamespaceObject()); - final var resultFuture = new CompletableFuture>(); - preparedFuture.thenAccept(inserted -> policyCacheInitMap.compute(namespace, (___, existingFuture) -> { - if (!inserted || existingFuture != null) { - final var partitionedTopicName = TopicName.get(topicName.getPartitionedTopicName()); - final var policies = Optional.ofNullable(switch (type) { - case DEFAULT -> Optional.ofNullable(policiesCache.get(partitionedTopicName)) - .orElseGet(() -> globalPoliciesCache.get(partitionedTopicName)); - case GLOBAL_ONLY -> globalPoliciesCache.get(partitionedTopicName); - case LOCAL_ONLY -> policiesCache.get(partitionedTopicName); - }); - resultFuture.complete(policies); - } else { - CompletableFuture.runAsync(() -> { - log.info("The future of {} has been removed from cache, retry getTopicPolicies again", namespace); - // Call it in another thread to avoid recursive update because getTopicPoliciesAsync() could call - // policyCacheInitMap.computeIfAbsent() - getTopicPoliciesAsync(topicName, type).whenComplete((result, e) -> { - if (e == null) { - resultFuture.complete(result); - } else { - resultFuture.completeExceptionally(e); - } - }); - }); + return preparedFuture.thenCompose(inserted -> { + @Data + class PoliciesFutureHolder { + CompletableFuture> future; } - return existingFuture; - })).exceptionally(e -> { - resultFuture.completeExceptionally(e); - return null; + final var policiesFutureHolder = new PoliciesFutureHolder(); + // notice: avoid using any callback with lock scope + policyCacheInitMap.compute(namespace, (___, existingFuture) -> { + if (!inserted || existingFuture != null) { + final var partitionedTopicName = TopicName.get(topicName.getPartitionedTopicName()); + final var policies = Optional.ofNullable(switch (type) { + case DEFAULT -> Optional.ofNullable(policiesCache.get(partitionedTopicName)) + .orElseGet(() -> globalPoliciesCache.get(partitionedTopicName)); + case GLOBAL_ONLY -> globalPoliciesCache.get(partitionedTopicName); + case LOCAL_ONLY -> policiesCache.get(partitionedTopicName); + }); + policiesFutureHolder.setFuture(CompletableFuture.completedFuture(policies)); + } else { + log.info("The future of {} has been removed from cache, retry getTopicPolicies again", namespace); + policiesFutureHolder.setFuture(getTopicPoliciesAsync(topicName, type)); + } + return existingFuture; + }); + return policiesFutureHolder.getFuture(); }); - return resultFuture; } public void addOwnedNamespaceBundleAsync(NamespaceBundle namespaceBundle) { From 7c7de9136963305d27caa7a393740ef2ff7f8200 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 17:04:30 +0800 Subject: [PATCH 2/9] add comments --- .../broker/service/SystemTopicBasedTopicPoliciesService.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) 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 b4fdac7d9b338..c8e578c152fc2 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 @@ -272,13 +272,14 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam return CompletableFuture.completedFuture(Optional.empty()); } final CompletableFuture preparedFuture = prepareInitPoliciesCacheAsync(topicName.getNamespaceObject()); - return preparedFuture.thenCompose(inserted -> { + // switch thread to avoid potential metadata thread cost and recursive deadlock + return preparedFuture.thenComposeAsync(inserted -> { @Data class PoliciesFutureHolder { CompletableFuture> future; } final var policiesFutureHolder = new PoliciesFutureHolder(); - // notice: avoid using any callback with lock scope + // NOTICE: avoid using any callback with lock scope to avoid deadlock policyCacheInitMap.compute(namespace, (___, existingFuture) -> { if (!inserted || existingFuture != null) { final var partitionedTopicName = TopicName.get(topicName.getPartitionedTopicName()); From a7f7de2b01802ac425a69e4f9113ea0838462b2f Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 17:08:48 +0800 Subject: [PATCH 3/9] fix checkstyle --- .../broker/service/SystemTopicBasedTopicPoliciesService.java | 3 --- 1 file changed, 3 deletions(-) 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 c8e578c152fc2..6dfad400d9bfe 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 @@ -34,10 +34,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.atomic.AtomicReference; import javax.annotation.Nonnull; - -import it.unimi.dsi.fastutil.Pair; import lombok.Data; import org.apache.commons.lang3.concurrent.ConcurrentInitializer; import org.apache.commons.lang3.concurrent.LazyInitializer; From 42cfb7f881776c39842d6da6a1a5de29ec1591f8 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 17:09:00 +0800 Subject: [PATCH 4/9] remove useless import --- .../broker/service/SystemTopicBasedTopicPoliciesService.java | 1 - 1 file changed, 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 6dfad400d9bfe..d64feac0f47ec 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 @@ -57,7 +57,6 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.SystemTopicNames; 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.util.FutureUtil; import org.slf4j.Logger; From 356c8767b71239504aa77e7727b934a734997801 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 17:14:41 +0800 Subject: [PATCH 5/9] add protection --- .../service/SystemTopicBasedTopicPoliciesService.java | 6 ++++++ 1 file changed, 6 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 d64feac0f47ec..02b0494911d0e 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 @@ -273,6 +273,12 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam @Data class PoliciesFutureHolder { CompletableFuture> future; + + public CompletableFuture> getFuture() { + return Objects.requireNonNullElseGet(future, + () -> CompletableFuture.failedFuture( + new IllegalStateException("BUG! unexpected topic policy init."))); + } } final var policiesFutureHolder = new PoliciesFutureHolder(); // NOTICE: avoid using any callback with lock scope to avoid deadlock From 52ab2d1df9c0b75e259bd2e58744d766f34ef204 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 21:49:52 +0800 Subject: [PATCH 6/9] apply comment --- .../SystemTopicBasedTopicPoliciesService.java | 23 +++++++------------ 1 file changed, 8 insertions(+), 15 deletions(-) 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 02b0494911d0e..60607c350461e 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 @@ -35,9 +35,10 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import javax.annotation.Nonnull; -import lombok.Data; import org.apache.commons.lang3.concurrent.ConcurrentInitializer; import org.apache.commons.lang3.concurrent.LazyInitializer; +import org.apache.commons.lang3.mutable.Mutable; +import org.apache.commons.lang3.mutable.MutableObject; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.namespace.NamespaceBundleOwnershipListener; @@ -270,17 +271,9 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam final CompletableFuture preparedFuture = prepareInitPoliciesCacheAsync(topicName.getNamespaceObject()); // switch thread to avoid potential metadata thread cost and recursive deadlock return preparedFuture.thenComposeAsync(inserted -> { - @Data - class PoliciesFutureHolder { - CompletableFuture> future; - - public CompletableFuture> getFuture() { - return Objects.requireNonNullElseGet(future, - () -> CompletableFuture.failedFuture( - new IllegalStateException("BUG! unexpected topic policy init."))); - } - } - final var policiesFutureHolder = new PoliciesFutureHolder(); + final Mutable>> policiesFutureHolder = + new MutableObject<>(CompletableFuture + .failedFuture(new IllegalStateException("BUG! unexpected topic policy init."))); // NOTICE: avoid using any callback with lock scope to avoid deadlock policyCacheInitMap.compute(namespace, (___, existingFuture) -> { if (!inserted || existingFuture != null) { @@ -291,14 +284,14 @@ public CompletableFuture> getFuture() { case GLOBAL_ONLY -> globalPoliciesCache.get(partitionedTopicName); case LOCAL_ONLY -> policiesCache.get(partitionedTopicName); }); - policiesFutureHolder.setFuture(CompletableFuture.completedFuture(policies)); + policiesFutureHolder.setValue(CompletableFuture.completedFuture(policies)); } else { log.info("The future of {} has been removed from cache, retry getTopicPolicies again", namespace); - policiesFutureHolder.setFuture(getTopicPoliciesAsync(topicName, type)); + policiesFutureHolder.setValue(getTopicPoliciesAsync(topicName, type)); } return existingFuture; }); - return policiesFutureHolder.getFuture(); + return policiesFutureHolder.getValue(); }); } From 3f2a838418fadfc410b4509670a9c29893d65caa Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Thu, 26 Dec 2024 21:51:34 +0800 Subject: [PATCH 7/9] avoid allocat useless object --- .../service/SystemTopicBasedTopicPoliciesService.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) 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 60607c350461e..7992f7fc735c1 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 @@ -271,9 +271,7 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam final CompletableFuture preparedFuture = prepareInitPoliciesCacheAsync(topicName.getNamespaceObject()); // switch thread to avoid potential metadata thread cost and recursive deadlock return preparedFuture.thenComposeAsync(inserted -> { - final Mutable>> policiesFutureHolder = - new MutableObject<>(CompletableFuture - .failedFuture(new IllegalStateException("BUG! unexpected topic policy init."))); + final Mutable>> policiesFutureHolder = new MutableObject<>(); // NOTICE: avoid using any callback with lock scope to avoid deadlock policyCacheInitMap.compute(namespace, (___, existingFuture) -> { if (!inserted || existingFuture != null) { @@ -291,7 +289,8 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam } return existingFuture; }); - return policiesFutureHolder.getValue(); + return Objects.requireNonNullElseGet(policiesFutureHolder.getValue(), ()-> CompletableFuture + .failedFuture(new IllegalStateException("BUG! unexpected topic policy init."))); }); } From bacc578da49b9832446adef84a22434345e3ca88 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Fri, 27 Dec 2024 00:21:12 +0800 Subject: [PATCH 8/9] fix recursive update --- .../SystemTopicBasedTopicPoliciesService.java | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) 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 7992f7fc735c1..5d3b09d9ced78 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 @@ -39,6 +39,7 @@ import org.apache.commons.lang3.concurrent.LazyInitializer; import org.apache.commons.lang3.mutable.Mutable; import org.apache.commons.lang3.mutable.MutableObject; +import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.namespace.NamespaceBundleOwnershipListener; @@ -271,7 +272,8 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam final CompletableFuture preparedFuture = prepareInitPoliciesCacheAsync(topicName.getNamespaceObject()); // switch thread to avoid potential metadata thread cost and recursive deadlock return preparedFuture.thenComposeAsync(inserted -> { - final Mutable>> policiesFutureHolder = new MutableObject<>(); + // initialized : policies + final Mutable>> policiesFutureHolder = new MutableObject<>(); // NOTICE: avoid using any callback with lock scope to avoid deadlock policyCacheInitMap.compute(namespace, (___, existingFuture) -> { if (!inserted || existingFuture != null) { @@ -282,15 +284,18 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam case GLOBAL_ONLY -> globalPoliciesCache.get(partitionedTopicName); case LOCAL_ONLY -> policiesCache.get(partitionedTopicName); }); - policiesFutureHolder.setValue(CompletableFuture.completedFuture(policies)); + policiesFutureHolder.setValue(Pair.of(true, policies)); } else { log.info("The future of {} has been removed from cache, retry getTopicPolicies again", namespace); - policiesFutureHolder.setValue(getTopicPoliciesAsync(topicName, type)); + policiesFutureHolder.setValue(Pair.of(false, null)); } return existingFuture; }); - return Objects.requireNonNullElseGet(policiesFutureHolder.getValue(), ()-> CompletableFuture - .failedFuture(new IllegalStateException("BUG! unexpected topic policy init."))); + final var p = policiesFutureHolder.getValue(); + if (!p.getLeft()) { + return getTopicPoliciesAsync(topicName, type); + } + return CompletableFuture.completedFuture(p.getRight()); }); } From 646772a6c26c53dd64e5801b9250af6610f6d9f2 Mon Sep 17 00:00:00 2001 From: Qiang Zhao Date: Fri, 27 Dec 2024 00:21:46 +0800 Subject: [PATCH 9/9] move the log --- .../broker/service/SystemTopicBasedTopicPoliciesService.java | 2 +- 1 file changed, 1 insertion(+), 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 5d3b09d9ced78..5488d5563f607 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 @@ -286,13 +286,13 @@ public CompletableFuture> getTopicPoliciesAsync(TopicNam }); policiesFutureHolder.setValue(Pair.of(true, policies)); } else { - log.info("The future of {} has been removed from cache, retry getTopicPolicies again", namespace); policiesFutureHolder.setValue(Pair.of(false, null)); } return existingFuture; }); final var p = policiesFutureHolder.getValue(); if (!p.getLeft()) { + log.info("The future of {} has been removed from cache, retry getTopicPolicies again", namespace); return getTopicPoliciesAsync(topicName, type); } return CompletableFuture.completedFuture(p.getRight());