From 314f70ccaab8869b5d34ee9830219521338edc6d Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Thu, 24 Sep 2026 23:46:39 +0800 Subject: [PATCH 1/4] [improve][broker] Add placement-aware BookKeeper client context Resolve topic placement configuration in one shared path and pair it with the cached BookKeeper client and matching ledger metadata. Assisted-by: Codex --- .../broker/ManagedLedgerClientFactory.java | 10 ++ .../apache/pulsar/broker/PulsarService.java | 38 ++++++ .../pulsar/broker/service/BrokerService.java | 45 ++----- .../storage/BookKeeperClientContext.java | 83 ++++++++++++ ...okKeeperPlacementPolicyConfigResolver.java | 81 ++++++++++++ .../BookkeeperManagedLedgerStorageClass.java | 24 ++++ .../ManagedLedgerClientFactoryTest.java | 37 ++++++ .../storage/BookKeeperClientContextTest.java | 106 ++++++++++++++++ ...eperPlacementPolicyConfigResolverTest.java | 118 ++++++++++++++++++ 9 files changed, 508 insertions(+), 34 deletions(-) create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperClientContext.java create mode 100644 pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolver.java create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperClientContextTest.java create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolverTest.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index e2915ceda3158..ee63af2eefb9c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -202,6 +202,16 @@ public StatsProvider getStatsProvider() { public BookKeeper getBookKeeperClient() { return defaultBkClient; } + + @Override + public CompletableFuture getBookKeeperClient( + EnsemblePlacementPolicyConfig ensemblePlacementPolicyConfig) { + try { + return bkFactory.get(ensemblePlacementPolicyConfig); + } catch (RuntimeException e) { + return CompletableFuture.failedFuture(e); + } + } }; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 83c05d02e02b3..96cd15ec51cdb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -129,6 +129,8 @@ import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsServlet; import org.apache.pulsar.broker.stats.prometheus.PrometheusRawMetricsProvider; import org.apache.pulsar.broker.stats.prometheus.PulsarPrometheusMetricsServlet; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.broker.storage.BookKeeperPlacementPolicyConfigResolver; import org.apache.pulsar.broker.storage.BookkeeperManagedLedgerStorageClass; import org.apache.pulsar.broker.storage.ManagedLedgerStorage; import org.apache.pulsar.broker.storage.ManagedLedgerStorageClass; @@ -171,6 +173,8 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterDataImpl; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig.ParseEnsemblePlacementPolicyConfigException; import org.apache.pulsar.common.policies.data.InactiveTopicDeleteMode; import org.apache.pulsar.common.policies.data.OffloadPoliciesImpl; import org.apache.pulsar.common.protocol.schema.SchemaStorage; @@ -1659,6 +1663,40 @@ public BookKeeper getBookKeeperClient() { } } + /** + * Resolve the placement policy for a topic and borrow the corresponding shared BookKeeper client. + * + * @param topicName the topic that owns the ledger + * @return a future that completes with the BookKeeper client and matching placement metadata + */ + public CompletableFuture getBookKeeperClientContext(TopicName topicName) { + return CompletableFuture.completedFuture(topicName) + .thenCompose(name -> { + Objects.requireNonNull(name, "topicName"); + return getPulsarResources().getLocalPolicies() + .getLocalPoliciesAsync(name.getNamespaceObject()); + }).thenCompose(localPolicies -> { + EnsemblePlacementPolicyConfig placementPolicyConfig = + BookKeeperPlacementPolicyConfigResolver.resolve(config, topicName, localPolicies) + .orElse(null); + ManagedLedgerStorageClass defaultStorageClass = + getManagedLedgerStorage().getDefaultStorageClass(); + if (!(defaultStorageClass instanceof BookkeeperManagedLedgerStorageClass bkStorageClass)) { + return CompletableFuture.failedFuture( + new UnsupportedOperationException("BookKeeper client is not available")); + } + return bkStorageClass.getBookKeeperClient(placementPolicyConfig) + .thenCompose(bookKeeper -> { + try { + return CompletableFuture.completedFuture( + BookKeeperClientContext.create(bookKeeper, placementPolicyConfig)); + } catch (ParseEnsemblePlacementPolicyConfigException e) { + return CompletableFuture.failedFuture(e); + } + }); + }); + } + public ManagedLedgerFactory getDefaultManagedLedgerFactory() { return getManagedLedgerStorage().getDefaultStorageClass().getManagedLedgerFactory(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 4e55c1a4295a0..e7c996e40fa83 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -86,6 +86,7 @@ import lombok.CustomLog; import lombok.Getter; import lombok.Setter; +import org.apache.bookkeeper.client.EnsemblePlacementPolicy; import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteLedgerCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenLedgerCallback; @@ -102,7 +103,6 @@ import org.apache.commons.lang3.mutable.MutableBoolean; import org.apache.commons.lang3.tuple.ImmutablePair; import org.apache.commons.lang3.tuple.Pair; -import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; @@ -146,6 +146,7 @@ import org.apache.pulsar.broker.stats.ClusterReplicationMetrics; import org.apache.pulsar.broker.stats.prometheus.metrics.ObserverGauge; import org.apache.pulsar.broker.stats.prometheus.metrics.Summary; +import org.apache.pulsar.broker.storage.BookKeeperPlacementPolicyConfigResolver; import org.apache.pulsar.broker.storage.ManagedLedgerStorage; import org.apache.pulsar.broker.storage.ManagedLedgerStorageClass; import org.apache.pulsar.broker.topiclistlimit.TopicListMemoryLimiter; @@ -180,6 +181,7 @@ import org.apache.pulsar.common.policies.data.AutoTopicCreationOverride; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import org.apache.pulsar.common.policies.data.LocalPolicies; import org.apache.pulsar.common.policies.data.OffloadPoliciesImpl; import org.apache.pulsar.common.policies.data.PersistencePolicies; @@ -2481,39 +2483,14 @@ private CompletableFuture getManagedLedgerConfig( managedLedgerConfig.setAckQuorumSize(persistencePolicies.getBookkeeperAckQuorum()); managedLedgerConfig.setStorageClassName(persistencePolicies.getManagedLedgerStorageClassName()); - if (serviceConfig.isStrictBookieAffinityEnabled()) { - managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyClassName( - IsolatedBookieEnsemblePlacementPolicy.class); - if (localPolicies.isPresent() && localPolicies.get().bookieAffinityGroup != null) { - Map properties = new HashMap<>(); - properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, - localPolicies.get().bookieAffinityGroup.getBookkeeperAffinityGroupPrimary()); - properties.put(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, - localPolicies.get().bookieAffinityGroup.getBookkeeperAffinityGroupSecondary()); - managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(properties); - } else if (isSystemTopic(topicName)) { - Map properties = new HashMap<>(); - properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "*"); - properties.put(IsolatedBookieEnsemblePlacementPolicy - .SECONDARY_ISOLATION_BOOKIE_GROUPS, "*"); - managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(properties); - } else { - Map properties = new HashMap<>(); - properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, ""); - properties.put(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, ""); - managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(properties); - } - } else { - if (localPolicies.isPresent() && localPolicies.get().bookieAffinityGroup != null) { - managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyClassName( - IsolatedBookieEnsemblePlacementPolicy.class); - Map properties = new HashMap<>(); - properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, - localPolicies.get().bookieAffinityGroup.getBookkeeperAffinityGroupPrimary()); - properties.put(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, - localPolicies.get().bookieAffinityGroup.getBookkeeperAffinityGroupSecondary()); - managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(properties); - } + Optional placementPolicy = + BookKeeperPlacementPolicyConfigResolver.resolve(serviceConfig, topicName, localPolicies); + if (placementPolicy.isPresent()) { + EnsemblePlacementPolicyConfig policyConfig = placementPolicy.get(); + Class policyClass = + policyConfig.getPolicyClass().asSubclass(EnsemblePlacementPolicy.class); + managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyClassName(policyClass); + managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(policyConfig.getProperties()); } managedLedgerConfig.setThrottleMarkDelete(persistencePolicies.getManagedLedgerMaxMarkDeleteRate() >= 0 diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperClientContext.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperClientContext.java new file mode 100644 index 0000000000000..c2d9ecee033e7 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperClientContext.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.storage; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig.ParseEnsemblePlacementPolicyConfigException; + +/** + * A borrowed BookKeeper client together with the placement metadata that must be used for new ledgers. + */ +public final class BookKeeperClientContext { + private final BookKeeper bookKeeper; + private final byte[] encodedPlacementPolicyConfig; + + private BookKeeperClientContext(BookKeeper bookKeeper, byte[] encodedPlacementPolicyConfig) { + this.bookKeeper = Objects.requireNonNull(bookKeeper); + this.encodedPlacementPolicyConfig = encodedPlacementPolicyConfig; + } + + /** + * Create a context for a BookKeeper client and placement policy. + * + * @param bookKeeper the borrowed BookKeeper client + * @param placementPolicyConfig the policy used by the client, or {@code null} for the default client + * @return a client context + * @throws ParseEnsemblePlacementPolicyConfigException if the policy cannot be encoded + */ + public static BookKeeperClientContext create(BookKeeper bookKeeper, + EnsemblePlacementPolicyConfig placementPolicyConfig) + throws ParseEnsemblePlacementPolicyConfigException { + byte[] encodedConfig = placementPolicyConfig != null + && placementPolicyConfig.getPolicyClass() != null + && placementPolicyConfig.getProperties() != null + ? placementPolicyConfig.encode() : null; + return new BookKeeperClientContext(bookKeeper, encodedConfig); + } + + /** + * Return the borrowed BookKeeper client. The caller must not close it. + */ + public BookKeeper getBookKeeper() { + return bookKeeper; + } + + /** + * Add the placement policy metadata required for ledger recovery to the supplied metadata. + * + * @param metadata component-specific ledger metadata + * @return metadata containing the placement policy configuration when one is configured + */ + public Map withPlacementMetadata(Map metadata) { + Objects.requireNonNull(metadata); + if (encodedPlacementPolicyConfig == null) { + return metadata; + } + Map metadataWithPlacementPolicy = new HashMap<>(metadata); + metadataWithPlacementPolicy.put( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG, + encodedPlacementPolicyConfig.clone()); + return Collections.unmodifiableMap(metadataWithPlacementPolicy); + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolver.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolver.java new file mode 100644 index 0000000000000..8698a033b30f2 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolver.java @@ -0,0 +1,81 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.storage; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.Optional; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.namespace.NamespaceService; +import org.apache.pulsar.common.naming.SystemTopicNames; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.policies.data.LocalPolicies; + +/** + * Resolves the BookKeeper ensemble placement policy for a topic. + */ +public final class BookKeeperPlacementPolicyConfigResolver { + + /** + * Resolve the placement policy from broker configuration and namespace local policies. + * + * @param configuration broker configuration + * @param topicName topic whose auxiliary ledger will be created + * @param localPolicies namespace local policies + * @return the custom placement policy, or empty when the default BookKeeper client should be used + */ + public static Optional resolve(ServiceConfiguration configuration, + TopicName topicName, + Optional localPolicies) { + BookieAffinityGroupData affinityGroup = localPolicies + .map(policies -> policies.bookieAffinityGroup) + .orElse(null); + if (!configuration.isStrictBookieAffinityEnabled() && affinityGroup == null) { + return Optional.empty(); + } + + Map properties = new HashMap<>(); + if (affinityGroup != null) { + properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, + affinityGroup.getBookkeeperAffinityGroupPrimary()); + properties.put(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, + affinityGroup.getBookkeeperAffinityGroupSecondary()); + } else if (isSystemTopic(topicName)) { + properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "*"); + properties.put(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "*"); + } else { + properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, ""); + properties.put(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, ""); + } + return Optional.of(new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, Collections.unmodifiableMap(properties))); + } + + private static boolean isSystemTopic(TopicName topicName) { + return NamespaceService.isSystemServiceNamespace(topicName.getNamespace()) + || SystemTopicNames.isSystemTopic(topicName); + } + + private BookKeeperPlacementPolicyConfigResolver() { + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookkeeperManagedLedgerStorageClass.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookkeeperManagedLedgerStorageClass.java index 1f05cde72a5b5..c29e761e34a3e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookkeeperManagedLedgerStorageClass.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/storage/BookkeeperManagedLedgerStorageClass.java @@ -18,8 +18,10 @@ */ package org.apache.pulsar.broker.storage; +import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; /** * ManagedLedgerStorageClass represents a configured instance of ManagedLedgerFactory for managed ledgers. @@ -33,6 +35,28 @@ public interface BookkeeperManagedLedgerStorageClass extends ManagedLedgerStorag */ BookKeeper getBookKeeperClient(); + /** + * Return the BookKeeper client for the specified ensemble placement policy. + * + *

The returned client is owned by the storage class and must not be closed by the caller. + * + * @param ensemblePlacementPolicyConfig the ensemble placement policy configuration + * @return a future that completes with the BookKeeper client + */ + default CompletableFuture getBookKeeperClient( + EnsemblePlacementPolicyConfig ensemblePlacementPolicyConfig) { + if (ensemblePlacementPolicyConfig == null + || ensemblePlacementPolicyConfig.getPolicyClass() == null) { + try { + return CompletableFuture.completedFuture(getBookKeeperClient()); + } catch (RuntimeException e) { + return CompletableFuture.failedFuture(e); + } + } + return CompletableFuture.failedFuture(new UnsupportedOperationException( + "Custom ensemble placement policies are not supported by this storage class")); + } + /** * Return the stats provider to expose the stats of the storage implementation. * diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/ManagedLedgerClientFactoryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/ManagedLedgerClientFactoryTest.java index 5b6c7d9cd0c23..1340142c0c925 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/ManagedLedgerClientFactoryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/ManagedLedgerClientFactoryTest.java @@ -20,18 +20,24 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.ArgumentMatchers.isNull; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import io.netty.channel.EventLoopGroup; import io.opentelemetry.api.OpenTelemetry; +import java.util.Map; import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.storage.BookkeeperManagedLedgerStorageClass; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; import org.mockito.ArgumentCaptor; import org.testng.annotations.DataProvider; @@ -72,4 +78,35 @@ public void testCacheExtensionSettingsAtStartup(Boolean extendRecentlyAccessed, .isEqualTo(maxExtensions); } } + + @Test + public void testPlacementPolicyClientIsSharedByEquivalentConfigurations() throws Exception { + ServiceConfiguration conf = new ServiceConfiguration(); + conf.setBookkeeperClientExposeStatsToPrometheus(false); + BookKeeper defaultClient = mock(BookKeeper.class); + BookKeeper policyClient = mock(BookKeeper.class); + BookKeeperClientFactory bookkeeperProvider = mock(BookKeeperClientFactory.class); + when(bookkeeperProvider.create(any(), any(), any(), any(), isNull(), any())) + .thenReturn(CompletableFuture.completedFuture(defaultClient)); + when(bookkeeperProvider.create(any(), any(), any(), any(), anyMap(), any())) + .thenReturn(CompletableFuture.completedFuture(policyClient)); + + try (ManagedLedgerClientFactory factory = spy(new ManagedLedgerClientFactory())) { + doReturn(mock(ManagedLedgerFactoryImpl.class)).when(factory) + .createManagedLedgerFactory(any(), any(), any(), any(), any()); + factory.initialize(conf, mock(MetadataStoreExtended.class), bookkeeperProvider, + mock(EventLoopGroup.class), OpenTelemetry.noop()); + BookkeeperManagedLedgerStorageClass storageClass = + (BookkeeperManagedLedgerStorageClass) factory.getDefaultStorageClass(); + EnsemblePlacementPolicyConfig firstConfig = new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, Map.of("group", "primary")); + EnsemblePlacementPolicyConfig equivalentConfig = new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, Map.of("group", "primary")); + + assertThat(storageClass.getBookKeeperClient(firstConfig).join()).isSameAs(policyClient); + assertThat(storageClass.getBookKeeperClient(equivalentConfig).join()).isSameAs(policyClient); + assertThat(factory.getBkEnsemblePolicyToBookKeeperMap()).hasSize(1); + verify(bookkeeperProvider, times(2)).create(any(), any(), any(), any(), any(), any()); + } + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperClientContextTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperClientContextTest.java new file mode 100644 index 0000000000000..49863730b3c8c --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperClientContextTest.java @@ -0,0 +1,106 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.storage; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; +import java.nio.charset.StandardCharsets; +import java.util.Map; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.mledger.ManagedLedgerFactory; +import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.testng.annotations.Test; + +@Test(groups = "broker") +public class BookKeeperClientContextTest { + + @Test + public void testDefaultClientDoesNotAddPlacementMetadata() throws Exception { + BookKeeper bookKeeper = mock(BookKeeper.class); + BookKeeperClientContext context = BookKeeperClientContext.create(bookKeeper, null); + Map metadata = Map.of("component", "schema".getBytes(StandardCharsets.UTF_8)); + + assertThat(context.getBookKeeper()).isSameAs(bookKeeper); + assertThat(context.withPlacementMetadata(metadata)).isSameAs(metadata); + assertThat(metadata).doesNotContainKey( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + } + + @Test + public void testPlacementMetadataMatchesClientPolicy() throws Exception { + EnsemblePlacementPolicyConfig policy = new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, + Map.of(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary", + IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "secondary")); + BookKeeperClientContext context = BookKeeperClientContext.create(mock(BookKeeper.class), policy); + byte[] stalePolicy = "stale".getBytes(StandardCharsets.UTF_8); + Map metadata = Map.of( + "component", "schema".getBytes(StandardCharsets.UTF_8), + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG, stalePolicy); + + Map firstResult = context.withPlacementMetadata(metadata); + byte[] firstEncodedPolicy = firstResult.get( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + + assertThat(EnsemblePlacementPolicyConfig.decode(firstEncodedPolicy)).isEqualTo(policy); + assertThat(metadata.get(EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG)) + .isSameAs(stalePolicy); + + firstEncodedPolicy[0] = (byte) (firstEncodedPolicy[0] + 1); + byte[] secondEncodedPolicy = context.withPlacementMetadata(metadata).get( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + assertThat(secondEncodedPolicy).isNotSameAs(firstEncodedPolicy); + assertThat(EnsemblePlacementPolicyConfig.decode(secondEncodedPolicy)).isEqualTo(policy); + } + + @Test + public void testStorageClassDoesNotSilentlyIgnoreCustomPolicy() { + BookKeeper defaultClient = mock(BookKeeper.class); + BookkeeperManagedLedgerStorageClass storageClass = new BookkeeperManagedLedgerStorageClass() { + @Override + public BookKeeper getBookKeeperClient() { + return defaultClient; + } + + @Override + public StatsProvider getStatsProvider() { + return null; + } + + @Override + public String getName() { + return "test"; + } + + @Override + public ManagedLedgerFactory getManagedLedgerFactory() { + return null; + } + }; + EnsemblePlacementPolicyConfig policy = new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, Map.of()); + + assertThat(storageClass.getBookKeeperClient(null).join()).isSameAs(defaultClient); + assertThatThrownBy(() -> storageClass.getBookKeeperClient(policy).join()) + .hasCauseInstanceOf(UnsupportedOperationException.class); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolverTest.java new file mode 100644 index 0000000000000..7b5d7e4d4f5a0 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/storage/BookKeeperPlacementPolicyConfigResolverTest.java @@ -0,0 +1,118 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.storage; + +import static org.assertj.core.api.Assertions.assertThat; +import java.util.Optional; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.policies.data.LocalPolicies; +import org.testng.annotations.Test; + +@Test(groups = "broker") +public class BookKeeperPlacementPolicyConfigResolverTest { + + @Test + public void testNoCustomPolicyWithoutStrictAffinity() { + ServiceConfiguration configuration = new ServiceConfiguration(); + + assertThat(resolve(configuration, "persistent://tenant/namespace/topic", Optional.empty())).isEmpty(); + assertThat(resolve(configuration, "persistent://tenant/namespace/__change_events", Optional.empty())) + .isEmpty(); + } + + @Test + public void testNamespaceAffinity() { + ServiceConfiguration configuration = new ServiceConfiguration(); + Optional localPolicies = affinityGroup("primary", null); + + EnsemblePlacementPolicyConfig policy = resolve( + configuration, "persistent://tenant/namespace/topic", localPolicies).orElseThrow(); + + assertThat(policy.getPolicyClass()).isEqualTo(IsolatedBookieEnsemblePlacementPolicy.class); + assertThat(policy.getProperties()) + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary") + .containsKey(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS); + assertThat(policy.getProperties() + .get(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS)).isNull(); + } + + @Test + public void testStrictAffinityForOrdinaryTopic() { + ServiceConfiguration configuration = new ServiceConfiguration(); + configuration.setStrictBookieAffinityEnabled(true); + + EnsemblePlacementPolicyConfig policy = resolve( + configuration, "persistent://tenant/namespace/topic", Optional.empty()).orElseThrow(); + + assertThat(policy.getProperties()) + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "") + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, ""); + } + + @Test + public void testStrictAffinityForSystemTopics() { + ServiceConfiguration configuration = new ServiceConfiguration(); + configuration.setStrictBookieAffinityEnabled(true); + + assertWildcardPolicy(resolve(configuration, + "persistent://tenant/namespace/__change_events", Optional.empty()).orElseThrow()); + assertWildcardPolicy(resolve(configuration, + "persistent://pulsar/system/ordinary-topic", Optional.empty()).orElseThrow()); + } + + @Test + public void testNamespaceAffinityOverridesStrictSystemTopicDefault() { + ServiceConfiguration configuration = new ServiceConfiguration(); + configuration.setStrictBookieAffinityEnabled(true); + + EnsemblePlacementPolicyConfig policy = resolve(configuration, + "persistent://tenant/namespace/__change_events", affinityGroup("primary", "secondary")) + .orElseThrow(); + + assertThat(policy.getProperties()) + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary") + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, + "secondary"); + } + + private static Optional resolve(ServiceConfiguration configuration, + String topic, + Optional localPolicies) { + return BookKeeperPlacementPolicyConfigResolver.resolve( + configuration, TopicName.get(topic), localPolicies); + } + + private static Optional affinityGroup(String primary, String secondary) { + BookieAffinityGroupData affinityGroup = BookieAffinityGroupData.builder() + .bookkeeperAffinityGroupPrimary(primary) + .bookkeeperAffinityGroupSecondary(secondary) + .build(); + return Optional.of(new LocalPolicies(null, affinityGroup, null)); + } + + private static void assertWildcardPolicy(EnsemblePlacementPolicyConfig policy) { + assertThat(policy.getProperties()) + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "*") + .containsEntry(IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "*"); + } +} From a7b2ccb439fd601c9902b3aab8e8ec117756eef7 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Fri, 25 Sep 2026 00:07:05 +0800 Subject: [PATCH 2/4] [fix][broker] Apply placement policy to schema and delayed bucket ledgers Resolve the owner topic before creating schema and delayed-delivery bucket ledgers, then use the shared placement-aware BookKeeper client and persist the matching policy metadata. Preserve arbitrary schema keys and bucket retry semantics. Assisted-by: Codex --- .../BookkeeperBucketSnapshotStorage.java | 53 +++++--- .../schema/BookkeeperSchemaStorage.java | 47 ++++++- .../BookkeeperBucketSnapshotStorageTest.java | 2 +- ...perBucketSnapshotStoragePlacementTest.java | 112 ++++++++++++++++ .../schema/BookkeeperSchemaStorageTest.java | 126 ++++++++++++++++++ 5 files changed, 316 insertions(+), 24 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java index 5c8b2d13dc00b..b1e4bc57aaa7e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java @@ -40,6 +40,7 @@ import org.apache.pulsar.broker.delayed.proto.SnapshotSegment; import org.apache.pulsar.broker.delayed.proto.SnapshotSegmentMetadata; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; import org.jspecify.annotations.NonNull; @@ -172,28 +173,40 @@ private List parseSnapshotSegmentEntries(Enumeration createLedger(String bucketKey, String topicName, String cursorName) { - CompletableFuture future = new CompletableFuture<>(); - Map metadata = LedgerMetadataUtils.buildMetadataForDelayedIndexBucket(bucketKey, - topicName, cursorName); - bookKeeper.newCreateLedgerOp() - .withEnsembleSize(config.getManagedLedgerDefaultEnsembleSize()) - .withWriteQuorumSize(config.getManagedLedgerDefaultWriteQuorum()) - .withAckQuorumSize(config.getManagedLedgerDefaultAckQuorum()) - .withDigestType(config.getManagedLedgerDigestType()) - .withPassword(LedgerPassword) - .withCustomMetadata(metadata) - .withLoggerContext(log.with().attr("topic", topicName).attr("cursor", cursorName).build()) - .execute() - .whenComplete((writeHandle, ex) -> { - if (ex != null) { - future.completeExceptionally(bkException("Create ledger", - BKException.getExceptionCode(ex), -1)); - } else { - future.complete((LedgerHandle) writeHandle); + @VisibleForTesting + CompletableFuture createLedger(String bucketKey, String topicName, String cursorName) { + return CompletableFuture.completedFuture(topicName) + .thenApply(TopicName::get) + .thenCompose(pulsar::getBookKeeperClientContext) + .thenCompose(bookKeeperClientContext -> { + CompletableFuture future = new CompletableFuture<>(); + Map metadata = bookKeeperClientContext.withPlacementMetadata( + LedgerMetadataUtils.buildMetadataForDelayedIndexBucket(bucketKey, topicName, cursorName)); + bookKeeperClientContext.getBookKeeper().newCreateLedgerOp() + .withEnsembleSize(config.getManagedLedgerDefaultEnsembleSize()) + .withWriteQuorumSize(config.getManagedLedgerDefaultWriteQuorum()) + .withAckQuorumSize(config.getManagedLedgerDefaultAckQuorum()) + .withDigestType(config.getManagedLedgerDigestType()) + .withPassword(LedgerPassword) + .withCustomMetadata(metadata) + .withLoggerContext(log.with().attr("topic", topicName).attr("cursor", cursorName).build()) + .execute() + .whenComplete((writeHandle, ex) -> { + if (ex != null) { + future.completeExceptionally(bkException("Create ledger", + BKException.getExceptionCode(ex), -1)); + } else { + future.complete((LedgerHandle) writeHandle); + } + }); + return future; + }).exceptionallyCompose(ex -> { + Throwable cause = FutureUtil.unwrapCompletionException(ex); + if (cause instanceof BucketSnapshotPersistenceException) { + return FutureUtil.failedFuture(cause); } + return FutureUtil.failedFuture(new BucketSnapshotPersistenceException(cause)); }); - return future; } private CompletableFuture getLedgerHandle(Long ledgerId) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java index 246da337f509a..2612dc8af585c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java @@ -49,11 +49,15 @@ import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException; import org.apache.pulsar.broker.service.schema.exceptions.SchemaException; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.common.naming.TopicDomain; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.SchemaMetadata; import org.apache.pulsar.common.protocol.schema.SchemaStorage; import org.apache.pulsar.common.protocol.schema.SchemaVersion; import org.apache.pulsar.common.protocol.schema.StoredSchema; import org.apache.pulsar.common.schema.LongSchemaVersion; +import org.apache.pulsar.common.util.Codec; import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.metadata.api.MetadataCache; import org.apache.pulsar.metadata.api.MetadataSerde; @@ -641,11 +645,48 @@ private CompletableFuture addEntry(LedgerHandle ledgerHandle, SchemaEntry } @NonNull - private CompletableFuture createLedger(String schemaId) { - Map metadata = LedgerMetadataUtils.buildMetadataForSchema(schemaId); + @VisibleForTesting + CompletableFuture createLedger(String schemaId) { + return CompletableFuture.completedFuture(schemaId) + .thenCompose(id -> { + Optional topicName = getOwnerTopicName(id); + if (topicName.isEmpty()) { + return createLedger(id, bookKeeper, LedgerMetadataUtils.buildMetadataForSchema(id)); + } + return pulsar.getBookKeeperClientContext(topicName.get()) + .thenCompose(clientContext -> createLedger(id, clientContext)); + }); + } + + private static Optional getOwnerTopicName(String schemaId) { + String[] parts = schemaId.split("/", -1); + if (parts.length != 3) { + return Optional.empty(); + } + try { + String localName = Codec.decode(parts[2]); + TopicDomain domain = localName.indexOf('/') >= 0 ? TopicDomain.topic : TopicDomain.persistent; + TopicName topicName = TopicName.get(domain.value(), parts[0], parts[1], localName); + return schemaId.equals(topicName.getSchemaName()) ? Optional.of(topicName) : Optional.empty(); + } catch (IllegalArgumentException e) { + // SchemaStorage historically accepts arbitrary keys, including legacy topic names. + // Non-canonical keys have no unambiguous owner namespace, so use the default schema client. + return Optional.empty(); + } + } + + private CompletableFuture createLedger(String schemaId, + BookKeeperClientContext clientContext) { + Map metadata = clientContext.withPlacementMetadata( + LedgerMetadataUtils.buildMetadataForSchema(schemaId)); + return createLedger(schemaId, clientContext.getBookKeeper(), metadata); + } + + private CompletableFuture createLedger(String schemaId, BookKeeper client, + Map metadata) { final CompletableFuture future = new CompletableFuture<>(); try { - bookKeeper.newCreateLedgerOp() + client.newCreateLedgerOp() .withEnsembleSize(config.getManagedLedgerDefaultEnsembleSize()) .withWriteQuorumSize(config.getManagedLedgerDefaultWriteQuorum()) .withAckQuorumSize(config.getManagedLedgerDefaultAckQuorum()) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BookkeeperBucketSnapshotStorageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BookkeeperBucketSnapshotStorageTest.java index 82396e66fa3d4..74066137a75a9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BookkeeperBucketSnapshotStorageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/BookkeeperBucketSnapshotStorageTest.java @@ -64,7 +64,7 @@ protected void cleanup() throws Exception { bucketSnapshotStorage.close(); } - private static final String TOPIC_NAME = "topicName"; + private static final String TOPIC_NAME = "persistent://public/default/bucket-snapshot-storage-test"; private static final String CURSOR_NAME = "sub"; @Test diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java new file mode 100644 index 0000000000000..bf8de4a62edc1 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java @@ -0,0 +1,112 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.delayed.bucket; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.RETURNS_SELF; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import java.nio.charset.StandardCharsets; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.LedgerHandle; +import org.apache.bookkeeper.client.api.CreateBuilder; +import org.apache.bookkeeper.mledger.impl.LedgerMetadataUtils; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.util.FutureUtil; +import org.mockito.ArgumentCaptor; +import org.testng.annotations.Test; + +@Test(groups = "broker") +public class BookkeeperBucketSnapshotStoragePlacementTest { + + private static final String TOPIC = "persistent://tenant/namespace/topic"; + + @Test + public void testCreateLedgerUsesTopicClientAndMatchingPlacementMetadata() throws Exception { + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getConfig()).thenReturn(new ServiceConfiguration()); + BookKeeper policyBookKeeper = mock(BookKeeper.class); + EnsemblePlacementPolicyConfig placementPolicy = new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, + Map.of(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary", + IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "secondary")); + BookKeeperClientContext context = BookKeeperClientContext.create(policyBookKeeper, placementPolicy); + when(pulsar.getBookKeeperClientContext(TopicName.get(TOPIC))) + .thenReturn(CompletableFuture.completedFuture(context)); + + CreateBuilder createBuilder = mock(CreateBuilder.class, RETURNS_SELF); + LedgerHandle ledgerHandle = mock(LedgerHandle.class); + when(policyBookKeeper.newCreateLedgerOp()).thenReturn(createBuilder); + doReturn(CompletableFuture.completedFuture(ledgerHandle)).when(createBuilder).execute(); + + BookkeeperBucketSnapshotStorage storage = new BookkeeperBucketSnapshotStorage(pulsar); + assertThat(storage.createLedger("bucket", TOPIC, "subscription").join()).isSameAs(ledgerHandle); + + verify(policyBookKeeper).newCreateLedgerOp(); + @SuppressWarnings("unchecked") + ArgumentCaptor> metadataCaptor = ArgumentCaptor.forClass(Map.class); + verify(createBuilder).withCustomMetadata(metadataCaptor.capture()); + Map metadata = metadataCaptor.getValue(); + assertThat(metadata.get(LedgerMetadataUtils.METADATA_PROPERTY_COMPONENT)) + .containsExactly(LedgerMetadataUtils.METADATA_PROPERTY_COMPONENT_DELAYED_INDEX_BUCKET); + assertThat(new String(metadata.get(LedgerMetadataUtils.METADATA_PROPERTY_DELAYED_INDEX_TOPIC), + StandardCharsets.UTF_8)).isEqualTo(TOPIC); + assertThat(EnsemblePlacementPolicyConfig.decode(metadata.get( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG))) + .isEqualTo(placementPolicy); + } + + @Test + public void testInvalidTopicCompletesFutureExceptionally() { + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getConfig()).thenReturn(new ServiceConfiguration()); + BookkeeperBucketSnapshotStorage storage = new BookkeeperBucketSnapshotStorage(pulsar); + + CompletableFuture future = storage.createLedger("bucket", null, "subscription"); + + assertThat(future).isCompletedExceptionally(); + assertThatThrownBy(future::join).hasCauseInstanceOf(BucketSnapshotPersistenceException.class); + } + + @Test + public void testPlacementLookupFailureIsRetriablePersistenceFailure() { + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getConfig()).thenReturn(new ServiceConfiguration()); + RuntimeException lookupFailure = new RuntimeException("Failed to read local policies"); + when(pulsar.getBookKeeperClientContext(TopicName.get(TOPIC))) + .thenReturn(FutureUtil.failedFuture(lookupFailure)); + BookkeeperBucketSnapshotStorage storage = new BookkeeperBucketSnapshotStorage(pulsar); + + CompletableFuture future = storage.createLedger("bucket", TOPIC, "subscription"); + + Throwable failure = future.handle((__, ex) -> FutureUtil.unwrapCompletionException(ex)).join(); + assertThat(failure).isInstanceOf(BucketSnapshotPersistenceException.class); + assertThat(failure.getCause()).isSameAs(lookupFailure); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java index 400e85403dd9f..f524d932f9bda 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java @@ -19,17 +19,39 @@ package org.apache.pulsar.broker.service.schema; import static org.apache.pulsar.broker.service.schema.BookkeeperSchemaStorage.bkException; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import java.nio.ByteBuffer; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.LedgerHandle; import org.apache.bookkeeper.client.api.BKException; +import org.apache.bookkeeper.client.api.CreateBuilder; +import org.apache.bookkeeper.client.api.WriteHandle; +import org.apache.bookkeeper.mledger.impl.LedgerMetadataUtils; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.BookKeeperClientFactory; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.schema.exceptions.SchemaException; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import org.apache.pulsar.common.schema.LongSchemaVersion; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import org.mockito.ArgumentCaptor; +import org.mockito.Mockito; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @@ -74,4 +96,108 @@ public void testVersionFromBytes() { assertEquals(new LongSchemaVersion(version), schemaStorage.versionFromBytes(versionBytesPre240)); assertEquals(new LongSchemaVersion(version), schemaStorage.versionFromBytes(versionBytesPost240)); } + + @DataProvider(name = "canonicalSchemaIds") + public static Object[][] canonicalSchemaIds() { + return new Object[][] { + {"tenant/namespace/topic", TopicName.get("persistent://tenant/namespace/topic")}, + {"tenant/namespace/a%3Ab", TopicName.get("persistent://tenant/namespace/a:b")}, + {"tenant/namespace/a%2Fb", TopicName.get("topic://tenant/namespace/a/b")} + }; + } + + @Test(dataProvider = "canonicalSchemaIds") + public void testCreateLedgerUsesTopicPlacementClientAndMetadata(String schemaId, TopicName ownerTopic) + throws Exception { + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getLocalMetadataStore()).thenReturn(mock(MetadataStoreExtended.class)); + when(pulsar.getConfiguration()).thenReturn(new ServiceConfiguration()); + + BookKeeper placementBookKeeper = mock(BookKeeper.class); + CreateBuilder createBuilder = mock(CreateBuilder.class, Mockito.RETURNS_SELF); + LedgerHandle ledgerHandle = mock(LedgerHandle.class); + CompletableFuture createResult = CompletableFuture.completedFuture(ledgerHandle); + doReturn(createBuilder).when(placementBookKeeper).newCreateLedgerOp(); + doReturn(createResult).when(createBuilder).execute(); + + EnsemblePlacementPolicyConfig placementPolicy = new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, + Map.of(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary", + IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "secondary")); + BookKeeperClientContext clientContext = + BookKeeperClientContext.create(placementBookKeeper, placementPolicy); + when(pulsar.getBookKeeperClientContext(ownerTopic)) + .thenReturn(CompletableFuture.completedFuture(clientContext)); + + BookkeeperSchemaStorage schemaStorage = new BookkeeperSchemaStorage(pulsar); + assertThat(schemaStorage.createLedger(schemaId).join()).isSameAs(ledgerHandle); + + verify(pulsar).getBookKeeperClientContext(ownerTopic); + verify(placementBookKeeper).newCreateLedgerOp(); + ArgumentCaptor> metadataCaptor = ArgumentCaptor.captor(); + verify(createBuilder).withCustomMetadata(metadataCaptor.capture()); + Map metadata = metadataCaptor.getValue(); + LedgerMetadataUtils.buildMetadataForSchema(schemaId) + .forEach((key, value) -> assertThat(metadata.get(key)).containsExactly(value)); + assertThat(EnsemblePlacementPolicyConfig.decode(metadata.get( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG))) + .isEqualTo(placementPolicy); + } + + @DataProvider(name = "nonCanonicalSchemaIds") + public static Object[][] nonCanonicalSchemaIds() { + return new Object[][] { + {"tenant/cluster/namespace/topic"}, + {"id2"}, + {"tenant/namespace/a:b"} + }; + } + + @Test(dataProvider = "nonCanonicalSchemaIds") + public void testCreateLedgerUsesDefaultClientForNonCanonicalSchemaId(String schemaId) throws Exception { + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getLocalMetadataStore()).thenReturn(mock(MetadataStoreExtended.class)); + when(pulsar.getConfiguration()).thenReturn(new ServiceConfiguration()); + BookKeeper defaultBookKeeper = mock(BookKeeper.class); + CreateBuilder createBuilder = mock(CreateBuilder.class, Mockito.RETURNS_SELF); + LedgerHandle ledgerHandle = mock(LedgerHandle.class); + CompletableFuture createResult = CompletableFuture.completedFuture(ledgerHandle); + doReturn(createBuilder).when(defaultBookKeeper).newCreateLedgerOp(); + doReturn(createResult).when(createBuilder).execute(); + BookKeeperClientFactory bookKeeperClientFactory = mock(BookKeeperClientFactory.class); + when(pulsar.getBookKeeperClientFactory()).thenReturn(bookKeeperClientFactory); + doReturn(CompletableFuture.completedFuture(defaultBookKeeper)) + .when(bookKeeperClientFactory).create(any(), any(), any(), any(), any()); + BookkeeperSchemaStorage schemaStorage = new BookkeeperSchemaStorage(pulsar); + schemaStorage.start(); + + CompletableFuture createFuture = schemaStorage.createLedger(schemaId); + + assertThat(createFuture.join()).isSameAs(ledgerHandle); + verify(pulsar, never()).getBookKeeperClientContext(any()); + ArgumentCaptor> metadataCaptor = ArgumentCaptor.captor(); + verify(createBuilder).withCustomMetadata(metadataCaptor.capture()); + Map metadata = metadataCaptor.getValue(); + LedgerMetadataUtils.buildMetadataForSchema(schemaId) + .forEach((key, value) -> assertThat(metadata.get(key)).containsExactly(value)); + assertThat(metadata).doesNotContainKey( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + } + + @Test + public void testCreateLedgerDoesNotFallbackWhenPlacementLookupFails() { + String schemaId = "tenant/namespace/topic"; + PulsarService pulsar = mock(PulsarService.class); + when(pulsar.getLocalMetadataStore()).thenReturn(mock(MetadataStoreExtended.class)); + when(pulsar.getConfiguration()).thenReturn(new ServiceConfiguration()); + RuntimeException failure = new RuntimeException("placement lookup failed"); + when(pulsar.getBookKeeperClientContext(TopicName.get(schemaId))) + .thenReturn(CompletableFuture.failedFuture(failure)); + BookkeeperSchemaStorage schemaStorage = new BookkeeperSchemaStorage(pulsar); + + CompletableFuture createFuture = schemaStorage.createLedger(schemaId); + + assertThat(createFuture).isCompletedExceptionally(); + assertThatThrownBy(createFuture::join).hasCause(failure); + } } From 1fa761e790b96fe4c1a4ddba6d57bd08d5bd911d Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Fri, 25 Sep 2026 00:22:19 +0800 Subject: [PATCH 3/4] [fix][broker] Apply namespace placement to compaction ledgers Resolve the BookKeeper client and ledger placement metadata from the compacted topic for publishing-order, event-time, strategic, and CLI compaction paths. Preserve default behavior when no policy is configured and fail compaction instead of silently bypassing a configured policy. Assisted-by: Codex --- .../apache/pulsar/broker/PulsarService.java | 2 +- .../compaction/AbstractTwoPhaseCompactor.java | 25 ++- .../apache/pulsar/compaction/Compactor.java | 15 +- .../pulsar/compaction/CompactorTool.java | 68 +++++- .../EventTimeCompactionServiceFactory.java | 2 +- .../compaction/EventTimeOrderCompactor.java | 14 +- .../compaction/PublishingOrderCompactor.java | 14 +- .../PulsarCompactionServiceFactory.java | 2 +- .../StrategicTwoPhaseCompactor.java | 14 ++ .../CompactorPlacementPolicyTest.java | 212 ++++++++++++++++++ .../pulsar/compaction/CompactorToolTest.java | 90 ++++++++ 11 files changed, 448 insertions(+), 10 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorPlacementPolicyTest.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 96cd15ec51cdb..af8fba739e38f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -1814,7 +1814,7 @@ public Compactor getNullableCompactor() { public StrategicTwoPhaseCompactor newStrategicCompactor() throws PulsarServerException { return new StrategicTwoPhaseCompactor(this.getConfiguration(), getClient(), getBookKeeperClient(), - getCompactorExecutor()); + getCompactorExecutor(), this::getBookKeeperClientContext); } public synchronized StrategicTwoPhaseCompactor getStrategicCompactor() throws PulsarServerException { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/AbstractTwoPhaseCompactor.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/AbstractTwoPhaseCompactor.java index 7255421b84ee6..c31e4d683d254 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/AbstractTwoPhaseCompactor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/AbstractTwoPhaseCompactor.java @@ -31,6 +31,7 @@ import java.util.concurrent.TimeUnit; import java.util.function.BiFunction; import java.util.function.BiPredicate; +import java.util.function.Function; import lombok.CustomLog; import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.client.BookKeeper; @@ -38,6 +39,7 @@ import org.apache.bookkeeper.mledger.impl.LedgerMetadataUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; @@ -46,6 +48,7 @@ import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.RawBatchConverter; import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.protocol.Markers; import org.apache.pulsar.common.util.Backoff; @@ -76,7 +79,15 @@ public AbstractTwoPhaseCompactor(ServiceConfiguration conf, PulsarClient pulsar, BookKeeper bk, ScheduledExecutorService scheduler) { - super(conf, pulsar, bk, scheduler); + this(conf, pulsar, bk, scheduler, null); + } + + public AbstractTwoPhaseCompactor(ServiceConfiguration conf, + PulsarClient pulsar, + BookKeeper bk, + ScheduledExecutorService scheduler, + Function> bookKeeperClientContextProvider) { + super(conf, pulsar, bk, scheduler, bookKeeperClientContextProvider); phaseOneLoopReadTimeout = Duration.ofSeconds( conf.getBrokerServiceCompactionPhaseOneLoopTimeInSeconds()); topicCompactionRetainNullKey = conf.isTopicCompactionRetainNullKey(); @@ -385,6 +396,18 @@ private void phaseTwoLoop(RawReader reader, MessageId to, Map protected CompletableFuture createLedger(BookKeeper bk, Map metadata, String topic) { + if (bookKeeperClientContextProvider == null) { + return createLedgerWithBookKeeper(bk, metadata, topic); + } + return CompletableFuture.completedFuture(topic) + .thenApply(TopicName::get) + .thenCompose(bookKeeperClientContextProvider) + .thenCompose(clientContext -> createLedgerWithBookKeeper( + clientContext.getBookKeeper(), clientContext.withPlacementMetadata(metadata), topic)); + } + + private CompletableFuture createLedgerWithBookKeeper(BookKeeper bk, + Map metadata, String topic) { CompletableFuture bkf = new CompletableFuture<>(); try { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/Compactor.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/Compactor.java index e8044de51b532..a46e3a49eda94 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/Compactor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/Compactor.java @@ -21,11 +21,14 @@ import static java.nio.charset.StandardCharsets.UTF_8; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; +import java.util.function.Function; import lombok.CustomLog; import org.apache.bookkeeper.client.BookKeeper; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.RawReader; +import org.apache.pulsar.common.naming.TopicName; /** * Compactor for Pulsar topics. @@ -41,16 +44,27 @@ public abstract class Compactor { protected final ScheduledExecutorService scheduler; protected final PulsarClient pulsar; protected final BookKeeper bk; + protected final Function> bookKeeperClientContextProvider; protected final CompactorMXBeanImpl mxBean; public Compactor(ServiceConfiguration conf, PulsarClient pulsar, BookKeeper bk, ScheduledExecutorService scheduler) { + this(conf, pulsar, bk, scheduler, null); + } + + public Compactor(ServiceConfiguration conf, + PulsarClient pulsar, + BookKeeper bk, + ScheduledExecutorService scheduler, + Function> + bookKeeperClientContextProvider) { this.conf = conf; this.scheduler = scheduler; this.pulsar = pulsar; this.bk = bk; + this.bookKeeperClientContextProvider = bookKeeperClientContextProvider; this.mxBean = new CompactorMXBeanImpl(); } @@ -90,4 +104,3 @@ public CompactorMXBean getStats() { return this.mxBean; } } - diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java index 34eee39dabcb2..60d67ab5349f7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java @@ -20,20 +20,29 @@ import static org.apache.commons.lang3.StringUtils.isBlank; import static org.apache.commons.lang3.StringUtils.isNotBlank; +import com.google.common.annotations.VisibleForTesting; import com.google.common.util.concurrent.ThreadFactoryBuilder; import io.netty.channel.EventLoopGroup; import io.netty.util.concurrent.DefaultThreadFactory; import java.nio.file.Path; +import java.util.Map; import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import lombok.Cleanup; import lombok.CustomLog; +import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.EnsemblePlacementPolicy; import org.apache.pulsar.broker.BookKeeperClientFactory; import org.apache.pulsar.broker.BookKeeperClientFactoryImpl; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.ServiceConfigurationUtils; +import org.apache.pulsar.broker.resources.LocalPoliciesResources; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.broker.storage.BookKeeperPlacementPolicyConfigResolver; import org.apache.pulsar.client.api.ClientBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; @@ -42,6 +51,10 @@ import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.client.internal.PropertiesUtils; import org.apache.pulsar.common.configuration.PulsarConfigurationLoader; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig.ParseEnsemblePlacementPolicyConfigException; +import org.apache.pulsar.common.policies.data.LocalPolicies; import org.apache.pulsar.common.util.netty.EventLoopUtil; import org.apache.pulsar.docs.tools.CmdGenerateDocs; import org.apache.pulsar.metadata.api.MetadataStoreConfig; @@ -144,6 +157,45 @@ public static PulsarClient createClient(ServiceConfiguration brokerConfig) throw return clientBuilder.build(); } + @VisibleForTesting + static CompletableFuture createBookKeeperClientContext( + ServiceConfiguration brokerConfig, + BookKeeperClientFactory bkClientFactory, + MetadataStoreExtended store, + EventLoopGroup eventLoopGroup, + TopicName topicName, + Optional localPolicies) { + return CompletableFuture.completedFuture(localPolicies) + .thenApply(policies -> BookKeeperPlacementPolicyConfigResolver.resolve( + brokerConfig, topicName, policies)) + .thenCompose(placementPolicy -> { + Optional> policyClass = placementPolicy + .map(EnsemblePlacementPolicyConfig::getPolicyClass) + .map(clazz -> clazz.asSubclass(EnsemblePlacementPolicy.class)); + Map policyProperties = placementPolicy + .map(EnsemblePlacementPolicyConfig::getProperties) + .orElse(null); + return bkClientFactory.create( + brokerConfig, store, eventLoopGroup, policyClass, policyProperties) + .thenApply(bookKeeper -> { + try { + return BookKeeperClientContext.create( + bookKeeper, placementPolicy.orElse(null)); + } catch (ParseEnsemblePlacementPolicyConfigException e) { + try { + bookKeeper.close(); + } catch (InterruptedException closeException) { + Thread.currentThread().interrupt(); + e.addSuppressed(closeException); + } catch (BKException closeException) { + e.addSuppressed(closeException); + } + throw new CompletionException(e); + } + }); + }); + } + public static void main(String[] args) throws Exception { Arguments arguments = new Arguments(); CommandLine commander = new CommandLine(arguments); @@ -180,6 +232,7 @@ public static void main(String[] args) throws Exception { """, filepath); throw new IllegalArgumentException(message); } + final TopicName topicName = TopicName.get(arguments.topic); @Cleanup(value = "shutdownNow") ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor( @@ -199,8 +252,15 @@ public static void main(String[] args) throws Exception { EventLoopGroup eventLoopGroup = EventLoopUtil.newEventLoopGroup(1, false, new DefaultThreadFactory("compactor-io")); + LocalPoliciesResources localPoliciesResources = new LocalPoliciesResources( + store, brokerConfig.getMetadataStoreOperationTimeoutSeconds()); + Optional localPolicies = + localPoliciesResources.getLocalPolicies(topicName.getNamespaceObject()); + BookKeeperClientContext bookKeeperClientContext = createBookKeeperClientContext( + brokerConfig, bkClientFactory, store, eventLoopGroup, topicName, localPolicies).get(); + @Cleanup - BookKeeper bk = bkClientFactory.create(brokerConfig, store, eventLoopGroup, Optional.empty(), null).get(); + BookKeeper bk = bookKeeperClientContext.getBookKeeper(); @Cleanup PulsarClient pulsar = createClient(brokerConfig); @@ -209,10 +269,12 @@ public static void main(String[] args) throws Exception { switch (arguments.compactorType) { case PUBLISHING: - compactor = new PublishingOrderCompactor(brokerConfig, pulsar, bk, scheduler); + compactor = new PublishingOrderCompactor(brokerConfig, pulsar, bk, scheduler, + ignored -> CompletableFuture.completedFuture(bookKeeperClientContext)); break; case EVENT_TIME: - compactor = new EventTimeOrderCompactor(brokerConfig, pulsar, bk, scheduler); + compactor = new EventTimeOrderCompactor(brokerConfig, pulsar, bk, scheduler, + ignored -> CompletableFuture.completedFuture(bookKeeperClientContext)); break; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeCompactionServiceFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeCompactionServiceFactory.java index 383c7b1aeedd6..f552af87d19fa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeCompactionServiceFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeCompactionServiceFactory.java @@ -28,6 +28,6 @@ protected Compactor newCompactor() throws PulsarServerException { PulsarService pulsarService = getPulsarService(); return new EventTimeOrderCompactor(pulsarService.getConfiguration(), pulsarService.getClient(), pulsarService.getBookKeeperClient(), - pulsarService.getCompactorExecutor()); + pulsarService.getCompactorExecutor(), pulsarService::getBookKeeperClientContext); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeOrderCompactor.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeOrderCompactor.java index 19269b2d52d33..f5c69457d48cc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeOrderCompactor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/EventTimeOrderCompactor.java @@ -22,18 +22,22 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; +import java.util.function.Function; import java.util.stream.Collectors; import lombok.CustomLog; import org.apache.bookkeeper.client.BookKeeper; import org.apache.commons.lang3.tuple.ImmutablePair; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.RawMessage; import org.apache.pulsar.client.impl.RawBatchConverter; import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.naming.TopicName; @CustomLog public class EventTimeOrderCompactor extends AbstractTwoPhaseCompactor> { @@ -44,6 +48,14 @@ public EventTimeOrderCompactor(ServiceConfiguration conf, super(conf, pulsar, bk, scheduler); } + public EventTimeOrderCompactor(ServiceConfiguration conf, + PulsarClient pulsar, + BookKeeper bk, + ScheduledExecutorService scheduler, + Function> bookKeeperClientContextProvider) { + super(conf, pulsar, bk, scheduler, bookKeeperClientContextProvider); + } + @Override protected Map toLatestMessageIdForKey( Map> latestForKey) { @@ -149,4 +161,4 @@ private List extractMessageCompactionDataFromBatch(RawMes throws IOException { return RawBatchConverter.extractMessageCompactionData(msg, metadata); } -} \ No newline at end of file +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PublishingOrderCompactor.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PublishingOrderCompactor.java index d2ab61da93bbc..83b2af150b53c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PublishingOrderCompactor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PublishingOrderCompactor.java @@ -21,17 +21,21 @@ import java.io.IOException; import java.util.List; import java.util.Map; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; +import java.util.function.Function; import lombok.CustomLog; import org.apache.bookkeeper.client.BookKeeper; import org.apache.commons.lang3.tuple.ImmutableTriple; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.RawMessage; import org.apache.pulsar.client.impl.RawBatchConverter; import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.naming.TopicName; @CustomLog public class PublishingOrderCompactor extends AbstractTwoPhaseCompactor { @@ -42,6 +46,14 @@ public PublishingOrderCompactor(ServiceConfiguration conf, super(conf, pulsar, bk, scheduler); } + public PublishingOrderCompactor(ServiceConfiguration conf, + PulsarClient pulsar, + BookKeeper bk, + ScheduledExecutorService scheduler, + Function> bookKeeperClientContextProvider) { + super(conf, pulsar, bk, scheduler, bookKeeperClientContextProvider); + } + @Override protected Map toLatestMessageIdForKey(Map latestForKey) { return latestForKey; @@ -121,4 +133,4 @@ protected List> extractIdsAndKeysAnd return RawBatchConverter.extractIdsAndKeysAndSize(msg, metadata); } -} \ No newline at end of file +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PulsarCompactionServiceFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PulsarCompactionServiceFactory.java index b44480378c62a..a3d70615f5029 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PulsarCompactionServiceFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/PulsarCompactionServiceFactory.java @@ -56,7 +56,7 @@ public Compactor getNullableCompactor() { protected Compactor newCompactor() throws PulsarServerException { return new PublishingOrderCompactor(pulsarService.getConfiguration(), pulsarService.getClient(), pulsarService.getBookKeeperClient(), - pulsarService.getCompactorExecutor()); + pulsarService.getCompactorExecutor(), pulsarService::getBookKeeperClientContext); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/StrategicTwoPhaseCompactor.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/StrategicTwoPhaseCompactor.java index 72290ab0c0127..2fe44b00de3b6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/StrategicTwoPhaseCompactor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/StrategicTwoPhaseCompactor.java @@ -29,12 +29,14 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Function; import lombok.CustomLog; import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.LedgerHandle; import org.apache.bookkeeper.mledger.impl.LedgerMetadataUtils; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.CryptoKeyReader; import org.apache.pulsar.client.api.Message; @@ -48,6 +50,7 @@ import org.apache.pulsar.client.impl.MessageImpl; import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.RawBatchMessageContainerImpl; +import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.topics.TopicCompactionStrategy; import org.apache.pulsar.common.util.FutureUtil; @@ -76,6 +79,17 @@ public StrategicTwoPhaseCompactor(ServiceConfiguration conf, phaseOneLoopReadTimeout = Duration.ofSeconds(conf.getBrokerServiceCompactionPhaseOneLoopTimeInSeconds()); } + public StrategicTwoPhaseCompactor(ServiceConfiguration conf, + PulsarClient pulsar, + BookKeeper bk, + ScheduledExecutorService scheduler, + Function> + bookKeeperClientContextProvider) { + super(conf, pulsar, bk, scheduler, bookKeeperClientContextProvider); + batchMessageContainer = new RawBatchMessageContainerImpl(); + phaseOneLoopReadTimeout = Duration.ofSeconds(conf.getBrokerServiceCompactionPhaseOneLoopTimeInSeconds()); + } + public CompletableFuture compact(String topic) { return FutureUtil.failedFuture(new UnsupportedOperationException()); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorPlacementPolicyTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorPlacementPolicyTest.java new file mode 100644 index 0000000000000..9c173ea84435b --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorPlacementPolicyTest.java @@ -0,0 +1,212 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.compaction; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.RETURNS_SELF; +import static org.mockito.Mockito.doCallRealMethod; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Function; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.LedgerHandle; +import org.apache.bookkeeper.client.api.CreateBuilder; +import org.apache.bookkeeper.mledger.impl.LedgerMetadataUtils; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.mockito.ArgumentCaptor; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Test; + +@Test(groups = "broker") +public class CompactorPlacementPolicyTest { + + private static final String TOPIC = "persistent://tenant/namespace/topic"; + + @DataProvider(name = "compactorTypes") + public static Object[][] compactorTypes() { + return new Object[][] { + {"publishing"}, + {"event-time"}, + {"strategic"} + }; + } + + @Test(dataProvider = "compactorTypes") + public void testCreateLedgerUsesPlacementClientAndMatchingMetadata(String compactorType) throws Exception { + BookKeeper defaultBookKeeper = mock(BookKeeper.class); + BookKeeper placementBookKeeper = mock(BookKeeper.class); + CreateBuilder createBuilder = mock(CreateBuilder.class, RETURNS_SELF); + LedgerHandle ledgerHandle = mock(LedgerHandle.class); + doReturn(createBuilder).when(placementBookKeeper).newCreateLedgerOp(); + doReturn(CompletableFuture.completedFuture(ledgerHandle)).when(createBuilder).execute(); + + EnsemblePlacementPolicyConfig placementPolicy = placementPolicy(); + BookKeeperClientContext clientContext = + BookKeeperClientContext.create(placementBookKeeper, placementPolicy); + AtomicReference requestedTopic = new AtomicReference<>(); + Function> provider = topicName -> { + requestedTopic.set(topicName); + return CompletableFuture.completedFuture(clientContext); + }; + AbstractTwoPhaseCompactor compactor = newCompactor(compactorType, defaultBookKeeper, provider); + Map compactionMetadata = LedgerMetadataUtils.buildMetadataForCompactedLedger( + TOPIC, new byte[] {1, 2, 3}); + + LedgerHandle result = compactor.createLedger(defaultBookKeeper, compactionMetadata, TOPIC).join(); + + assertThat(result).isSameAs(ledgerHandle); + assertThat(requestedTopic.get()).isEqualTo(TopicName.get(TOPIC)); + verify(placementBookKeeper).newCreateLedgerOp(); + verify(defaultBookKeeper, never()).newCreateLedgerOp(); + + @SuppressWarnings("unchecked") + ArgumentCaptor> metadataCaptor = ArgumentCaptor.forClass(Map.class); + verify(createBuilder).withCustomMetadata(metadataCaptor.capture()); + Map actualMetadata = metadataCaptor.getValue(); + compactionMetadata.forEach((key, value) -> assertThat(actualMetadata.get(key)).containsExactly(value)); + assertThat(EnsemblePlacementPolicyConfig.decode(actualMetadata.get( + EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG))) + .isEqualTo(placementPolicy); + } + + @Test + public void testFailedProviderCompletesCreateFutureExceptionally() { + RuntimeException failure = new RuntimeException("placement lookup failed"); + assertProviderFailureIsAsynchronous(topicName -> CompletableFuture.failedFuture(failure), failure); + } + + @Test + public void testThrowingProviderCompletesCreateFutureExceptionally() { + RuntimeException failure = new RuntimeException("provider threw"); + assertProviderFailureIsAsynchronous(topicName -> { + throw failure; + }, failure); + } + + @DataProvider(name = "brokerFactoryTypes") + public static Object[][] brokerFactoryTypes() { + return new Object[][] { + {false, PublishingOrderCompactor.class}, + {true, EventTimeOrderCompactor.class} + }; + } + + @Test(dataProvider = "brokerFactoryTypes") + public void testBrokerFactoryWiresTopicClientProvider( + boolean eventTime, Class expectedType) throws Exception { + PulsarService pulsarService = mock(PulsarService.class); + ServiceConfiguration configuration = new ServiceConfiguration(); + BookKeeper defaultBookKeeper = mock(BookKeeper.class); + PulsarClient pulsarClient = mock(PulsarClient.class); + ScheduledExecutorService scheduler = mock(ScheduledExecutorService.class); + BookKeeperClientContext expectedContext = BookKeeperClientContext.create( + mock(BookKeeper.class), placementPolicy()); + TopicName topicName = TopicName.get(TOPIC); + when(pulsarService.getConfiguration()).thenReturn(configuration); + when(pulsarService.getClient()).thenReturn(pulsarClient); + when(pulsarService.getBookKeeperClient()).thenReturn(defaultBookKeeper); + when(pulsarService.getCompactorExecutor()).thenReturn(scheduler); + when(pulsarService.getBookKeeperClientContext(topicName)) + .thenReturn(CompletableFuture.completedFuture(expectedContext)); + PulsarCompactionServiceFactory factory = eventTime + ? new EventTimeCompactionServiceFactory() : new PulsarCompactionServiceFactory(); + factory.initialize(pulsarService).join(); + + Compactor compactor = factory.getCompactor(); + + assertThat(compactor).isInstanceOf(expectedType); + assertThat(compactor.bookKeeperClientContextProvider.apply(topicName).join()).isSameAs(expectedContext); + verify(pulsarService).getBookKeeperClientContext(topicName); + } + + @Test + public void testPulsarServiceWiresStrategicCompactorTopicClientProvider() throws Exception { + PulsarService pulsarService = mock(PulsarService.class); + TopicName topicName = TopicName.get(TOPIC); + BookKeeperClientContext expectedContext = BookKeeperClientContext.create( + mock(BookKeeper.class), placementPolicy()); + when(pulsarService.getConfiguration()).thenReturn(new ServiceConfiguration()); + when(pulsarService.getClient()).thenReturn(mock(PulsarClient.class)); + when(pulsarService.getBookKeeperClient()).thenReturn(mock(BookKeeper.class)); + when(pulsarService.getCompactorExecutor()).thenReturn(mock(ScheduledExecutorService.class)); + when(pulsarService.getBookKeeperClientContext(topicName)) + .thenReturn(CompletableFuture.completedFuture(expectedContext)); + doCallRealMethod().when(pulsarService).newStrategicCompactor(); + + StrategicTwoPhaseCompactor compactor = pulsarService.newStrategicCompactor(); + + assertThat(compactor.bookKeeperClientContextProvider.apply(topicName).join()).isSameAs(expectedContext); + verify(pulsarService).getBookKeeperClientContext(topicName); + } + + private static void assertProviderFailureIsAsynchronous( + Function> provider, + RuntimeException failure) { + BookKeeper defaultBookKeeper = mock(BookKeeper.class); + AbstractTwoPhaseCompactor compactor = newCompactor("publishing", defaultBookKeeper, provider); + AtomicReference> result = new AtomicReference<>(); + + assertThatCode(() -> result.set(compactor.createLedger(defaultBookKeeper, Map.of(), TOPIC))) + .doesNotThrowAnyException(); + + assertThat(result.get()).isCompletedExceptionally(); + assertThatThrownBy(result.get()::join).hasCause(failure); + verify(defaultBookKeeper, never()).newCreateLedgerOp(); + } + + private static AbstractTwoPhaseCompactor newCompactor( + String compactorType, + BookKeeper defaultBookKeeper, + Function> provider) { + ServiceConfiguration configuration = new ServiceConfiguration(); + PulsarClient pulsarClient = mock(PulsarClient.class); + ScheduledExecutorService scheduler = mock(ScheduledExecutorService.class); + return switch (compactorType) { + case "publishing" -> new PublishingOrderCompactor( + configuration, pulsarClient, defaultBookKeeper, scheduler, provider); + case "event-time" -> new EventTimeOrderCompactor( + configuration, pulsarClient, defaultBookKeeper, scheduler, provider); + case "strategic" -> new StrategicTwoPhaseCompactor( + configuration, pulsarClient, defaultBookKeeper, scheduler, provider); + default -> throw new IllegalArgumentException("Unknown compactor type: " + compactorType); + }; + } + + private static EnsemblePlacementPolicyConfig placementPolicy() { + return new EnsemblePlacementPolicyConfig( + IsolatedBookieEnsemblePlacementPolicy.class, + Map.of(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary", + IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "secondary")); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorToolTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorToolTest.java index c089ba7520d34..bd3328e905d8f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorToolTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactorToolTest.java @@ -18,23 +18,46 @@ */ package org.apache.pulsar.compaction; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.same; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertSame; import static org.testng.Assert.assertTrue; +import io.netty.channel.EventLoopGroup; import java.io.ByteArrayOutputStream; import java.io.PrintStream; import java.lang.reflect.Constructor; import java.lang.reflect.Field; import java.util.Arrays; +import java.util.Collections; +import java.util.Map; import java.util.Optional; import java.util.Properties; +import java.util.concurrent.CompletableFuture; import lombok.Cleanup; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.EnsemblePlacementPolicy; +import org.apache.pulsar.bookie.rackawareness.IsolatedBookieEnsemblePlacementPolicy; +import org.apache.pulsar.broker.BookKeeperClientFactory; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.policies.data.LocalPolicies; import org.apache.pulsar.docs.tools.CmdGenerateDocs; +import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import org.mockito.ArgumentCaptor; import org.testng.annotations.Test; import picocli.CommandLine.Option; @@ -43,6 +66,73 @@ */ public class CompactorToolTest { + @Test + @SuppressWarnings({"unchecked", "rawtypes"}) + public void testBookKeeperClientContextUsesNamespaceAffinity() throws Exception { + ServiceConfiguration brokerConfig = new ServiceConfiguration(); + BookKeeperClientFactory bkClientFactory = mock(BookKeeperClientFactory.class); + MetadataStoreExtended store = mock(MetadataStoreExtended.class); + EventLoopGroup eventLoopGroup = mock(EventLoopGroup.class); + BookKeeper bookKeeper = mock(BookKeeper.class); + TopicName topicName = TopicName.get("persistent://tenant/namespace/topic"); + BookieAffinityGroupData affinityGroup = BookieAffinityGroupData.builder() + .bookkeeperAffinityGroupPrimary("primary") + .bookkeeperAffinityGroupSecondary("secondary") + .build(); + Optional localPolicies = Optional.of(new LocalPolicies(null, affinityGroup, null)); + + when(bkClientFactory.create(same(brokerConfig), same(store), same(eventLoopGroup), any(), any())) + .thenReturn(CompletableFuture.completedFuture(bookKeeper)); + + BookKeeperClientContext context = CompactorTool.createBookKeeperClientContext( + brokerConfig, bkClientFactory, store, eventLoopGroup, topicName, localPolicies).get(); + + ArgumentCaptor>> policyClass = + ArgumentCaptor.forClass(Optional.class); + ArgumentCaptor> policyProperties = ArgumentCaptor.forClass(Map.class); + verify(bkClientFactory).create(same(brokerConfig), same(store), same(eventLoopGroup), + policyClass.capture(), policyProperties.capture()); + assertEquals(policyClass.getValue(), Optional.of(IsolatedBookieEnsemblePlacementPolicy.class)); + assertEquals(policyProperties.getValue().get( + IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS), "primary"); + assertEquals(policyProperties.getValue().get( + IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS), "secondary"); + assertSame(context.getBookKeeper(), bookKeeper); + + byte[] encodedPolicy = context.withPlacementMetadata(Collections.emptyMap()) + .get(EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + EnsemblePlacementPolicyConfig decodedPolicy = EnsemblePlacementPolicyConfig.decode(encodedPolicy); + assertEquals(decodedPolicy.getPolicyClass(), IsolatedBookieEnsemblePlacementPolicy.class); + assertEquals(decodedPolicy.getProperties(), policyProperties.getValue()); + } + + @Test + @SuppressWarnings({"unchecked", "rawtypes"}) + public void testBookKeeperClientContextPreservesDefaultPlacementWithoutAffinity() throws Exception { + ServiceConfiguration brokerConfig = new ServiceConfiguration(); + BookKeeperClientFactory bkClientFactory = mock(BookKeeperClientFactory.class); + MetadataStoreExtended store = mock(MetadataStoreExtended.class); + EventLoopGroup eventLoopGroup = mock(EventLoopGroup.class); + BookKeeper bookKeeper = mock(BookKeeper.class); + when(bkClientFactory.create(same(brokerConfig), same(store), same(eventLoopGroup), any(), any())) + .thenReturn(CompletableFuture.completedFuture(bookKeeper)); + + BookKeeperClientContext context = CompactorTool.createBookKeeperClientContext( + brokerConfig, bkClientFactory, store, eventLoopGroup, + TopicName.get("persistent://tenant/namespace/topic"), Optional.empty()).get(); + + ArgumentCaptor>> policyClass = + ArgumentCaptor.forClass(Optional.class); + ArgumentCaptor> policyProperties = ArgumentCaptor.forClass(Map.class); + verify(bkClientFactory).create(same(brokerConfig), same(store), same(eventLoopGroup), + policyClass.capture(), policyProperties.capture()); + assertEquals(policyClass.getValue(), Optional.empty()); + assertNull(policyProperties.getValue()); + assertSame(context.getBookKeeper(), bookKeeper); + assertFalse(context.withPlacementMetadata(Collections.emptyMap()) + .containsKey(EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG)); + } + /** * Test broker-tool generate docs. * From 6498854d74802e69c061a53c7e21f5bb8a26f27b Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Thu, 1 Oct 2026 22:45:50 +0800 Subject: [PATCH 4/4] [fix][broker] Preserve auxiliary clients and await schema close Wait for schema ledger closure before publishing its locator, and clean up the orphan ledger when closure fails. Use the existing schema and bucket clients when no custom placement policy applies. Document that custom placement uses the default BookKeeper backend. Assisted-by: Codex --- .../apache/pulsar/broker/PulsarService.java | 28 +++- .../BookkeeperBucketSnapshotStorage.java | 2 +- .../schema/BookkeeperSchemaStorage.java | 13 +- ...sarServiceBookKeeperClientContextTest.java | 121 ++++++++++++++++++ ...perBucketSnapshotStoragePlacementTest.java | 6 +- .../schema/BookkeeperSchemaStorageTest.java | 85 +++++++++++- 6 files changed, 242 insertions(+), 13 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceBookKeeperClientContextTest.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index af8fba739e38f..b178866c87e9f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -1664,12 +1664,28 @@ public BookKeeper getBookKeeperClient() { } /** - * Resolve the placement policy for a topic and borrow the corresponding shared BookKeeper client. + * Resolve the namespace placement policy for an auxiliary ledger using the broker's default BookKeeper storage. + * The topic's managed-ledger storage class is not used to select the client. * * @param topicName the topic that owns the ledger * @return a future that completes with the BookKeeper client and matching placement metadata */ public CompletableFuture getBookKeeperClientContext(TopicName topicName) { + return getBookKeeperClientContext(topicName, this::getBookKeeperClient); + } + + /** + * Resolve the namespace placement policy for an auxiliary ledger, using the caller's BookKeeper client when there + * is no custom policy. A custom policy uses the broker's default BookKeeper storage class, not the topic's + * managed-ledger storage class. The caller's client must access the same BookKeeper backend for later reads and + * deletion. + * + * @param topicName the topic that owns the ledger + * @param fallbackDefaultClient the caller's existing BookKeeper client, used only when there is no custom policy + * @return a future that completes with the BookKeeper client and matching placement metadata + */ + public CompletableFuture getBookKeeperClientContext( + TopicName topicName, Supplier fallbackDefaultClient) { return CompletableFuture.completedFuture(topicName) .thenCompose(name -> { Objects.requireNonNull(name, "topicName"); @@ -1677,8 +1693,16 @@ public CompletableFuture getBookKeeperClientContext(Top .getLocalPoliciesAsync(name.getNamespaceObject()); }).thenCompose(localPolicies -> { EnsemblePlacementPolicyConfig placementPolicyConfig = - BookKeeperPlacementPolicyConfigResolver.resolve(config, topicName, localPolicies) + BookKeeperPlacementPolicyConfigResolver.resolve(getConfig(), topicName, localPolicies) .orElse(null); + if (placementPolicyConfig == null) { + try { + return CompletableFuture.completedFuture( + BookKeeperClientContext.create(fallbackDefaultClient.get(), null)); + } catch (ParseEnsemblePlacementPolicyConfigException e) { + return CompletableFuture.failedFuture(e); + } + } ManagedLedgerStorageClass defaultStorageClass = getManagedLedgerStorage().getDefaultStorageClass(); if (!(defaultStorageClass instanceof BookkeeperManagedLedgerStorageClass bkStorageClass)) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java index b1e4bc57aaa7e..a84f91a4650a7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStorage.java @@ -177,7 +177,7 @@ private List parseSnapshotSegmentEntries(Enumeration createLedger(String bucketKey, String topicName, String cursorName) { return CompletableFuture.completedFuture(topicName) .thenApply(TopicName::get) - .thenCompose(pulsar::getBookKeeperClientContext) + .thenCompose(name -> pulsar.getBookKeeperClientContext(name, () -> bookKeeper)) .thenCompose(bookKeeperClientContext -> { CompletableFuture future = new CompletableFuture<>(); Map metadata = bookKeeperClientContext.withPlacementMetadata( diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java index 2612dc8af585c..a16d7634d431e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorage.java @@ -468,10 +468,13 @@ private CompletableFuture addNewSchemaEntryToStore( return createLedger(schemaId).thenCompose(ledgerHandle -> { final long ledgerId = ledgerHandle.getId(); return addEntry(ledgerHandle, schemaEntry) - .thenApply(entryId -> { - ledgerHandle.closeAsync(); - return Functions.newPositionInfo(ledgerId, entryId); - }); + .thenCompose(entryId -> FutureUtil.supplySafely(ledgerHandle::closeAsync) + .thenApply(__ -> Functions.newPositionInfo(ledgerId, entryId)) + .exceptionallyCompose(ex -> { + Throwable cause = FutureUtil.unwrapCompletionException(ex); + return deleteLedgerAsync(schemaId, ledgerId, "schema ledger close failed") + .thenCompose(__ -> FutureUtil.failedFuture(cause)); + })); }); } @@ -653,7 +656,7 @@ CompletableFuture createLedger(String schemaId) { if (topicName.isEmpty()) { return createLedger(id, bookKeeper, LedgerMetadataUtils.buildMetadataForSchema(id)); } - return pulsar.getBookKeeperClientContext(topicName.get()) + return pulsar.getBookKeeperClientContext(topicName.get(), () -> bookKeeper) .thenCompose(clientContext -> createLedger(id, clientContext)); }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceBookKeeperClientContextTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceBookKeeperClientContextTest.java new file mode 100644 index 0000000000000..5362ba7a2b106 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceBookKeeperClientContextTest.java @@ -0,0 +1,121 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.CALLS_REAL_METHODS; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.pulsar.broker.resources.LocalPoliciesResources; +import org.apache.pulsar.broker.resources.PulsarResources; +import org.apache.pulsar.broker.storage.BookKeeperClientContext; +import org.apache.pulsar.broker.storage.BookkeeperManagedLedgerStorageClass; +import org.apache.pulsar.broker.storage.ManagedLedgerStorage; +import org.apache.pulsar.broker.storage.ManagedLedgerStorageClass; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.testng.annotations.Test; + +@Test(groups = "broker") +public class PulsarServiceBookKeeperClientContextTest { + + private static final TopicName TOPIC = TopicName.get("persistent://tenant/namespace/topic"); + + @Test + public void testNoPolicyUsesCallerClientWithNonBookKeeperDefaultStorage() { + ManagedLedgerStorage managedLedgerStorage = mock(ManagedLedgerStorage.class); + when(managedLedgerStorage.getDefaultStorageClass()).thenReturn(mock(ManagedLedgerStorageClass.class)); + PulsarService pulsar = mockPulsar(new ServiceConfiguration(), managedLedgerStorage); + BookKeeper callerClient = mock(BookKeeper.class); + + BookKeeperClientContext context = pulsar.getBookKeeperClientContext(TOPIC, () -> callerClient).join(); + + assertThat(context.getBookKeeper()).isSameAs(callerClient); + assertThat(context.withPlacementMetadata(Map.of())) + .doesNotContainKey(EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + verify(managedLedgerStorage, never()).getDefaultStorageClass(); + } + + @Test + public void testCustomPolicyUsesPolicyClientWithoutCallingFallback() { + ServiceConfiguration configuration = new ServiceConfiguration(); + configuration.setStrictBookieAffinityEnabled(true); + ManagedLedgerStorage managedLedgerStorage = mock(ManagedLedgerStorage.class); + BookkeeperManagedLedgerStorageClass storageClass = mock(BookkeeperManagedLedgerStorageClass.class); + BookKeeper policyClient = mock(BookKeeper.class); + when(managedLedgerStorage.getDefaultStorageClass()).thenReturn(storageClass); + when(storageClass.getBookKeeperClient(any(EnsemblePlacementPolicyConfig.class))) + .thenReturn(CompletableFuture.completedFuture(policyClient)); + PulsarService pulsar = mockPulsar(configuration, managedLedgerStorage); + AtomicBoolean fallbackCalled = new AtomicBoolean(); + + BookKeeperClientContext context = pulsar.getBookKeeperClientContext(TOPIC, () -> { + fallbackCalled.set(true); + return mock(BookKeeper.class); + }).join(); + + assertThat(context.getBookKeeper()).isSameAs(policyClient); + assertThat(context.withPlacementMetadata(Map.of())) + .containsKey(EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); + assertThat(fallbackCalled.get()).isFalse(); + verify(storageClass).getBookKeeperClient(any(EnsemblePlacementPolicyConfig.class)); + } + + @Test + public void testCustomPolicyFailsWithNonBookKeeperDefaultStorage() { + ServiceConfiguration configuration = new ServiceConfiguration(); + configuration.setStrictBookieAffinityEnabled(true); + ManagedLedgerStorage managedLedgerStorage = mock(ManagedLedgerStorage.class); + when(managedLedgerStorage.getDefaultStorageClass()).thenReturn(mock(ManagedLedgerStorageClass.class)); + PulsarService pulsar = mockPulsar(configuration, managedLedgerStorage); + AtomicBoolean fallbackCalled = new AtomicBoolean(); + + CompletableFuture future = pulsar.getBookKeeperClientContext(TOPIC, () -> { + fallbackCalled.set(true); + return mock(BookKeeper.class); + }); + + assertThatThrownBy(future::join).hasCauseInstanceOf(UnsupportedOperationException.class); + assertThat(fallbackCalled.get()).isFalse(); + } + + private static PulsarService mockPulsar(ServiceConfiguration configuration, + ManagedLedgerStorage managedLedgerStorage) { + PulsarService pulsar = mock(PulsarService.class, CALLS_REAL_METHODS); + PulsarResources resources = mock(PulsarResources.class); + LocalPoliciesResources localPolicies = mock(LocalPoliciesResources.class); + doReturn(configuration).when(pulsar).getConfig(); + doReturn(resources).when(pulsar).getPulsarResources(); + doReturn(managedLedgerStorage).when(pulsar).getManagedLedgerStorage(); + when(resources.getLocalPolicies()).thenReturn(localPolicies); + when(localPolicies.getLocalPoliciesAsync(TOPIC.getNamespaceObject())) + .thenReturn(CompletableFuture.completedFuture(Optional.empty())); + return pulsar; + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java index bf8de4a62edc1..1ac44e5323fd0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BookkeeperBucketSnapshotStoragePlacementTest.java @@ -20,6 +20,8 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.RETURNS_SELF; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; @@ -57,7 +59,7 @@ public void testCreateLedgerUsesTopicClientAndMatchingPlacementMetadata() throws Map.of(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "primary", IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "secondary")); BookKeeperClientContext context = BookKeeperClientContext.create(policyBookKeeper, placementPolicy); - when(pulsar.getBookKeeperClientContext(TopicName.get(TOPIC))) + when(pulsar.getBookKeeperClientContext(eq(TopicName.get(TOPIC)), any())) .thenReturn(CompletableFuture.completedFuture(context)); CreateBuilder createBuilder = mock(CreateBuilder.class, RETURNS_SELF); @@ -99,7 +101,7 @@ public void testPlacementLookupFailureIsRetriablePersistenceFailure() { PulsarService pulsar = mock(PulsarService.class); when(pulsar.getConfig()).thenReturn(new ServiceConfiguration()); RuntimeException lookupFailure = new RuntimeException("Failed to read local policies"); - when(pulsar.getBookKeeperClientContext(TopicName.get(TOPIC))) + when(pulsar.getBookKeeperClientContext(eq(TopicName.get(TOPIC)), any())) .thenReturn(FutureUtil.failedFuture(lookupFailure)); BookkeeperBucketSnapshotStorage storage = new BookkeeperBucketSnapshotStorage(pulsar); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java index f524d932f9bda..9e548af07a8e9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/schema/BookkeeperSchemaStorageTest.java @@ -22,6 +22,8 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; @@ -32,7 +34,9 @@ import static org.testng.Assert.assertTrue; import java.nio.ByteBuffer; import java.util.Map; +import java.util.Optional; import java.util.concurrent.CompletableFuture; +import org.apache.bookkeeper.client.AsyncCallback; import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.LedgerHandle; import org.apache.bookkeeper.client.api.BKException; @@ -47,9 +51,14 @@ import org.apache.pulsar.broker.storage.BookKeeperClientContext; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; +import org.apache.pulsar.common.protocol.schema.SchemaVersion; import org.apache.pulsar.common.schema.LongSchemaVersion; +import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.metadata.api.MetadataSerde; +import org.apache.pulsar.metadata.api.Stat; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; import org.mockito.ArgumentCaptor; +import org.mockito.ArgumentMatchers; import org.mockito.Mockito; import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @@ -126,13 +135,13 @@ public void testCreateLedgerUsesTopicPlacementClientAndMetadata(String schemaId, IsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, "secondary")); BookKeeperClientContext clientContext = BookKeeperClientContext.create(placementBookKeeper, placementPolicy); - when(pulsar.getBookKeeperClientContext(ownerTopic)) + when(pulsar.getBookKeeperClientContext(eq(ownerTopic), any())) .thenReturn(CompletableFuture.completedFuture(clientContext)); BookkeeperSchemaStorage schemaStorage = new BookkeeperSchemaStorage(pulsar); assertThat(schemaStorage.createLedger(schemaId).join()).isSameAs(ledgerHandle); - verify(pulsar).getBookKeeperClientContext(ownerTopic); + verify(pulsar).getBookKeeperClientContext(eq(ownerTopic), any()); verify(placementBookKeeper).newCreateLedgerOp(); ArgumentCaptor> metadataCaptor = ArgumentCaptor.captor(); verify(createBuilder).withCustomMetadata(metadataCaptor.capture()); @@ -184,6 +193,76 @@ public void testCreateLedgerUsesDefaultClientForNonCanonicalSchemaId(String sche EnsemblePlacementPolicyConfig.ENSEMBLE_PLACEMENT_POLICY_CONFIG); } + @DataProvider(name = "ledgerCloseResults") + public static Object[][] ledgerCloseResults() { + return new Object[][] {{true}, {false}}; + } + + @Test(dataProvider = "ledgerCloseResults") + public void testSchemaLocatorWaitsForLedgerClose(boolean closeSucceeds) throws Exception { + String schemaId = "id2"; + PulsarService pulsar = mock(PulsarService.class); + MetadataStoreExtended store = mock(MetadataStoreExtended.class); + when(pulsar.getLocalMetadataStore()).thenReturn(store); + when(pulsar.getConfiguration()).thenReturn(new ServiceConfiguration()); + @SuppressWarnings("unchecked") + MetadataCache locatorCache = mock(MetadataCache.class); + doReturn(locatorCache).when(store).getMetadataCache( + ArgumentMatchers.>any()); + when(locatorCache.getWithStats("/schemas/" + schemaId)) + .thenReturn(CompletableFuture.completedFuture(Optional.empty())); + when(store.put(eq("/schemas/" + schemaId), any(byte[].class), eq(Optional.of(-1L)))) + .thenReturn(CompletableFuture.completedFuture(mock(Stat.class))); + + BookKeeper bookKeeper = mock(BookKeeper.class); + BookKeeperClientFactory bookKeeperClientFactory = mock(BookKeeperClientFactory.class); + when(pulsar.getBookKeeperClientFactory()).thenReturn(bookKeeperClientFactory); + doReturn(CompletableFuture.completedFuture(bookKeeper)) + .when(bookKeeperClientFactory).create(any(), any(), any(), any(), any()); + CreateBuilder createBuilder = mock(CreateBuilder.class, Mockito.RETURNS_SELF); + LedgerHandle ledgerHandle = mock(LedgerHandle.class); + when(ledgerHandle.getId()).thenReturn(42L); + doReturn(createBuilder).when(bookKeeper).newCreateLedgerOp(); + doReturn(CompletableFuture.completedFuture(ledgerHandle)).when(createBuilder).execute(); + doAnswer(invocation -> { + AsyncCallback.AddCallback callback = invocation.getArgument(1); + callback.addComplete(BKException.Code.OK, ledgerHandle, 0L, null); + return null; + }).when(ledgerHandle).asyncAddEntry(any(byte[].class), any(), any()); + doAnswer(invocation -> { + AsyncCallback.DeleteCallback callback = invocation.getArgument(1); + callback.deleteComplete(BKException.Code.OK, null); + return null; + }).when(bookKeeper).asyncDeleteLedger(eq(42L), any(), any()); + CompletableFuture closeFuture = new CompletableFuture<>(); + when(ledgerHandle.closeAsync()).thenReturn(closeFuture); + + BookkeeperSchemaStorage schemaStorage = new BookkeeperSchemaStorage(pulsar); + schemaStorage.start(); + try { + CompletableFuture putFuture = schemaStorage.put(schemaId, new byte[] {1}, new byte[] {2}); + + verify(ledgerHandle).closeAsync(); + assertThat(putFuture.isDone()).isFalse(); + verify(store, never()).put(eq("/schemas/" + schemaId), any(byte[].class), eq(Optional.of(-1L))); + + if (closeSucceeds) { + closeFuture.complete(null); + assertThat(putFuture.join()).isEqualTo(new LongSchemaVersion(0)); + verify(store).put(eq("/schemas/" + schemaId), any(byte[].class), eq(Optional.of(-1L))); + verify(bookKeeper, never()).asyncDeleteLedger(eq(42L), any(), any()); + } else { + RuntimeException failure = new RuntimeException("ledger close failed"); + closeFuture.completeExceptionally(failure); + assertThatThrownBy(putFuture::join).hasCause(failure); + verify(store, never()).put(eq("/schemas/" + schemaId), any(byte[].class), eq(Optional.of(-1L))); + verify(bookKeeper).asyncDeleteLedger(eq(42L), any(), any()); + } + } finally { + schemaStorage.close(); + } + } + @Test public void testCreateLedgerDoesNotFallbackWhenPlacementLookupFails() { String schemaId = "tenant/namespace/topic"; @@ -191,7 +270,7 @@ public void testCreateLedgerDoesNotFallbackWhenPlacementLookupFails() { when(pulsar.getLocalMetadataStore()).thenReturn(mock(MetadataStoreExtended.class)); when(pulsar.getConfiguration()).thenReturn(new ServiceConfiguration()); RuntimeException failure = new RuntimeException("placement lookup failed"); - when(pulsar.getBookKeeperClientContext(TopicName.get(schemaId))) + when(pulsar.getBookKeeperClientContext(eq(TopicName.get(schemaId)), any())) .thenReturn(CompletableFuture.failedFuture(failure)); BookkeeperSchemaStorage schemaStorage = new BookkeeperSchemaStorage(pulsar);