Skip to content

Pin a ledger's callbacks to a caller-chosen worker thread via withOrderingKey - #4881

Merged
merlimat merged 4 commits into
apache:masterfrom
merlimat:wt/kind-rhodes-ecbb71
Sep 10, 2026
Merged

merlimat merged 4 commits into
apache:masterfrom
merlimat:wt/kind-rhodes-ecbb71

Conversation

@merlimat

@merlimat merlimat commented Sep 9, 2026

Copy link
Copy Markdown
Contributor

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 ManagedLedgerImpl runs each managed ledger on bookKeeper.getMainWorkerPool().chooseThread(mlName), where mlName is a String. The BookKeeper client pins every LedgerHandle to chooseThread(ledgerId) and dispatches every bookie response with executor.executeOrdered(ledgerId, ...). Both threads come from the same OrderedExecutor, but they hash differently, so on the broker:

  • every add completion costs an extra cross-thread hop (ledger thread to managed-ledger thread, done by ml.getExecutor().execute(this) in OpAddEntry.addComplete), and
  • each of the E bookie responses per entry wakes a thread that is not the one doing that ledger's work.

With 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 with OrderedExecutor.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.asyncCreateLedger and from the ledger open call sites in ManagedLedgerImpl / ManagedCursorImpl) is a separate follow-up in apache/pulsar.

Changes

Public API

  • CreateBuilder.withOrderingKey(Object), CreateAdvBuilder.withOrderingKey(Object) and OpenBuilder.withOrderingKey(Object), implemented in LedgerCreateOp.CreateBuilderImpl / CreateAdvBuilderImpl and OpenBuilderBase. They are default no-op methods on the @Public interfaces, like withLoggerContext, so out-of-tree implementors keep compiling. null (the default) means "select the thread by ledger id".

Handle-side plumbing

  • LedgerHandle, LedgerHandleAdv and ReadOnlyLedgerHandle get a constructor overload taking the key; the existing constructors delegate with null. LedgerHandle.executor is mainWorkerPool.chooseThread(orderingKey) when a key is set and chooseThread(ledgerId) otherwise. The two overload calls are kept on purpose: chooseThread(Object) hashes through Long.hashCode, which folds the high bits, so boxing the ledger id would move ledgers with ids >= 2^31 to a different thread than today.
  • Every other place that picked a thread by ledger id for a handle now goes through the handle's executor: the four whenCompleteAsync / addCallback sites in LedgerHandle, the metadata updater in ReadOnlyLedgerHandle, ReadOpBase.submit() (so PendingReadOp / BatchedReadOp start on the handle thread), the speculative-request tasks of ReadOpBase and ReadLastConfirmedAndEntryOp (through a new LedgerHandle.submitOrdered(Callable) that mirrors OrderedExecutor.submitOrdered) and, in LedgerOpenOp, both the scheduler thread that runs openWithMetadata and the recovery completion callback. The default path keeps the ledger-id keyed OrderedGenericCallback verbatim; with a key the completion is submitted to the handle's executor (the body moved to recoveryComplete).
  • LedgerCreateOp / LedgerOpenOp carry the key from the builders to the handle constructors.

Response dispatch

  • BookieClient gets an Executor callbackExecutor overload for addEntry, readEntry, batchReadEntries, readLac, writeLac, forceLedger and readEntryWaitForLACUpdate. The previous signatures become default methods delegating with null, so the admin, replication, checker, benchmark and distributedlog callers are untouched.
  • BookieClientImpl, PerChannelBookieClient and the CompletionValue hierarchy (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, otherwise executor.executeOrdered(ledgerId, ...) exactly as before.
  • The client ops pass the handle's executor: PendingAddOp, PendingReadOp, BatchedReadOp, PendingReadLacOp, PendingWriteLacOp, ForceLedgerOp, TryReadLastConfirmedOp, ReadLastConfirmedAndEntryOp, and ReadLastConfirmedOp (new constructor parameter, used by LedgerHandle and LedgerRecoveryOp).
  • MockBookieClient implements the new signatures; the Mockito stubs in MockBookKeeperTestCase, BookKeeperBuildersOpenLedgerTest, PendingWriteLacOpTest, ReadLastConfirmedAndEntryOpTest and the direct PerChannelBookieClient call in TestPerChannelBookieClient are 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, so sendAddSuccessCallbacks and the rest of the unsynchronized handle state keep their single-writer assumption. PendingAddOp and PendingReadOp semantics are unchanged; they only forward the executor.

Decisions

  • Builders only, no legacy overloads. BookKeeper.asyncOpenLedger, asyncOpenLedgerNoRecovery and asyncCreateLedger are 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 uses newCreateLedgerOp(); its open call sites can move to newOpenLedgerOp() (which now also supports withKeepUpdateMetadata, Add OpenBuilder.withKeepUpdateMetadata and DeleteBuilder.withLoggerContext to the builder API #4834) and cast the returned ReadOnlyLedgerHandle, a LedgerHandle subclass.
  • Delete stays keyed by ledger id. LedgerDeleteOp has no handle and no callbacks that need to be ordered with a handle's callbacks, so the option is scoped to create and open.
  • Executor rather than key through the proto layer. BookieClient takes 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 the PerChannelBookieClient executor being the same pool as the client's main worker pool.
  • No ledger-id to executor registry. Plumbing through BookieClient is explicit, has no global state and no cleanup-on-close lifecycle; the interface grows by one nullable parameter per ledger operation.
  • Scheduler thread for openWithMetadata is keyed by the ordering key when one is set, for consistency; the client scheduler is a single thread, so this has no observable effect.
  • Left untouched: the periodic explicit-LAC timer in ExplicitLacFlushPolicy (scheduleAtFixedRateOrdered(ledgerId, ...) on that same single-thread scheduler) and the two unkeyed mainWorkerPool.submit(...) calls in ExplicitLacFlushPolicy and LedgerHandleAdv. None of them delivers a callback to the application, and re-routing them would change the default path.

Verification

  • New 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 checks asyncReadLastConfirmed; 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.
  • New BookieClientTest.testCallbackExecutorV2 / V3: at the BookieClientImpl level, add and read responses as well as a connection failure to an unreachable bookie complete on the supplied executor.
  • Existing tests run locally, all green: 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. (testSequenceReadLocalEnsemble in 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) and spotbugs:check on bookkeeper-server pass.

…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.
@merlimat
merlimat requested a review from lhotari September 9, 2026 22:26

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Great improvement! LGTM

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.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

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 callbackExecutor overloads across BookieClient request 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 callbackExecutor but doesn’t wait for termination. This can leave a non-daemon thread running and cause intermittent test hangs/thread-leak failures in CI. After shutdown(), also awaitTermination(...) (and optionally shutdownNow() 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 callbackExecutor but doesn’t wait for termination. This can leave a non-daemon thread running and cause intermittent test hangs/thread-leak failures in CI. After shutdown(), also awaitTermination(...) (and optionally shutdownNow() 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.

@merlimat
merlimat merged commit badd5f8 into apache:master Sep 10, 2026
20 checks passed
merlimat added a commit that referenced this pull request Sep 11, 2026
…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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants