Repository navigation
Pin a ledger's callbacks to a caller-chosen worker thread via withOrderingKey - #4881
Conversation
…thread Add CreateBuilder/CreateAdvBuilder/OpenBuilder.withOrderingKey(Object). When set, LedgerHandle.executor is mainWorkerPool.chooseThread(key) instead of chooseThread(ledgerId), and every place that used to pick a thread by ledger id for the handle (handle callbacks, read op submission, speculative requests, metadata updates, open and recovery completion) goes through that executor. Bookie response dispatch follows the same thread: BookieClient gains an Executor callbackExecutor overload for addEntry, readEntry, batchReadEntries, readLac, writeLac, forceLedger and readEntryWaitForLACUpdate, carried through BookieClientImpl, PerChannelBookieClient and the CompletionValue hierarchy so that responses (V2 and V3), connection failures, error-outs and timeouts all run on the handle's thread. The previous signatures remain as default methods. Without a key nothing changes: null resolves to the ledger-id thread exactly as before. This lets Pulsar pass its managed-ledger name as the key so the ledger thread and the managed-ledger thread coincide, removing a cross-thread hop per add completion and per bookie response.
…thread chooseThread(long) hashes the raw id while chooseThread(Object) goes through hashCode(), and Long.hashCode folds the high 32 bits. Resolving the ordering key to a single Object and making one chooseThread call would therefore move ledgers with ids >= 2^31 to a different worker thread than today and than the other ledger-id keyed dispatches, so the two overload calls stay.
The op now issues its speculative reads through the handle's executor (LedgerHandle.submitOrdered), which on the Mockito-mocked handle of this test returned null, so the speculative reads were never sent and the test waited forever for them. Route the stubbed method to the test's ordered scheduler, keyed by the ledger id as the op did before.
There was a problem hiding this comment.
Warning
Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.
Pull request overview
Introduces a caller-provided “ordering key” that pins all LedgerHandle callbacks to an application-chosen worker thread (instead of hashing by ledgerId), and plumbs an optional callbackExecutor through BookieClient so bookie responses/failures/timeout completions execute on the handle’s resolved executor.
Changes:
- Add
withOrderingKey(Object)to create/open builder APIs and carry it into handle construction to select the handle’s executor thread. - Add
Executor callbackExecutoroverloads acrossBookieClientrequest methods and dispatch all completions via the supplied executor when present. - Add/extend tests to validate callback-thread affinity across protocols and failure paths.
Reviewed changes
Copilot reviewed 38 out of 38 changed files in this pull request and generated 2 comments.
Show a summary per file
| File | Description |
|---|---|
| bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookieClientTest.java | Adds BookieClient-level tests asserting completions run on a supplied callback executor (V2/V3). |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/TestPerChannelBookieClient.java | Updates direct test invocation to new readEntry signature with callback executor. |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/MockBookieClient.java | Implements new BookieClient callback-executor signatures for mock behavior. |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/client/api/OrderingKeyTest.java | New integration test validating ordering-key thread pinning for create/open/recovery and default ledgerId behavior. |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/client/api/BookKeeperBuildersOpenLedgerTest.java | Updates Mockito stubs for new BookieClient method signatures. |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/client/ReadLastConfirmedAndEntryOpTest.java | Updates tests to reflect speculative reads now submit via handle executor. |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/client/PendingWriteLacOpTest.java | Updates mock expectations for new writeLac(..., Executor) signature. |
| bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java | Extends stubs to cover both legacy and new callback-executor overloads. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/WriteLacCompletion.java | Carries callback executor through completion value. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadLacCompletion.java | Carries callback executor through completion value. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadCompletion.java | Carries callback executor through completion value. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/PerChannelBookieClient.java | Plumbs callback executor through requests and dispatches responses via callback executor when provided. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ForceLedgerCompletion.java | Carries callback executor through completion value. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/CompletionValue.java | Centralizes ordered execution to respect callback executor when present. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieClientImpl.java | Adds callback-executor parameter plumbing and uses unified ordered-execution helper. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieClient.java | Adds default methods + new overloads with Executor callbackExecutor and documents semantics. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BatchedReadCompletion.java | Carries callback executor through completion value. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/AddCompletion.java | Carries callback executor through pooled completion object lifecycle. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/impl/OpenBuilderBase.java | Stores ordering key in open builder implementation. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/OpenBuilder.java | Adds withOrderingKey(Object) default no-op to public API. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/CreateBuilder.java | Adds withOrderingKey(Object) default no-op to public API. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/api/CreateAdvBuilder.java | Adds withOrderingKey(Object) default no-op to public API. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/TryReadLastConfirmedOp.java | Routes read callbacks via handle executor. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ReadOpBase.java | Submits ops/speculative tasks via handle executor rather than ledgerId-keyed pool. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ReadOnlyLedgerHandle.java | Adds orderingKey-aware constructor and uses handle executor for metadata updates. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ReadLastConfirmedOp.java | Adds callback executor parameter and forwards it to bookie reads. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ReadLastConfirmedAndEntryOp.java | Uses handle executor for speculative tasks and passes handle executor to long-poll reads. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingWriteLacOp.java | Passes handle executor to bookie writeLAC requests. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingReadOp.java | Passes handle executor to bookie read requests. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingReadLacOp.java | Passes handle executor to bookie readLAC requests. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/PendingAddOp.java | Passes handle executor to bookie addEntry requests. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerRecoveryOp.java | Ensures recovery LAC read callbacks execute on handle executor. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerOpenOp.java | Adds orderingKey plumbing and dispatches open/recovery completions respecting handle thread. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandleAdv.java | Adds orderingKey-aware constructor overload. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java | Adds orderingKey-aware constructor, uses selected executor for callbacks, and adds submitOrdered. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerCreateOp.java | Carries orderingKey from builders into handle construction. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/ForceLedgerOp.java | Passes handle executor to forceLedger requests. |
| bookkeeper-server/src/main/java/org/apache/bookkeeper/client/BatchedReadOp.java | Passes handle executor to batchRead requests. |
Suppressed comments (4)
bookkeeper-server/src/test/java/org/apache/bookkeeper/client/api/OrderingKeyTest.java:1
- This new test is written using JUnit Jupiter (
org.junit.jupiter.*) while the surrounding BookKeeper test suite in this PR uses JUnit 4 (org.junit.*). If the module isn’t configured to run/compile with JUnit 5, this will fail compilation or not execute. Align the test with the project’s test framework for this module (e.g., convert to JUnit 4 annotations/assertions) or ensure the module explicitly supports Jupiter for test compilation/execution.
bookkeeper-server/src/test/java/org/apache/bookkeeper/client/api/OrderingKeyTest.java:1 - This new test is written using JUnit Jupiter (
org.junit.jupiter.*) while the surrounding BookKeeper test suite in this PR uses JUnit 4 (org.junit.*). If the module isn’t configured to run/compile with JUnit 5, this will fail compilation or not execute. Align the test with the project’s test framework for this module (e.g., convert to JUnit 4 annotations/assertions) or ensure the module explicitly supports Jupiter for test compilation/execution.
bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookieClientTest.java:1 - The test shuts down
callbackExecutorbut doesn’t wait for termination. This can leave a non-daemon thread running and cause intermittent test hangs/thread-leak failures in CI. Aftershutdown(), alsoawaitTermination(...)(and optionallyshutdownNow()on timeout) to ensure the executor is fully stopped.
bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookieClientTest.java:1 - The test shuts down
callbackExecutorbut doesn’t wait for termination. This can leave a non-daemon thread running and cause intermittent test hangs/thread-leak failures in CI. Aftershutdown(), alsoawaitTermination(...)(and optionallyshutdownNow()on timeout) to ensure the executor is fully stopped.
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
…eringKey (#4881) * Client: option to pin a ledger's callbacks to a caller-chosen worker thread Add CreateBuilder/CreateAdvBuilder/OpenBuilder.withOrderingKey(Object). When set, LedgerHandle.executor is mainWorkerPool.chooseThread(key) instead of chooseThread(ledgerId), and every place that used to pick a thread by ledger id for the handle (handle callbacks, read op submission, speculative requests, metadata updates, open and recovery completion) goes through that executor. Bookie response dispatch follows the same thread: BookieClient gains an Executor callbackExecutor overload for addEntry, readEntry, batchReadEntries, readLac, writeLac, forceLedger and readEntryWaitForLACUpdate, carried through BookieClientImpl, PerChannelBookieClient and the CompletionValue hierarchy so that responses (V2 and V3), connection failures, error-outs and timeouts all run on the handle's thread. The previous signatures remain as default methods. Without a key nothing changes: null resolves to the ledger-id thread exactly as before. This lets Pulsar pass its managed-ledger name as the key so the ledger thread and the managed-ledger thread coincide, removing a cross-thread hop per add completion and per bookie response. * Comment why LedgerHandle never boxes the ledger id when choosing its thread chooseThread(long) hashes the raw id while chooseThread(Object) goes through hashCode(), and Long.hashCode folds the high 32 bits. Resolving the ordering key to a single Object and making one chooseThread call would therefore move ledgers with ids >= 2^31 to a different worker thread than today and than the other ledger-id keyed dispatches, so the two overload calls stay. * Stub LedgerHandle.submitOrdered in ReadLastConfirmedAndEntryOpTest The op now issues its speculative reads through the handle's executor (LedgerHandle.submitOrdered), which on the Mockito-mocked handle of this test returned null, so the speculative reads were never sent and the test waited forever for them. Route the stubbed method to the test's ordered scheduler, keyed by the ledger id as the op did before.
Descriptions of the changes in this PR:
Adds an optional ordering key to ledger creation and open, so that every callback of a ledger handle runs on a worker thread chosen by the application instead of the one hashed from the ledger id.
Motivation
Apache Pulsar's
ManagedLedgerImplruns each managed ledger onbookKeeper.getMainWorkerPool().chooseThread(mlName), wheremlNameis a String. The BookKeeper client pins everyLedgerHandletochooseThread(ledgerId)and dispatches every bookie response withexecutor.executeOrdered(ledgerId, ...). Both threads come from the sameOrderedExecutor, but they hash differently, so on the broker:ml.getExecutor().execute(this)inOpAddEntry.addComplete), andWith this change the application passes the same key it uses for its own
chooseThread(key)call and the ledger's callbacks are delivered on that very thread. The key is resolved withOrderedExecutor.chooseThread(Object)on the client's main worker pool. When no key is given nothing changes: the thread is still selected by ledger id, exactly as before.The Pulsar-side wiring (passing the managed ledger name from
ManagedLedgerImpl.asyncCreateLedgerand from the ledger open call sites inManagedLedgerImpl/ManagedCursorImpl) is a separate follow-up in apache/pulsar.Changes
Public API
CreateBuilder.withOrderingKey(Object),CreateAdvBuilder.withOrderingKey(Object)andOpenBuilder.withOrderingKey(Object), implemented inLedgerCreateOp.CreateBuilderImpl/CreateAdvBuilderImplandOpenBuilderBase. They aredefaultno-op methods on the@Publicinterfaces, likewithLoggerContext, so out-of-tree implementors keep compiling.null(the default) means "select the thread by ledger id".Handle-side plumbing
LedgerHandle,LedgerHandleAdvandReadOnlyLedgerHandleget a constructor overload taking the key; the existing constructors delegate withnull.LedgerHandle.executorismainWorkerPool.chooseThread(orderingKey)when a key is set andchooseThread(ledgerId)otherwise. The two overload calls are kept on purpose:chooseThread(Object)hashes throughLong.hashCode, which folds the high bits, so boxing the ledger id would move ledgers with ids >= 2^31 to a different thread than today.whenCompleteAsync/addCallbacksites inLedgerHandle, the metadata updater inReadOnlyLedgerHandle,ReadOpBase.submit()(soPendingReadOp/BatchedReadOpstart on the handle thread), the speculative-request tasks ofReadOpBaseandReadLastConfirmedAndEntryOp(through a newLedgerHandle.submitOrdered(Callable)that mirrorsOrderedExecutor.submitOrdered) and, inLedgerOpenOp, both the scheduler thread that runsopenWithMetadataand the recovery completion callback. The default path keeps the ledger-id keyedOrderedGenericCallbackverbatim; with a key the completion is submitted to the handle's executor (the body moved torecoveryComplete).LedgerCreateOp/LedgerOpenOpcarry the key from the builders to the handle constructors.Response dispatch
BookieClientgets anExecutor callbackExecutoroverload foraddEntry,readEntry,batchReadEntries,readLac,writeLac,forceLedgerandreadEntryWaitForLACUpdate. The previous signatures becomedefaultmethods delegating withnull, so the admin, replication, checker, benchmark and distributedlog callers are untouched.BookieClientImpl,PerChannelBookieClientand theCompletionValuehierarchy (AddCompletion,ReadCompletion,BatchedReadCompletion,ReadLacCompletion,WriteLacCompletion,ForceLedgerCompletion) carry the executor with the request. Responses on both wire protocols (readV2Response/readV3Response), connection failures, error-outs and timeouts all dispatch through one helper: the caller's executor when present, otherwiseexecutor.executeOrdered(ledgerId, ...)exactly as before.PendingAddOp,PendingReadOp,BatchedReadOp,PendingReadLacOp,PendingWriteLacOp,ForceLedgerOp,TryReadLastConfirmedOp,ReadLastConfirmedAndEntryOp, andReadLastConfirmedOp(new constructor parameter, used byLedgerHandleandLedgerRecoveryOp).MockBookieClientimplements the new signatures; the Mockito stubs inMockBookKeeperTestCase,BookKeeperBuildersOpenLedgerTest,PendingWriteLacOpTest,ReadLastConfirmedAndEntryOpTestand the directPerChannelBookieClientcall inTestPerChannelBookieClientare extended to the new signatures.Ordering guarantee
All callbacks of a handle, on both wire protocols and on every failure path, are submitted to the single thread behind
LedgerHandle.executor, sosendAddSuccessCallbacksand the rest of the unsynchronized handle state keep their single-writer assumption.PendingAddOpandPendingReadOpsemantics are unchanged; they only forward the executor.Decisions
BookKeeper.asyncOpenLedger,asyncOpenLedgerNoRecoveryandasyncCreateLedgerare not extended. The builder API is the designated extension point for optional parameters and the legacy API already carries several positional overloads per operation. Pulsar's create path already usesnewCreateLedgerOp(); its open call sites can move tonewOpenLedgerOp()(which now also supportswithKeepUpdateMetadata, Add OpenBuilder.withKeepUpdateMetadata and DeleteBuilder.withLoggerContext to the builder API #4834) and cast the returnedReadOnlyLedgerHandle, aLedgerHandlesubclass.LedgerDeleteOphas no handle and no callbacks that need to be ordered with a handle's callbacks, so the option is scoped to create and open.BookieClienttakes the handle's resolved executor (its single worker thread) instead of the raw key: the proto layer never re-hashes a key per response, and correctness does not depend on thePerChannelBookieClientexecutor being the same pool as the client's main worker pool.BookieClientis explicit, has no global state and no cleanup-on-close lifecycle; the interface grows by one nullable parameter per ledger operation.openWithMetadatais keyed by the ordering key when one is set, for consistency; the client scheduler is a single thread, so this has no observable effect.ExplicitLacFlushPolicy(scheduleAtFixedRateOrdered(ledgerId, ...)on that same single-thread scheduler) and the two unkeyedmainWorkerPool.submit(...)calls inExplicitLacFlushPolicyandLedgerHandleAdv. None of them delivers a callback to the application, and re-routing them would change the default path.Verification
OrderingKeyTest(3-bookie cluster, run on both the V2 and V3 wire protocols): creates a ledger with an explicit id and a key that maps to a different worker thread than the id, and checks that add and read callbacks run on the key's thread; opens the closed ledger under another key and checks reads; opens an unclosed ledger without recovery and checksasyncReadLastConfirmed; opens an unclosed ledger with recovery and checks the reads that follow; and, without a key, checks that add and read callbacks still run on the thread selected by ledger id.BookieClientTest.testCallbackExecutorV2/V3: at theBookieClientImpllevel, add and read responses as well as a connection failure to an unreachable bookie complete on the supplied executor.BookieClientTest,TestPerChannelBookieClient,PendingAddOpTest,PendingWriteLacOpTest,ReadLastConfirmedAndEntryOpTest,ReadLastConfirmedOpTest,TestPendingReadLacOp,LoggerContextTest,BookKeeperApiTest,BookKeeperBuildersTest,BookKeeperBuildersOpenLedgerTest,DeferredSyncTest,TestMaxEnsembleChangeNum,HandleFailuresTest,LedgerClose2Test,LedgerRecovery2Test,MockBookKeeperTest,TestLedgerFragmentReplicationWithMock,DataIntegrityCheckTest,EntryCopierTest,BookieWriteLedgerTest(V2 and V3, 96 tests),BookieReadWriteTest,BookKeeperTest,BookKeeperCloseTest,LedgerCloseTest,TestFencing,ExplicitLacTest,TestBatchedRead,TestSpeculativeRead,TestSpeculativeBatchRead,ParallelLedgerRecoveryTest,LedgerRecoveryTest,BookieRecoveryTest,ConcurrentV2RecoveryTest,TestReadLastConfirmedAndEntry,TestReadLastConfirmedLongPoll,TestReadLastEntry,TestLedgerFragmentReplication,TestLedgerChecker. (testSequenceReadLocalEnsemblein the two speculative-read classes fails on the development machine on any branch because its hostname does not resolve; unrelated to this change.)checkstyle:check(main and test sources) andspotbugs:checkonbookkeeper-serverpass.