This is an automated email from the ASF dual-hosted git repository. merlimat pushed a commit to branch branch-4.18 in repository https://gitbox.apache.org/repos/asf/bookkeeper.git
commit c9c25f2a89501705108ec5d3322cc2bdb5bfec97 Author: Matteo Merli <[email protected]> AuthorDate: Fri Sep 11 09:46:59 2026 -0700 Add SingleThreadExecutor.executeOrRun and issue reads inline from the ledger's own thread (#4883) * Run reads inline when issued from the ledger's own thread Add SingleThreadExecutor.isCurrentThread() and executeOrRun(Runnable): the task runs inline when called from the executor's own thread and is queued otherwise. The shortcut is opt-in because it bypasses the queue: a task submitted this way from the executor thread runs before the tasks already queued, nested inside the submitting task. execute() keeps its FIFO contract. Inline failures are logged and counted like those of queued tasks. LedgerHandle uses it for reads on a writable handle, which were queued onto the ledger's thread even when the caller was already on it. Adds already initiate inline on the caller's thread and reads on read-only handles already bypass the executor, so this was the remaining same-thread hop. Applications that run their own work on a thread from OrderedExecutor.chooseThread(...) can use the same check for their re-dispatches. * Type LedgerHandle.executor as SingleThreadExecutor and share the safe-run path Review follow-ups: the inline path of SingleThreadExecutor.executeOrRun calls safeRunTask directly, so inline and queued tasks share the same failure isolation, logging and counters; the queued path only adds the pending-count decrement, which must not apply to inline runs. LedgerHandle keeps its worker thread as a SingleThreadExecutor instead of checking with instanceof: the client's main worker pool is a plain OrderedExecutor, whose threads are SingleThreadExecutor instances. The two test fixtures that handed the client an OrderedScheduler as its worker pool (MockClientContext, MockBookKeeperTestCase) now use an OrderedExecutor, with the mock bookie client dispatching on that same pool. * Fix checkstyle line length in MockClientContext and reject executeOrRun after shutdown The scheduler line in MockClientContext.create was one character over the limit; the PR validation job runs checkstyle over all modules and failed on it. executeOrRun now rejects the task once the executor is shut down, like execute() does, instead of still running it inline when called from the executor thread. * Expose the thread identity through the OrderedExecutor decorators OrderedExecutor wraps each worker thread in a forwarding decorator when task tracing or MDC preservation is enabled, so chooseThread() returned the wrapper rather than the SingleThreadExecutor and the cast in the LedgerHandle constructor failed for such clients. The auditor and the replication worker run with task tracing enabled in the tests, so every ledger open there failed with UnexpectedConditionException and re-replication never completed. Add the ThreadBoundExecutor interface (isCurrentThread, executeOrRun), implemented by SingleThreadExecutor and by the decorator, which applies the timing and MDC wrappers to inline runs as well. LedgerHandle keeps its executor as a ThreadBoundExecutor. --- .../bookkeeper/common/util/OrderedExecutor.java | 111 ++++++++++------ .../common/util/SingleThreadExecutor.java | 46 ++++++- .../common/util/ThreadBoundExecutor.java | 46 +++++++ .../common/util/TestOrderedExecutorDecorators.java | 26 ++++ .../common/util/TestSingleThreadExecutor.java | 105 +++++++++++++++ .../org/apache/bookkeeper/client/LedgerHandle.java | 20 +-- .../client/LedgerHandleInlineReadTest.java | 148 +++++++++++++++++++++ .../bookkeeper/client/MockBookKeeperTestCase.java | 2 +- .../bookkeeper/client/MockClientContext.java | 8 +- 9 files changed, 450 insertions(+), 62 deletions(-) diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java index f32d25c3ce..699c1b1395 100644 --- a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/OrderedExecutor.java @@ -306,58 +306,81 @@ public class OrderedExecutor implements ExecutorService { } protected ExecutorService addExecutorDecorators(ExecutorService executor) { - return new ForwardingExecutorService() { - @Override - protected ExecutorService delegate() { - return executor; - } + checkArgument(executor instanceof ThreadBoundExecutor, "Expected a ThreadBoundExecutor, got %s", executor); + return new DecoratedThread((ThreadBoundExecutor) executor); + } - @Override - public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) - throws InterruptedException { - return super.invokeAll(timedCallables(tasks)); - } + /** + * One of the pool's threads with the task tracing and MDC decorators applied, which keeps exposing the thread + * identity of the underlying executor. + */ + private class DecoratedThread extends ForwardingExecutorService implements ThreadBoundExecutor { + private final ThreadBoundExecutor executor; - @Override - public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, - long timeout, TimeUnit unit) - throws InterruptedException { - return super.invokeAll(timedCallables(tasks), timeout, unit); - } + DecoratedThread(ThreadBoundExecutor executor) { + this.executor = executor; + } - @Override - public <T> T invokeAny(Collection<? extends Callable<T>> tasks) - throws InterruptedException, ExecutionException { - return super.invokeAny(timedCallables(tasks)); - } + @Override + protected ExecutorService delegate() { + return executor; + } - @Override - public <T> T invokeAny(Collection<? extends Callable<T>> tasks, - long timeout, TimeUnit unit) - throws InterruptedException, ExecutionException, TimeoutException { - return super.invokeAny(timedCallables(tasks), timeout, unit); - } + @Override + public boolean isCurrentThread() { + return executor.isCurrentThread(); + } - @Override - public void execute(Runnable command) { - super.execute(timedRunnable(command)); - } + @Override + public void executeOrRun(Runnable r) { + executor.executeOrRun(timedRunnable(r)); + } - @Override - public <T> Future<T> submit(Callable<T> task) { - return super.submit(timedCallable(task)); - } + @Override + public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) + throws InterruptedException { + return super.invokeAll(timedCallables(tasks)); + } - @Override - public Future<?> submit(Runnable task) { - return super.submit(timedRunnable(task)); - } + @Override + public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, + long timeout, TimeUnit unit) + throws InterruptedException { + return super.invokeAll(timedCallables(tasks), timeout, unit); + } - @Override - public <T> Future<T> submit(Runnable task, T result) { - return super.submit(timedRunnable(task), result); - } - }; + @Override + public <T> T invokeAny(Collection<? extends Callable<T>> tasks) + throws InterruptedException, ExecutionException { + return super.invokeAny(timedCallables(tasks)); + } + + @Override + public <T> T invokeAny(Collection<? extends Callable<T>> tasks, + long timeout, TimeUnit unit) + throws InterruptedException, ExecutionException, TimeoutException { + return super.invokeAny(timedCallables(tasks), timeout, unit); + } + + @Override + public void execute(Runnable command) { + super.execute(timedRunnable(command)); + } + + @Override + public <T> Future<T> submit(Callable<T> task) { + return super.submit(timedCallable(task)); + } + + @Override + public Future<?> submit(Runnable task) { + return super.submit(timedRunnable(task)); + } + + @Override + public <T> Future<T> submit(Runnable task, T result) { + return super.submit(timedRunnable(task), result); + } } /** diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java index 840a4c6dec..e26a3da5bc 100644 --- a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/SingleThreadExecutor.java @@ -24,7 +24,6 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.AbstractExecutorService; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; @@ -45,7 +44,7 @@ import org.apache.bookkeeper.stats.StatsLogger; * proceed with the next tasks. */ @CustomLog -public class SingleThreadExecutor extends AbstractExecutorService implements ExecutorService, Runnable { +public class SingleThreadExecutor extends AbstractExecutorService implements ThreadBoundExecutor, Runnable { private static final int MAX_DRAIN_BATCH_SIZE = 1024; @@ -120,7 +119,7 @@ public class SingleThreadExecutor extends AbstractExecutorService implements Exe for (int i = 0; i < n; i++) { Runnable task = localTasks[i]; localTasks[i] = null; - if (!safeRunTask(task)) { + if (!runQueuedTask(task)) { return; } } @@ -129,7 +128,7 @@ public class SingleThreadExecutor extends AbstractExecutorService implements Exe // Clear the queue in orderly shutdown Runnable task; while ((task = queue.poll()) != null) { - safeRunTask(task); + runQueuedTask(task); } } catch (InterruptedException ie) { // Exit loop when interrupted @@ -142,6 +141,19 @@ public class SingleThreadExecutor extends AbstractExecutorService implements Exe } } + private boolean runQueuedTask(Runnable r) { + try { + return safeRunTask(r); + } finally { + decrementPendingTaskCount(1); + } + } + + /** + * Runs a task, logging and counting a failure instead of propagating it. + * + * @return false when the task was interrupted + */ private boolean safeRunTask(Runnable r) { try { r.run(); @@ -154,8 +166,6 @@ public class SingleThreadExecutor extends AbstractExecutorService implements Exe tasksFailed.increment(); log.error().exception(t).log("Error while running task"); } - } finally { - decrementPendingTaskCount(1); } return true; @@ -220,6 +230,30 @@ public class SingleThreadExecutor extends AbstractExecutorService implements Exe executeRunnableOrList(r, null); } + @Override + public boolean isCurrentThread() { + return Thread.currentThread() == runner; + } + + /** + * {@inheritDoc} + * + * <p>Failures of an inline run are logged and counted like those of queued tasks. + */ + @Override + public void executeOrRun(Runnable r) { + if (state != State.Running) { + throw new RejectedExecutionException("Executor is shutting down"); + } + + if (isCurrentThread()) { + tasksCount.increment(); + safeRunTask(r); + } else { + execute(r); + } + } + @VisibleForTesting void executeRunnableOrList(Runnable runnable, List<Runnable> runnableList) { if (state != State.Running) { diff --git a/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/ThreadBoundExecutor.java b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/ThreadBoundExecutor.java new file mode 100644 index 0000000000..7d9098a3a2 --- /dev/null +++ b/bookkeeper-common/src/main/java/org/apache/bookkeeper/common/util/ThreadBoundExecutor.java @@ -0,0 +1,46 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.common.util; + +import java.util.concurrent.ExecutorService; + +/** + * An {@link ExecutorService} backed by a single thread, which can tell whether the caller is already on that + * thread and then run a task inline instead of queueing it. + * + * <p>Implemented by {@link SingleThreadExecutor} and by the threads an {@link OrderedExecutor} hands out from + * {@link OrderedExecutor#chooseThread(long)}, which are decorated when task tracing or MDC preservation is enabled. + */ +public interface ThreadBoundExecutor extends ExecutorService { + + /** + * Whether the calling thread is the thread of this executor. + */ + boolean isCurrentThread(); + + /** + * Runs the task inline when called from this executor's own thread, otherwise submits it like + * {@link #execute(Runnable)}. Like {@code execute}, it rejects the task once the executor is shut down. + * + * <p>The inline run bypasses the queue: a task submitted this way from the executor thread runs before the + * tasks already queued, nested inside the task that submitted it. Use it only where that reordering is + * acceptable. + */ + void executeOrRun(Runnable r); +} diff --git a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java index 69b0570899..475215313d 100644 --- a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java +++ b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestOrderedExecutorDecorators.java @@ -22,7 +22,9 @@ package org.apache.bookkeeper.common.util; import static org.hamcrest.Matchers.hasItem; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertTrue; import static org.mockito.AdditionalAnswers.answerVoid; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doAnswer; @@ -30,6 +32,7 @@ import static org.mockito.Mockito.spy; import java.util.Queue; import java.util.UUID; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; import lombok.CustomLog; @@ -82,6 +85,29 @@ public class TestOrderedExecutorDecorators { ThreadContext.clearMap(); } + @Test + public void testDecoratedThreadsAreThreadBound() throws Exception { + for (OrderedExecutor executor : new OrderedExecutor[] { + OrderedExecutor.newBuilder().name("traced").numThreads(2).traceTaskExecution(true).build(), + OrderedExecutor.newBuilder().name("mdc").numThreads(2).preserveMdcForTaskExecution(true).build()}) { + try { + ThreadBoundExecutor thread = (ThreadBoundExecutor) executor.chooseThread(10); + assertFalse(thread.isCurrentThread()); + + // From the thread itself, executeOrRun runs the task before returning + CompletableFuture<Boolean> ranInline = new CompletableFuture<>(); + thread.execute(() -> { + boolean[] ran = new boolean[1]; + thread.executeOrRun(() -> ran[0] = thread.isCurrentThread()); + ranInline.complete(ran[0]); + }); + assertTrue(ranInline.get(10, TimeUnit.SECONDS)); + } finally { + executor.shutdown(); + } + } + } + @Test public void testMDCInvokeOrdered() throws Exception { OrderedExecutor executor = OrderedExecutor.newBuilder() diff --git a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java index ed72704e97..0f6cdadf76 100644 --- a/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java +++ b/bookkeeper-common/src/test/java/org/apache/bookkeeper/common/util/TestSingleThreadExecutor.java @@ -20,6 +20,8 @@ package org.apache.bookkeeper.common.util; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotSame; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -30,6 +32,8 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.Callable; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.Future; @@ -38,6 +42,7 @@ import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import lombok.Cleanup; import org.awaitility.Awaitility; import org.junit.Test; @@ -369,4 +374,104 @@ public class TestSingleThreadExecutor { future.get(); } + + @Test + public void testIsCurrentThread() throws Exception { + @Cleanup("shutdown") + SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY); + + assertFalse(ste.isCurrentThread()); + Callable<Boolean> onExecutorThread = ste::isCurrentThread; + assertTrue(ste.submit(onExecutorThread).get(10, TimeUnit.SECONDS)); + } + + @Test + public void testExecuteOrRunInlineOnOwnThread() throws Exception { + @Cleanup("shutdown") + SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY); + + // From the executor thread the task runs before executeOrRun returns, ahead of the queued task + List<String> events = Collections.synchronizedList(new ArrayList<>()); + CompletableFuture<List<String>> seenBeforeReturn = new CompletableFuture<>(); + ste.execute(() -> { + ste.execute(() -> events.add("queued")); + ste.executeOrRun(() -> events.add("inline")); + seenBeforeReturn.complete(new ArrayList<>(events)); + }); + + assertEquals(Lists.newArrayList("inline"), seenBeforeReturn.get(10, TimeUnit.SECONDS)); + Awaitility.await().until(() -> events.size() == 2); + assertEquals(Lists.newArrayList("inline", "queued"), events); + Awaitility.await().until(() -> ste.getCompletedTasksCount() == 3); + assertEquals(3, ste.getSubmittedTasksCount()); + assertEquals(0, ste.getQueuedTasksCount()); + } + + @Test + public void testExecuteOrRunFromOtherThreadIsQueued() throws Exception { + @Cleanup("shutdown") + SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY); + + // Hold the executor thread so that a queued task cannot run yet + CountDownLatch release = new CountDownLatch(1); + ste.execute(() -> { + try { + release.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + + AtomicReference<Thread> ranOn = new AtomicReference<>(); + ste.executeOrRun(() -> ranOn.set(Thread.currentThread())); + assertNull(ranOn.get()); + + release.countDown(); + Awaitility.await().until(() -> ranOn.get() != null); + assertNotSame(Thread.currentThread(), ranOn.get()); + Callable<Boolean> ranOnExecutorThread = () -> ranOn.get() == Thread.currentThread(); + assertTrue(ste.submit(ranOnExecutorThread).get(10, TimeUnit.SECONDS)); + } + + @Test + public void testExecuteOrRunRejectedAfterShutdown() throws Exception { + SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY); + + // Shut down from the executor thread itself: the inline path must reject like execute() does + CompletableFuture<Boolean> rejectedInline = new CompletableFuture<>(); + ste.execute(() -> { + ste.shutdown(); + try { + ste.executeOrRun(() -> { + }); + rejectedInline.complete(false); + } catch (RejectedExecutionException e) { + rejectedInline.complete(true); + } + }); + assertTrue(rejectedInline.get(10, TimeUnit.SECONDS)); + + ste.awaitTermination(10, TimeUnit.SECONDS); + assertThrows(RejectedExecutionException.class, () -> ste.executeOrRun(() -> { + })); + } + + @Test + public void testExecuteOrRunInlineFailureIsIsolated() throws Exception { + @Cleanup("shutdown") + SingleThreadExecutor ste = new SingleThreadExecutor(THREAD_FACTORY); + + CompletableFuture<Boolean> submitterCompleted = new CompletableFuture<>(); + ste.execute(() -> { + ste.executeOrRun(() -> { + throw new RuntimeException("test"); + }); + submitterCompleted.complete(true); + }); + + assertTrue(submitterCompleted.get(10, TimeUnit.SECONDS)); + Awaitility.await().until(() -> ste.getCompletedTasksCount() == 1); + assertEquals(1, ste.getFailedTasksCount()); + assertEquals(2, ste.getSubmittedTasksCount()); + } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java index a4699a15b7..71116dd3d1 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/client/LedgerHandle.java @@ -51,7 +51,6 @@ import java.util.Set; import java.util.concurrent.Callable; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.ExecutorService; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -82,6 +81,7 @@ import org.apache.bookkeeper.client.impl.LedgerEntryImpl; import org.apache.bookkeeper.common.concurrent.FutureEventListener; import org.apache.bookkeeper.common.concurrent.FutureUtils; import org.apache.bookkeeper.common.util.MathUtils; +import org.apache.bookkeeper.common.util.ThreadBoundExecutor; import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.proto.BookieProtocol; import org.apache.bookkeeper.proto.checksum.DigestManager; @@ -106,7 +106,7 @@ public class LedgerHandle implements WriteHandle { final byte[] ledgerKey; private Versioned<LedgerMetadata> versionedMetadata; final long ledgerId; - final ExecutorService executor; + final ThreadBoundExecutor executor; long lastAddPushed; boolean notSupportBatch; @@ -233,9 +233,11 @@ public class LedgerHandle implements WriteHandle { // Two calls on purpose: chooseThread(long) hashes the raw id while chooseThread(Object) goes through // hashCode(), and Long.hashCode folds the high bits. Boxing the id would move ledgers with ids >= 2^31 // to a different thread than the other ledger-id keyed dispatches (e.g. OrderedGenericCallback). - this.executor = orderingKey == null + // The main worker pool is an OrderedExecutor, whose threads implement ThreadBoundExecutor whether or not + // they are decorated for task tracing or MDC preservation. + this.executor = (ThreadBoundExecutor) (orderingKey == null ? clientCtx.getMainWorkerPool().chooseThread(ledgerId) - : clientCtx.getMainWorkerPool().chooseThread(orderingKey); + : clientCtx.getMainWorkerPool().chooseThread(orderingKey)); if (clientCtx.getConf().enableStickyReads && getLedgerMetadata().getEnsembleSize() == getLedgerMetadata().getWriteQuorumSize()) { @@ -1119,8 +1121,9 @@ public class LedgerHandle implements WriteHandle { } if (isHandleWritable()) { - // Ledger handle in read/write mode: submit to OSE for ordered execution. - executeOrdered(op); + // Ledger handle in read/write mode: submit to OSE for ordered execution, unless the + // caller is already on the ledger's thread. + executor.executeOrRun(op); } else { // Read-only ledger handle: bypass OSE and execute read directly in client thread. // This avoids a context-switch to OSE thread and thus reduces latency. @@ -1299,8 +1302,9 @@ public class LedgerHandle implements WriteHandle { } if (isHandleWritable()) { - // Ledger handle in read/write mode: submit to OSE for ordered execution. - executeOrdered(op); + // Ledger handle in read/write mode: submit to OSE for ordered execution, unless the + // caller is already on the ledger's thread. + executor.executeOrRun(op); } else { // Read-only ledger handle: bypass OSE and execute read directly in client thread. // This avoids a context-switch to OSE thread and thus reduces latency. diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/LedgerHandleInlineReadTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/LedgerHandleInlineReadTest.java new file mode 100644 index 0000000000..a6b76b0136 --- /dev/null +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/LedgerHandleInlineReadTest.java @@ -0,0 +1,148 @@ +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ +package org.apache.bookkeeper.client; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertTrue; + +import com.google.common.collect.Lists; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; +import org.apache.bookkeeper.client.api.LedgerEntries; +import org.apache.bookkeeper.client.api.LedgerMetadata; +import org.apache.bookkeeper.client.api.WriteFlag; +import org.apache.bookkeeper.common.concurrent.FutureUtils; +import org.apache.bookkeeper.common.util.OrderedExecutor; +import org.apache.bookkeeper.net.BookieId; +import org.apache.bookkeeper.net.BookieSocketAddress; +import org.apache.bookkeeper.proto.MockBookieClient; +import org.apache.bookkeeper.proto.MockBookies; +import org.apache.bookkeeper.versioning.Versioned; +import org.junit.Test; + +/** + * Reads issued from the ledger's own worker thread are initiated inline instead of being queued on it. + */ +public class LedgerHandleInlineReadTest { + + private static final BookieId b1 = new BookieSocketAddress("b1", 3181).toBookieId(); + private static final BookieId b2 = new BookieSocketAddress("b2", 3181).toBookieId(); + private static final BookieId b3 = new BookieSocketAddress("b3", 3181).toBookieId(); + + private MockClientContext clientCtx; + private LedgerHandle lh; + private final AtomicReference<Thread> readIssuedOn = new AtomicReference<>(); + + private void setup(boolean decoratedThreads) throws Exception { + if (decoratedThreads) { + // Task tracing wraps the pool's threads in decorators; the handle must see through them + OrderedExecutor pool = OrderedExecutor.newBuilder().name("inline-read-test").numThreads(1) + .traceTaskExecution(true).build(); + MockBookies mockBookies = new MockBookies(); + clientCtx = MockClientContext.create(mockBookies) + .setMainWorkerPool(pool) + .setBookieClient(new MockBookieClient(pool, mockBookies)); + } else { + clientCtx = MockClientContext.create(); + } + Versioned<LedgerMetadata> md = ClientUtil.setupLedger(clientCtx, 10L, + LedgerMetadataBuilder.create().newEnsembleEntry(0L, Lists.newArrayList(b1, b2, b3))); + lh = new LedgerHandle(clientCtx, 10L, md, BookKeeper.DigestType.CRC32C, ClientUtil.PASSWD, + WriteFlag.NONE); + lh.append("entry".getBytes(UTF_8)); + + // The hook runs synchronously inside the bookie client's readEntry, so it records the thread + // that initiates the read request. + clientCtx.getMockBookieClient().setPreReadHook((bookie, ledgerId, entryId) -> { + readIssuedOn.compareAndSet(null, Thread.currentThread()); + return FutureUtils.value(null); + }); + } + + @Test(timeout = 30000) + public void testReadFromLedgerThreadIsInitiatedInline() throws Exception { + setup(false); + assertReadsFromLedgerThreadAreInitiatedInline(); + } + + @Test(timeout = 30000) + public void testReadFromLedgerThreadIsInitiatedInlineWithDecoratedThreads() throws Exception { + setup(true); + assertReadsFromLedgerThreadAreInitiatedInline(); + } + + private void assertReadsFromLedgerThreadAreInitiatedInline() throws Exception { + assertTrue(issuedBeforeReturnOnLedgerThread(result -> { + lh.readAsync(0, 0).whenComplete((entries, ex) -> complete(result, entries, ex)); + })); + + readIssuedOn.set(null); + assertTrue(issuedBeforeReturnOnLedgerThread(result -> { + lh.asyncReadEntries(0, 0, (rc, handle, entries, ctx) -> { + if (rc == BKException.Code.OK) { + result.complete(null); + } else { + result.completeExceptionally(BKException.create(rc)); + } + }, null); + })); + } + + @Test(timeout = 30000) + public void testReadFromOtherThreadIsQueuedOnLedgerThread() throws Exception { + setup(false); + try (LedgerEntries entries = lh.readAsync(0, 0).get(10, TimeUnit.SECONDS)) { + assertFalse(Thread.currentThread() == readIssuedOn.get()); + } + CompletableFuture<Thread> ledgerThread = new CompletableFuture<>(); + lh.executor.execute(() -> ledgerThread.complete(Thread.currentThread())); + assertSame(ledgerThread.get(10, TimeUnit.SECONDS), readIssuedOn.get()); + } + + /** + * Runs {@code read} on the ledger thread and returns whether the read request had reached the bookie + * client, on that same thread, by the time the call returned. Also waits for the read to complete. + */ + private boolean issuedBeforeReturnOnLedgerThread(Consumer<CompletableFuture<Void>> read) throws Exception { + CompletableFuture<Void> completed = new CompletableFuture<>(); + CompletableFuture<Boolean> issuedBeforeReturn = new CompletableFuture<>(); + lh.executor.execute(() -> { + read.accept(completed); + issuedBeforeReturn.complete(readIssuedOn.get() == Thread.currentThread()); + }); + boolean inline = issuedBeforeReturn.get(10, TimeUnit.SECONDS); + completed.get(10, TimeUnit.SECONDS); + return inline; + } + + private static void complete(CompletableFuture<Void> result, LedgerEntries entries, Throwable ex) { + if (ex != null) { + result.completeExceptionally(ex); + } else { + entries.close(); + result.complete(null); + } + } +} diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java index ffa59835ba..7a3ff584ba 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockBookKeeperTestCase.java @@ -210,7 +210,7 @@ public abstract class MockBookKeeperTestCase { @Override public OrderedExecutor getMainWorkerPool() { - return scheduler; + return executor; } @Override diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java index 93078a0512..8f0e9365d8 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/client/MockClientContext.java @@ -56,7 +56,9 @@ public class MockClientContext implements ClientContext { static MockClientContext create(MockBookies mockBookies) throws Exception { ClientConfiguration conf = new ClientConfiguration(); - OrderedScheduler scheduler = OrderedScheduler.newSchedulerBuilder().name("mock-executor").numThreads(1).build(); + OrderedExecutor executor = OrderedExecutor.newBuilder().name("mock-executor").numThreads(1).build(); + OrderedScheduler scheduler = OrderedScheduler.newSchedulerBuilder() + .name("mock-scheduler").numThreads(1).build(); MockRegistrationClient regClient = new MockRegistrationClient(); EnsemblePlacementPolicy placementPolicy = new DefaultEnsemblePlacementPolicy(); BookieWatcherImpl bookieWatcherImpl = new BookieWatcherImpl(conf, placementPolicy, @@ -71,9 +73,9 @@ public class MockClientContext implements ClientContext { .setBookieWatcher(bookieWatcherImpl) .setPlacementPolicy(placementPolicy) .setRegistrationClient(regClient) - .setBookieClient(new MockBookieClient(scheduler, mockBookies)) + .setBookieClient(new MockBookieClient(executor, mockBookies)) .setByteBufAllocator(UnpooledByteBufAllocator.DEFAULT) - .setMainWorkerPool(scheduler) + .setMainWorkerPool(executor) .setScheduler(scheduler) .setClientStats(BookKeeperClientStats.newInstance(NullStatsLogger.INSTANCE)) .setIsClientClosed(() -> false);
