From f1952c58af0f2dffe11551e877412b58cf6ceaa7 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Wed, 12 Aug 2026 03:00:13 +0800 Subject: [PATCH 1/2] Fix netty io thread block issue. --- .../proto/BatchedReadEntryProcessor.java | 1 - .../proto/BookieRequestProcessor.java | 2 - .../bookkeeper/proto/ReadEntryProcessor.java | 7 +- .../proto/ReadEntryProcessorV3.java | 3 +- .../proto/ReadEntryProcessorTest.java | 84 +++++++++++++++++++ 5 files changed, 91 insertions(+), 6 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BatchedReadEntryProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BatchedReadEntryProcessor.java index aa73e1986fb..bc95cb0ff91 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BatchedReadEntryProcessor.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BatchedReadEntryProcessor.java @@ -42,7 +42,6 @@ public static BatchedReadEntryProcessor create(BatchedReadRequest request, rep.fenceThreadPool = fenceThreadPool; rep.throttleReadResponses = throttleReadResponses; rep.maxBatchReadSize = maxBatchReadSize; - requestProcessor.onReadRequestStart(requestHandler.ctx().channel()); return rep; } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java index 2e1b0ab74fc..4b83c936857 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/BookieRequestProcessor.java @@ -551,7 +551,6 @@ private void processReadRequestV3(final Request r, final BookieRequestHandler re .setEntryId(r.getReadRequest().getEntryId()) .setStatus(StatusCode.ETOOMANYREQUESTS); read.sendResponse(StatusCode.ETOOMANYREQUESTS, resp, requestStats.getReadRequestStats()); - onReadRequestFinish(); } } } @@ -707,7 +706,6 @@ private void processReadRequest(final BookieProtocol.ReadRequest r, final Bookie BookieProtocol.ETOOMANYREQUESTS, ResponseBuilder.buildErrorResponse(BookieProtocol.ETOOMANYREQUESTS, r), requestStats.getReadRequestStats()); - onReadRequestFinish(); read.recycle(); } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java index 29cd159788b..7307bd44a95 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessor.java @@ -51,10 +51,15 @@ public static ReadEntryProcessor create(ReadRequest request, rep.init(request, requestHandler, requestProcessor); rep.fenceThreadPool = fenceThreadPool; rep.throttleReadResponses = throttleReadResponses; - requestProcessor.onReadRequestStart(requestHandler.ctx().channel()); return rep; } + @Override + public void run() { + requestProcessor.onReadRequestStart(requestHandler.ctx().channel()); + super.run(); + } + @Override protected void processPacket() { log.debug().attr("request", request).log("Received new read request"); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java index 1bf99cbd0d6..5630d75bbbf 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/proto/ReadEntryProcessorV3.java @@ -54,7 +54,6 @@ public ReadEntryProcessorV3(Request request, BookieRequestProcessor requestProcessor, ExecutorService fenceThreadPool) { super(request, requestHandler, requestProcessor); - requestProcessor.onReadRequestStart(requestHandler.ctx().channel()); this.readRequest = request.getReadRequest(); this.ledgerId = readRequest.getLedgerId(); @@ -267,6 +266,7 @@ protected ReadResponse getReadResponse() { @Override public void run() { + requestProcessor.onReadRequestStart(requestHandler.ctx().channel()); requestProcessor.getRequestStats().getReadEntrySchedulingDelayStats().registerSuccessfulEvent( MathUtils.elapsedNanos(enqueueNanos), TimeUnit.NANOSECONDS); if (!requestHandler.ctx().channel().isOpen()) { @@ -374,4 +374,3 @@ public String toString() { return RequestUtils.toSafeString(request); } } - diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java index 52bb1eeeb86..47ccf238c21 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java @@ -20,11 +20,13 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.Mockito.RETURNS_SELF; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -33,12 +35,18 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelPromise; import io.netty.channel.DefaultChannelPromise; +import io.netty.channel.DefaultEventLoopGroup; import io.netty.channel.EventLoop; +import io.netty.channel.EventLoopGroup; import java.io.IOException; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicReference; import org.apache.bookkeeper.bookie.Bookie; import org.apache.bookkeeper.bookie.BookieException; @@ -197,4 +205,80 @@ public void testNonFenceRequest() throws Exception { assertEquals(BookieProtocol.READENTRY, response.getOpCode()); assertEquals(BookieProtocol.EOK, response.getErrorCode()); } + + /** + * Blocking request creation on the event loop can starve write completion that would release a read permit. + */ + @Test + public void testCreateDoesNotStarveWriteCompletionThatReleasesReadPermit() throws Exception { + EventLoopGroup eventLoopGroup = new DefaultEventLoopGroup(1); + ExecutorService service = Executors.newSingleThreadExecutor(); + Semaphore readsSemaphore = new Semaphore(1); + CountDownLatch secondReadBlocked = new CountDownLatch(1); + CountDownLatch writeCompleted = new CountDownLatch(1); + CountDownLatch secondCreateReturned = new CountDownLatch(1); + ReadEntryProcessor[] processor = new ReadEntryProcessor[1]; + try { + EventLoop eventLoop = eventLoopGroup.next(); + when(channel.eventLoop()).thenReturn(eventLoop); + ChannelPromise writePromise = new DefaultChannelPromise(channel, eventLoop); + when(channel.writeAndFlush(any())).thenReturn(writePromise); + + doAnswer(inv -> { + if (!readsSemaphore.tryAcquire()) { + secondReadBlocked.countDown(); + readsSemaphore.acquireUninterruptibly(); + } + return null; + }).when(requestProcessor).onReadRequestStart(any(Channel.class)); + doAnswer(inv -> { + readsSemaphore.release(); + return null; + }).when(requestProcessor).onReadRequestFinish(); + + requestProcessor.onReadRequestStart(channel); + Future firstResponse = service.submit(() -> { + channel.writeAndFlush(new Object()).get(); + requestProcessor.onReadRequestFinish(); + writeCompleted.countDown(); + return null; + }); + + long ledgerId = System.currentTimeMillis(); + ReadRequest request = ReadRequest.create( + BookieProtocol.CURRENT_PROTOCOL_VERSION, ledgerId, 1, (short) 0, new byte[]{}); + eventLoop.execute(() -> { + processor[0] = ReadEntryProcessor.create(request, requestHandler, requestProcessor, null, true); + secondCreateReturned.countDown(); + }); + + long waitUntilNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(1); + while (secondReadBlocked.getCount() > 0 + && secondCreateReturned.getCount() > 0 + && System.nanoTime() < waitUntilNanos) { + TimeUnit.MILLISECONDS.sleep(10); + } + assertTrue("second read create should either return or start waiting for a permit", + secondReadBlocked.getCount() == 0 || secondCreateReturned.getCount() == 0); + eventLoop.execute(() -> writePromise.setSuccess()); + + try { + assertTrue("write completion must not be starved behind a blocked read create", + writeCompleted.await(1, TimeUnit.SECONDS)); + } finally { + if (writeCompleted.getCount() > 0) { + readsSemaphore.release(); + } + } + assertTrue("second read create should return without blocking the event loop", + secondCreateReturned.await(1, TimeUnit.SECONDS)); + firstResponse.get(1, TimeUnit.SECONDS); + } finally { + if (processor[0] != null) { + processor[0].recycle(); + } + service.shutdownNow(); + eventLoopGroup.shutdownGracefully(0, 0, TimeUnit.MILLISECONDS).sync(); + } + } } From d32c12df0d34df829bbfe8228880c62f83fcc597 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Wed, 12 Aug 2026 11:38:29 +0800 Subject: [PATCH 2/2] code clean. --- .../apache/bookkeeper/proto/ReadEntryProcessorTest.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java index 47ccf238c21..44abe2d3f53 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/proto/ReadEntryProcessorTest.java @@ -207,10 +207,11 @@ public void testNonFenceRequest() throws Exception { } /** - * Blocking request creation on the event loop can starve write completion that would release a read permit. + * Blocking request creation on the event loop can starve the write future + * that a throttled V2 read waits on before releasing a read permit. */ @Test - public void testCreateDoesNotStarveWriteCompletionThatReleasesReadPermit() throws Exception { + public void testCreateDoesNotStarveV2WriteCompletionNeededToReleaseReadPermit() throws Exception { EventLoopGroup eventLoopGroup = new DefaultEventLoopGroup(1); ExecutorService service = Executors.newSingleThreadExecutor(); Semaphore readsSemaphore = new Semaphore(1); @@ -263,7 +264,7 @@ public void testCreateDoesNotStarveWriteCompletionThatReleasesReadPermit() throw eventLoop.execute(() -> writePromise.setSuccess()); try { - assertTrue("write completion must not be starved behind a blocked read create", + assertTrue("V2 write future completion must not be starved behind a blocked read create", writeCompleted.await(1, TimeUnit.SECONDS)); } finally { if (writeCompleted.getCount() > 0) {