Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,16 @@ public StatsProvider getStatsProvider() {
public BookKeeper getBookKeeperClient() {
return defaultBkClient;
}

@Override
public CompletableFuture<BookKeeper> getBookKeeperClient(
EnsemblePlacementPolicyConfig ensemblePlacementPolicyConfig) {
try {
return bkFactory.get(ensemblePlacementPolicyConfig);
} catch (RuntimeException e) {
return CompletableFuture.failedFuture(e);
}
}
};
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -1659,6 +1663,64 @@ public BookKeeper getBookKeeperClient() {
}
}

/**
* 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<BookKeeperClientContext> 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<BookKeeperClientContext> getBookKeeperClientContext(
TopicName topicName, Supplier<BookKeeper> fallbackDefaultClient) {
return CompletableFuture.completedFuture(topicName)
.thenCompose(name -> {
Objects.requireNonNull(name, "topicName");
return getPulsarResources().getLocalPolicies()
.getLocalPoliciesAsync(name.getNamespaceObject());
}).thenCompose(localPolicies -> {
EnsemblePlacementPolicyConfig placementPolicyConfig =
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();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One related question: this always borrows the client from getDefaultStorageClass(), while normal topic loading honors PersistencePolicies.managedLedgerStorageClassName. With a multi-storage-class ManagedLedgerStorage, the owner topic can therefore use a different storage class/client than this "topic-aware" context.

Is the intended contract only "apply the namespace affinity to the broker's default BookKeeper storage", or should auxiliary ledgers follow the topic's actual storage class as well? If the former is intentional, I think it would be worth documenting that limitation explicitly.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The intended contract is to apply namespace affinity within the existing BookKeeper auxiliary-ledger backend. This PR does not route auxiliary ledgers by the topic's managed-ledger storage class. With no custom policy, schema and bucket retain their existing clients; with a custom policy, the context borrows from the default BookKeeper storage class. Following the topic's actual storage class would also require persisting backend identity with each auxiliary-ledger reference and routing reads and deletes accordingly. I documented this boundary in the getBookKeeperClientContext Javadoc in 6498854.

if (!(defaultStorageClass instanceof BookkeeperManagedLedgerStorageClass bkStorageClass)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

placementPolicyConfig can be null here, but we still require the default managed-ledger storage class to be BookkeeperManagedLedgerStorageClass. For schema/bucket this looks like a behavior regression unrelated to affinity: both components already own a BookKeeper client created through BookKeeperClientFactory, so they could work with a custom/non-BK ManagedLedgerStorage before this change.

Could we require a policy-aware storage class only when an effective custom placement policy is present, and preserve the existing component/default client path when there is no custom policy? The PR description says unsupported storage should fail when it cannot honor a custom placement policy, but the current check also fails the no-policy case.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 6498854. When there is no effective custom placement policy, the new overload uses the caller's existing BookKeeper client; schema and bucket pass the clients they already own. A custom policy still requires a BookkeeperManagedLedgerStorageClass so that we fail rather than silently ignore it. I added tests for the no-policy/non-BK-default case and both custom-policy cases.

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();
}
Expand Down Expand Up @@ -1776,7 +1838,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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -172,28 +173,40 @@ private List<SnapshotSegment> parseSnapshotSegmentEntries(Enumeration<LedgerEntr
}

@NonNull
private CompletableFuture<LedgerHandle> createLedger(String bucketKey, String topicName, String cursorName) {
CompletableFuture<LedgerHandle> future = new CompletableFuture<>();
Map<String, byte[]> 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<LedgerHandle> createLedger(String bucketKey, String topicName, String cursorName) {
return CompletableFuture.completedFuture(topicName)
.thenApply(TopicName::get)
.thenCompose(name -> pulsar.getBookKeeperClientContext(name, () -> bookKeeper))
.thenCompose(bookKeeperClientContext -> {
CompletableFuture<LedgerHandle> future = new CompletableFuture<>();
Map<String, byte[]> 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<LedgerHandle> getLedgerHandle(Long ledgerId) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -2481,39 +2483,14 @@ private CompletableFuture<ManagedLedgerConfig> getManagedLedgerConfig(
managedLedgerConfig.setAckQuorumSize(persistencePolicies.getBookkeeperAckQuorum());
managedLedgerConfig.setStorageClassName(persistencePolicies.getManagedLedgerStorageClassName());

if (serviceConfig.isStrictBookieAffinityEnabled()) {
managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyClassName(
IsolatedBookieEnsemblePlacementPolicy.class);
if (localPolicies.isPresent() && localPolicies.get().bookieAffinityGroup != null) {
Map<String, Object> 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<String, Object> properties = new HashMap<>();
properties.put(IsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, "*");
properties.put(IsolatedBookieEnsemblePlacementPolicy
.SECONDARY_ISOLATION_BOOKIE_GROUPS, "*");
managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(properties);
} else {
Map<String, Object> 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<String, Object> 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<EnsemblePlacementPolicyConfig> placementPolicy =
BookKeeperPlacementPolicyConfigResolver.resolve(serviceConfig, topicName, localPolicies);
if (placementPolicy.isPresent()) {
EnsemblePlacementPolicyConfig policyConfig = placementPolicy.get();
Class<? extends EnsemblePlacementPolicy> policyClass =
policyConfig.getPolicyClass().asSubclass(EnsemblePlacementPolicy.class);
managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyClassName(policyClass);
managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(policyConfig.getProperties());
}

managedLedgerConfig.setThrottleMarkDelete(persistencePolicies.getManagedLedgerMaxMarkDeleteRate() >= 0
Expand Down
Loading