Skip to content

Commit e006cf8

Browse files
committed
[test][ml] Cover common-pool read completion boundaries and fallback
Document and log common-pool eligibility, callback affinity, cached-chain continuation, and startup property parsing. Pin the initialized depth and its supported radix forms. Verify a fully cached nested chain returns to the ledger worker across a ledger boundary, use managed blocking in the gated common-pool test, and exercise inline-mode rejected completion when common-pool parallelism selects the ledger fallback. Assisted-by: Codex
1 parent 20b8a96 commit e006cf8

8 files changed

Lines changed: 146 additions & 8 deletions

File tree

‎conf/broker.conf‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1515,8 +1515,11 @@ managedLedgerBatchReadEnabled=true
15151515
# False also restores the Exclusive/Failover cache-hit handoff used before PR #26619.
15161516
# The JVM-wide property -Dpulsar.managedLedger.maxReadCompletionDepth limits nested inline
15171517
# callbacks in both modes (default 10, values below 1 use 1). Set it at JVM startup.
1518+
# The depth accepts Integer.decode syntax, including hexadecimal and leading-zero octal.
15181519
# At the limit, enabled mode queues to the JVM common ForkJoinPool; disabled mode queues to
15191520
# the ledger executor. If common-pool parallelism is at most 1, both use the ledger executor.
1521+
# Common-pool parallelism normally uses available processors minus one (at least one).
1522+
# Override it with -Djava.util.concurrent.ForkJoinPool.common.parallelism.
15201523
# A limit of 1 queues every subsequent completion in a nested cached-read chain.
15211524
# This setting is not dynamic: restart the broker to apply a changed value. The policy is captured
15221525
# when a managed ledger opens and cannot change for already loaded topics. Failure callbacks,

‎conf/standalone.conf‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1010,8 +1010,11 @@ managedLedgerBatchReadEnabled=true
10101010
# False also restores the Exclusive/Failover cache-hit handoff used before PR #26619.
10111011
# The JVM-wide property -Dpulsar.managedLedger.maxReadCompletionDepth limits nested inline
10121012
# callbacks in both modes (default 10, values below 1 use 1). Set it at JVM startup.
1013+
# The depth accepts Integer.decode syntax, including hexadecimal and leading-zero octal.
10131014
# At the limit, enabled mode queues to the JVM common ForkJoinPool; disabled mode queues to
10141015
# the ledger executor. If common-pool parallelism is at most 1, both use the ledger executor.
1016+
# Common-pool parallelism normally uses available processors minus one (at least one).
1017+
# Override it with -Djava.util.concurrent.ForkJoinPool.common.parallelism.
10151018
# A limit of 1 queues every subsequent completion in a nested cached-read chain.
10161019
# This setting is not dynamic: restart the broker to apply a changed value. The policy is captured
10171020
# when a managed ledger opens and cannot change for already loaded topics. Failure callbacks,

