[fix][broker] Prevent BookkeeperSchemaStorage, BookkeeperBucketSnapshotStorage, and compactors from bypassing namespace bookie affinity - #26705
Conversation
Resolve topic placement configuration in one shared path and pair it with the cached BookKeeper client and matching ledger metadata. Assisted-by: Codex
…gers 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
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
| store, brokerConfig.getMetadataStoreOperationTimeoutSeconds()); | ||
| Optional<LocalPolicies> localPolicies = | ||
| localPoliciesResources.getLocalPolicies(topicName.getNamespaceObject()); | ||
| BookKeeperClientContext bookKeeperClientContext = createBookKeeperClientContext( |
There was a problem hiding this comment.
This looks good for selecting the right placement-aware client, but I think this PR still depends on #26704 for the guarantee it is trying to provide. On current master, IsolatedBookieEnsemblePlacementPolicy finishes client initialization before the initial /bookies mapping load completes, and getExcludedBookiesWithIsolationGroups() returns an empty exclusion set while cachedRackConfiguration is still null. In particular, CompactorTool creates a fresh isolated client and can immediately create the compacted ledger in that window.
Should we make #26704 an explicit dependency (or include the initialization fix here)? Otherwise #26705 alone can still place a newly created auxiliary ledger outside the configured affinity group.
There was a problem hiding this comment.
Agreed. #26704 (or an equivalent initialization fix) is required for the first-ledger isolation guarantee. This PR can be reviewed independently, but I consider that fix a merge dependency: otherwise a newly created client can still place an auxiliary ledger before the initial /bookies mapping is available.
| return pulsar.getBookKeeperClientContext(topicName.get()) | ||
| .thenCompose(clientContext -> createLedger(id, clientContext)); |
There was a problem hiding this comment.
One thing I'm still worried about is the schema recovery path. We now create the schema ledger with the placement-aware client, but BookkeeperSchemaStorage#openLedger() still always uses the storage's default bookKeeper. BookKeeper#asyncOpenLedger() performs fencing/recovery when the ledger is still open, and ensemble replacement uses the placement policy of the client doing that recovery; the custom metadata alone does not switch a Rackaware client to IsolatedBookieEnsemblePlacementPolicy.
This seems reachable because addNewSchemaEntryToStore() calls ledgerHandle.closeAsync() without waiting for it before the new position can be published in the schema locator. If the broker dies in that window, the next read can recover this new placement-aware ledger through the default client and move an ensemble member outside the namespace affinity group.
Would it make sense to await closeAsync() before returning the position / publishing the locator? That looks simpler than making every schema open path resolve the owner placement context, and it would also make sure a locator only points to a fully closed ledger.
There was a problem hiding this comment.
Fixed in 6498854. addNewSchemaEntryToStore now composes closeAsync() before returning the position, so the locator is published only after the ledger is closed. If close fails, it does not publish the locator and makes a best-effort attempt to delete the orphan ledger. I added a deterministic test for both close outcomes.
| .orElse(null); | ||
| ManagedLedgerStorageClass defaultStorageClass = | ||
| getManagedLedgerStorage().getDefaultStorageClass(); | ||
| if (!(defaultStorageClass instanceof BookkeeperManagedLedgerStorageClass bkStorageClass)) { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| BookKeeperPlacementPolicyConfigResolver.resolve(config, topicName, localPolicies) | ||
| .orElse(null); | ||
| ManagedLedgerStorageClass defaultStorageClass = | ||
| getManagedLedgerStorage().getDefaultStorageClass(); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
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
Motivation
Managed-ledger topic data resolves an effective BookKeeper placement policy from the owner namespace's
bookieAffinityGroupandstrictBookieAffinityEnabledsettings. It then uses a BookKeeper client initialized with that policy and stores the sameEnsemblePlacementPolicyConfigin each ledger's custom metadata.Several Pulsar-owned auxiliary-ledger creation paths did not follow that flow:
BookkeeperSchemaStoragecreated schema ledgers through its default BookKeeper client.BookkeeperBucketSnapshotStoragecreated delayed-delivery bucket ledgers through its default BookKeeper client.PublishingOrderCompactor,EventTimeOrderCompactor,StrategicTwoPhaseCompactor, and the standaloneCompactorToolcreated compacted ledgers through a fixed default client.This matters when the owner namespace has a bookie-affinity group, or when strict bookie affinity selects
IsolatedBookieEnsemblePlacementPolicy. In those cases, the auxiliary ledger could be created outside the bookie groups selected for the topic even though the namespace policy was configured successfully. The ledger also lacked placement-policy metadata, so a later BookKeeper ensemble replacement could not reliably preserve the intended placement.The failure was silent: topic data ledgers honored the policy, while schema, delayed bucket, and compaction ledgers for the same topic could bypass it.
Modifications
Shared topic-aware BookKeeper context
BrokerServiceaffinity, strict-affinity, and system-topic decision matrix intoBookKeeperPlacementPolicyConfigResolverand continue using the same resolver for managed-ledger configuration.BookkeeperManagedLedgerStorageClass.ManagedLedgerClientFactoryreuse its existingEnsemblePlacementPolicyConfigclient cache for auxiliary ledgers.BookKeeperClientContext, which keeps the borrowed BookKeeper client and the matching encoded placement metadata together, preventing the client configuration and ledger metadata from diverging.PulsarService#getBookKeeperClientContext(TopicName)to resolve the owner namespace's local policy without requiring a loaded topic instance.Schema and delayed-delivery bucket ledgers
BookkeeperSchemaStorageorBookkeeperBucketSnapshotStoragecreates a ledger.BucketSnapshotPersistenceException.Schema storage does not gain a separate placement-policy setting. A canonical schema ledger follows the policy of its owner topic's namespace.
Compaction ledgers
CompactorToolsynchronously read the target namespace's local policy and create exactly one owned BookKeeper client configured with the matching policy.Behavior and impact
ManagedLedgerClientFactory, one per distinct placement-policy configuration; auxiliary components only borrow them.CompactorToolstill owns exactly one BookKeeper client for its single target topic.Commit structure
314f70ccaab— introduce the shared resolver, cached client lookup, andBookKeeperClientContext.a7b2ccb439f— apply the context to schema and delayed-delivery bucket ledger creation.1fa761e790b— apply the context to publishing-order, event-time, strategic, and CLI compaction ledger creation.Tests
./gradlew :pulsar-broker:test --tests org.apache.pulsar.broker.ManagedLedgerClientFactoryTest --tests org.apache.pulsar.broker.storage.BookKeeperPlacementPolicyConfigResolverTest --tests org.apache.pulsar.broker.storage.BookKeeperClientContextTest./gradlew :pulsar-broker:test --no-build-cache --tests org.apache.pulsar.broker.service.schema.BookkeeperSchemaStorageTest --tests org.apache.pulsar.broker.service.schema.SchemaServiceTest --tests org.apache.pulsar.broker.delayed.BookkeeperBucketSnapshotStorageTest --tests org.apache.pulsar.broker.delayed.bucket.BookkeeperBucketSnapshotStoragePlacementTest./gradlew :pulsar-broker:test --no-build-cache --tests org.apache.pulsar.compaction.CompactorPlacementPolicyTest --tests org.apache.pulsar.compaction.CompactorToolTest --tests org.apache.pulsar.compaction.CompactorToolProviderPinsTest./gradlew :pulsar-broker:test --tests org.apache.pulsar.compaction.CompactorTest.testCompaction./gradlew :pulsar-broker:clean :pulsar-broker:compileJava --no-build-cache --no-configuration-cache./gradlew :pulsar-broker:checkstyleMain :pulsar-broker:checkstyleTestAssisted-by: Codex