‎managed-ledger/src/main/java/org/apache/bookkeeper/mledger/AsyncCallbacks.java‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -95,7 +95,9 @@ interface ReadEntriesCallback {
9595
* May be invoked inline when enabled in the ledger configuration, including on the calling thread for a cache
9696
* hit. The default ledger configuration restricts ordinary cursor read completions to the ledger executor,
9797
* with bounded inline completion when already on that executor.
98-
* The broker enables inline completion on other threads by default.
98+
* The broker enables inline completion on other threads by default. At the nesting limit, enabled mode
99+
* may continue on a JVM common-pool worker without Netty or ledger-executor thread affinity; it falls back
100+
* to the ledger executor when common-pool parallelism is at most one.
99101
* The recipient owns the returned entries and must release each entry after processing or discarding it,
100102
* including when its own shutdown or cancellation makes the result unnecessary.
101103
* Callers chaining cursor reads must coordinate result processing as described in {@link ManagedCursor}.

‎managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerConfig.java‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -427,7 +427,8 @@ public boolean isReadEntriesCallbackInline() {
427427
* queued handoff to the JVM common ForkJoinPool when enabled and to the ledger executor when disabled.
428428
* If common-pool parallelism is at most 1, enabled mode also uses the ledger executor. The JVM-wide system
429429
* property {@code pulsar.managedLedger.maxReadCompletionDepth} controls this limit (default 10, values below
430-
* 1 use 1).
430+
* 1 use 1). Values accept decimal, hexadecimal ({@code 0x10} or {@code #10}), and octal ({@code 010})
431+
* notation, following {@link Integer#decode(String)}.
431432
* A limit of 1 queues every subsequent read completion in a nested cached-read chain. Set the property at
432433
* JVM startup; later changes have no effect.
433434
*

‎managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ class OpReadEntry implements ReadEntriesCallback {
5252

5353
static {
5454
log.debug().attr("maxReadCompletionDepth", MAX_NESTED_INLINE_COMPLETIONS)
55+
.attr("useCommonPool", USE_COMMON_POOL)
5556
.log("Initialized managed-ledger read completion depth limit");
5657
}
5758

@@ -335,7 +336,9 @@ private void completeWithDepthLimit(Object ctx) {
335336
}
336337
} else {
337338
try {
338-
// Queue so the current callback stack can unwind. Legacy mode retains ledger-executor affinity.
339+
// Queue so the current callback stack can unwind. An inline cached-read chain can then continue
340+
// on common-pool workers until a cross-ledger read, cache miss, cursor wait, or caller handoff
341+
// changes its execution context. Legacy mode retains ledger-executor affinity.
339342
if (cursor.ledger.isReadEntriesCallbackInline() && USE_COMMON_POOL) {
340343
ForkJoinPool.commonPool().execute(() -> completeWithDepthLimit(ctx));
341344
} else {

‎managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OpReadEntryConfigTest.java‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,12 @@ public void testDefaultReadCompletionDepth() {
3030
assertThat(OpReadEntry.readMaxNestedInlineCompletions(new Properties())).isEqualTo(10);
3131
}
3232

33+
@Test
34+
public void testInitializedReadCompletionDepthMatchesSystemProperty() {
35+
assertThat(OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS)
36+
.isEqualTo(OpReadEntry.readMaxNestedInlineCompletions(System.getProperties()));
37+
}
38+
3339
@DataProvider
3440
public Object[][] readCompletionDepths() {
3541
return new Object[][] {
@@ -38,6 +44,8 @@ public Object[][] readCompletionDepths() {
3844
{"1", 1},
3945
{"37", 37},
4046
{"0x10", 16},
47+
{"#10", 16},
48+
{"010", 8},
4149
{"invalid", 10},
4250
{"", 10},
4351
{"2147483648", 10}

‎managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ReadCompletionAffinityTest.java‎

Lines changed: 120 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,7 @@
5858
import org.apache.bookkeeper.mledger.ScanOutcome;
5959
import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils;
6060
import org.apache.bookkeeper.test.MockedBookKeeperTestCase;
61+
import org.testng.annotations.DataProvider;
6162
import org.testng.annotations.Test;
6263

6364
public class ReadCompletionAffinityTest extends MockedBookKeeperTestCase {
@@ -205,7 +206,7 @@ public void readEntriesComplete(List<Entry> entries, Object ctx) {
205206
queuedPool.set(ForkJoinTask.getPool());
206207
queuedCompletionStarted.complete(Thread.currentThread());
207208
try {
208-
if (!releaseQueuedCompletion.await(20, TimeUnit.SECONDS)) {
209+
if (!awaitQueuedCompletion(releaseQueuedCompletion)) {
209210
queuedCompletion.completeExceptionally(
210211
new IllegalStateException("Timed out waiting to release queued completion"));
211212
return;
@@ -275,6 +276,88 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) {
275276
}
276277
}
277278

279+
@Test(timeOut = 30000)
280+
public void testNestedCachedChainReturnsToLedgerExecutorAcrossLedgers() throws Exception {
281+
int firstLedgerCount = OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 2;
282+
ManagedLedgerConfig config = inlineConfig().setMaxEntriesPerLedger(firstLedgerCount);
283+
config.setMinimumRolloverTime(0, TimeUnit.SECONDS);
284+
ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("completion-cached-chain-rollover", config);
285+
AtomicInteger storageReads = new AtomicInteger();
286+
try {
287+
ManagedCursor cursor = ledger.openCursor("cursor");
288+
List<Position> positions = new ArrayList<>();
289+
for (int i = 0; i <= firstLedgerCount; i++) {
290+
positions.add(ledger.addEntry(new byte[] {(byte) i}));
291+
}
292+
assertThat(positions.get(firstLedgerCount).getLedgerId())
293+
.isNotEqualTo(positions.get(0).getLedgerId());
294+
// Resolve the closed ledger handle before starting the chain so opening it cannot cause a handoff.
295+
ledger.getLedgerHandle(positions.get(0).getLedgerId()).get(10, TimeUnit.SECONDS);
296+
ledger.entryCache.clear();
297+
for (int i = 0; i < positions.size(); i++) {
298+
cacheEntry(ledger, positions.get(i), (byte) i);
299+
}
300+
bkc.setReadHandleInterceptor((ledgerId, first, last, entries) -> {
301+
storageReads.incrementAndGet();
302+
return CompletableFuture.completedFuture(entries);
303+
});
304+
CompletableFuture<Thread> ledgerWorker = new CompletableFuture<>();
305+
ledger.getExecutor().execute(() -> ledgerWorker.complete(Thread.currentThread()));
306+
Thread worker = ledgerWorker.get(10, TimeUnit.SECONDS);
307+
Thread caller = Thread.currentThread();
308+
309+
AtomicInteger completions = new AtomicInteger();
310+
AtomicReference<Thread> overflowThread = new AtomicReference<>();
311+
AtomicReference<ForkJoinPool> overflowPool = new AtomicReference<>();
312+
AtomicReference<Thread> boundaryThread = new AtomicReference<>();
313+
AtomicReference<ForkJoinPool> boundaryPool = new AtomicReference<>();
314+
CompletableFuture<List<Position>> crossedEntries = new CompletableFuture<>();
315+
cursor.asyncReadEntries(1, new ReadEntriesCallback() {
316+
@Override
317+
public void readEntriesComplete(List<Entry> entries, Object ctx) {
318+
List<Position> returnedPositions = entries.stream().map(Entry::getPosition).toList();
319+
entries.forEach(Entry::release);
320+
int completion = completions.incrementAndGet();
321+
if (completion <= OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS) {
322+
cursor.asyncReadEntries(1, this, null, PositionFactory.LATEST);
323+
} else if (completion == OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 1) {
324+
overflowThread.set(Thread.currentThread());
325+
overflowPool.set(ForkJoinTask.getPool());
326+
// One entry remains in the first ledger. Reading two forces checkReadCompletion to
327+
// schedule the continuation on the ledger executor, even though both entries are cached.
328+
cursor.asyncReadEntries(2, this, null, PositionFactory.LATEST);
329+
} else {
330+
boundaryThread.set(Thread.currentThread());
331+
boundaryPool.set(ForkJoinTask.getPool());
332+
crossedEntries.complete(returnedPositions);
333+
}
334+
}
335+
336+
@Override
337+
public void readEntriesFailed(ManagedLedgerException exception, Object ctx) {
338+
crossedEntries.completeExceptionally(exception);
339+
}
340+
}, null, PositionFactory.LATEST);
341+
342+
assertThat(crossedEntries.get(10, TimeUnit.SECONDS))
343+
.containsExactly(positions.get(firstLedgerCount - 1), positions.get(firstLedgerCount));
344+
assertThat(overflowThread.get()).isNotSameAs(caller);
345+
if (ForkJoinPool.getCommonPoolParallelism() > 1) {
346+
assertThat(overflowPool.get()).isSameAs(ForkJoinPool.commonPool());
347+
assertThat(overflowThread.get()).isNotSameAs(worker);
348+
} else {
349+
assertThat(overflowThread.get()).isSameAs(worker);
350+
}
351+
assertThat(boundaryThread.get()).isSameAs(worker);
352+
assertThat(boundaryPool.get()).isNull();
353+
assertThat(completions.get()).isEqualTo(OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 2);
354+
assertThat(storageReads.get()).isZero();
355+
} finally {
356+
bkc.setReadHandleInterceptor(null);
357+
ledger.close();
358+
}
359+
}
360+
278361
@Test
279362
public void testLegacyCachedReadQueuesCompletionOffLedgerAndRunsItInlineOnLedger() throws Exception {
280363
ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("completion-default-legacy", rawEntryConfig());
@@ -419,11 +502,18 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) {
419502
}
420503
}
421504

422-
@Test
423-
public void testRejectedLedgerExecutorAfterDepthLimitReleasesEntriesAndCompletesFailureOnce()
505+
@DataProvider
506+
public Object[][] rejectedDepthLimitCompletionModes() {
507+
// Run in a fresh JVM with common-pool parallelism at most one to also cover the inline fallback.
508+
return ForkJoinPool.getCommonPoolParallelism() > 1
509+
? new Object[][] {{false}} : new Object[][] {{false}, {true}};
510+
}
511+
512+
@Test(dataProvider = "rejectedDepthLimitCompletionModes")
513+
public void testRejectedLedgerExecutorAfterDepthLimitReleasesEntriesAndCompletesFailureOnce(boolean inline)
424514
throws Exception {
425-
ManagedLedgerImpl ledger = spy((ManagedLedgerImpl) factory.open("completion-depth-rejection",
426-
rawEntryConfig()));
515+
ManagedLedgerImpl ledger = spy((ManagedLedgerImpl) factory.open("completion-depth-rejection-" + inline,
516+
rawEntryConfig().setReadEntriesCallbackInline(inline)));
427517
try {
428518
ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("cursor");
429519
int count = OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 1;
@@ -509,6 +599,31 @@ public void testReadCompletionPolicyIsCapturedWhenLedgerOpens() throws Exception
509599
}
510600
}
511601

602+
private static boolean awaitQueuedCompletion(CountDownLatch latch) throws InterruptedException {
603+
if (!ForkJoinTask.inForkJoinPool()) {
604+
return latch.await(20, TimeUnit.SECONDS);
605+
}
606+
ForkJoinPool.ManagedBlocker blocker = new ForkJoinPool.ManagedBlocker() {
607+
private final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(20);
608+
609+
@Override
610+
public boolean block() throws InterruptedException {
611+
long remaining = deadline - System.nanoTime();
612+
if (remaining > 0) {
613+
latch.await(remaining, TimeUnit.NANOSECONDS);
614+
}
615+
return true;
616+
}
617+
618+
@Override
619+
public boolean isReleasable() {
620+
return latch.getCount() == 0 || System.nanoTime() >= deadline;
621+
}
622+
};
623+
ForkJoinPool.managedBlock(blocker);
624+
return latch.getCount() == 0;
625+
}
626+
512627
private static ByteBuf cacheEntry(ManagedLedgerImpl ledger, Position position, byte value) {
513628
ByteBuf data = Unpooled.wrappedBuffer(new byte[] {value});
514629
EntryImpl entry = EntryImpl.create(position, data, 0);

‎pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2740,8 +2740,11 @@ in the Admin API message inspection endpoints (getMessageById, peekNthMessage,
27402740
+ "False also restores the Exclusive/Failover cache-hit handoff used before PR #26619. "
27412741
+ "The JVM-wide property pulsar.managedLedger.maxReadCompletionDepth limits nested inline "
27422742
+ "callbacks in both modes (default 10, values below 1 use 1); set it at JVM startup. "
2743+
+ "The depth accepts Integer.decode syntax, including hexadecimal and leading-zero octal. "
27432744
+ "At the limit, enabled mode queues to the JVM common ForkJoinPool; disabled mode queues to "
27442745
+ "the ledger executor. If common-pool parallelism is at most 1, both use the ledger executor. "
2746+
+ "Common-pool parallelism normally uses available processors minus one (at least one); "
2747+
+ "override it with -Djava.util.concurrent.ForkJoinPool.common.parallelism. "
27452748
+ "A limit of 1 queues every subsequent completion in a nested cached-read chain. "
27462749
+ "This is not a dynamic setting: the completion policy is captured when a managed "
27472750
+ "ledger opens and does not change for already loaded topics. Failure callbacks, single-entry "

0 commit comments

Comments
 (0)