From 5325257d2623a4c3bd0753e11959bb658e565f13 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Sat, 18 Jul 2026 12:15:01 +0800 Subject: [PATCH 1/2] Fix directIO reader progress after short reads --- .../directentrylogger/DirectReader.java | 19 ++++++-- .../directentrylogger/TestDirectReader.java | 43 +++++++++++++++++++ 2 files changed, 59 insertions(+), 3 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java index 707bf307c05..f9a12f4b396 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java @@ -157,7 +157,7 @@ private int readBytesIntoBuf(ByteBuf buf, long offset, int size) throws IOExcept .kv("offset", offset) .kv("size", size).toString()); } - return nativeBuffer.readByteBuf(buf, offsetInBuffer, size); + return nativeBuffer.readByteBuf(buf, offsetInBuffer, sizeInBuffer); } } @@ -226,8 +226,21 @@ void readBlock(long offset) throws IOException { if ((bytesOutstanding - bytesRead) <= 0) { break; } - bytesOutstanding -= bytesRead & Buffer.ALIGNMENT; - bufferOffset += bytesRead & Buffer.ALIGNMENT; + long alignedBytesRead = bytesRead & ~(Buffer.ALIGNMENT - 1L); + if (alignedBytesRead <= 0) { + readBlockStats.registerFailedEvent(System.nanoTime() - startNs, TimeUnit.NANOSECONDS); + throw new EOFException(exMsg("Short read did not make aligned progress") + .kv("requestedBytes", blockSize) + .kv("offset", blockStart) + .kv("expectedBytes", Math.min(blockSize, bytesAvailable)) + .kv("bytesOutstanding", bytesOutstanding) + .kv("bufferOffset", bufferOffset) + .kv("bytesRead", bytesRead) + .kv("file", filename) + .kv("fd", fd).toString()); + } + bytesOutstanding -= alignedBytesRead; + bufferOffset += alignedBytesRead; } } catch (NativeIOException ne) { readBlockStats.registerFailedEvent(System.nanoTime() - startNs, TimeUnit.NANOSECONDS); diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java index 04cd4a0c429..05b9848e887 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java @@ -445,6 +445,49 @@ buffers, new NativeIOImpl(), Logger.get(TestDirectReader.class))) { } } + @Test + public void testPartialReadProgressesByAlignedBytes() throws Exception { + File ledgerDir = tmpDirs.createNew("partialReadAligned", "logs"); + + writeFileWithPattern(ledgerDir, 1234, 0xbeefcafe, 1, 1 << 20); + + class ShortReadNativeIO extends NativeIOImpl { + int calls; + + @Override + public long pread(int fd, long buf, long size, long offset) throws NativeIOException { + if (calls == 0) { + assertThat(offset, equalTo(0L)); + } else if (calls == 1) { + assertThat(offset, equalTo((long) Buffer.ALIGNMENT * 2)); + } + calls++; + + long read = super.pread(fd, buf, size, offset); + return Math.min(read, Buffer.ALIGNMENT * 2L); + } + } + + ShortReadNativeIO nativeIO = new ShortReadNativeIO(); + try (LogReader reader = new DirectReader(1234, logFilename(ledgerDir, 1234), + ByteBufAllocator.DEFAULT, + nativeIO, Buffer.ALIGNMENT * 4, + 1 << 20, opLogger)) { + ByteBuf bb = reader.readBufferAt(0, Buffer.ALIGNMENT * 4); + try { + for (int block = 0; block < 4; block++) { + for (int i = 0; i < Buffer.ALIGNMENT / Integer.BYTES; i++) { + assertThat(bb.readInt(), equalTo(0xbeefcafe + block)); + } + } + assertThat(bb.readableBytes(), equalTo(0)); + } finally { + bb.release(); + } + } + assertThat(nativeIO.calls, equalTo(2)); + } + @Test public void testLargeEntry() throws Exception { File ledgerDir = tmpDirs.createNew("largeEntries", "logs"); From ab86f2874ece296192a231b443b38aecd7e428a8 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Sun, 19 Jul 2026 11:07:30 +0800 Subject: [PATCH 2/2] Handle failed short reads safely --- .../directentrylogger/DirectReader.java | 9 ++-- .../directentrylogger/TestDirectReader.java | 43 +++++++++++++++++-- 2 files changed, 45 insertions(+), 7 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java index f9a12f4b396..29c1c52fcb5 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/directentrylogger/DirectReader.java @@ -207,6 +207,9 @@ void readBlock(long offset) throws IOException { long bytesToRead = Math.min(blockSize, bytesAvailable); long bytesOutstanding = bytesToRead; long bytesRead = -1; + // A failed read may still overwrite nativeBuffer, so invalidate the + // cached block before loading a new one. + clearCache(); try { while (true) { long readSize = blockSize - bufferOffset; @@ -229,9 +232,9 @@ void readBlock(long offset) throws IOException { long alignedBytesRead = bytesRead & ~(Buffer.ALIGNMENT - 1L); if (alignedBytesRead <= 0) { readBlockStats.registerFailedEvent(System.nanoTime() - startNs, TimeUnit.NANOSECONDS); - throw new EOFException(exMsg("Short read did not make aligned progress") - .kv("requestedBytes", blockSize) - .kv("offset", blockStart) + throw new IOException(exMsg("Short read did not make aligned progress") + .kv("requestedBytes", readSize) + .kv("offset", blockStart + bufferOffset) .kv("expectedBytes", Math.min(blockSize, bytesAvailable)) .kv("bytesOutstanding", bytesOutstanding) .kv("bufferOffset", bufferOffset) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java index 05b9848e887..b98f8cf7567 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/directentrylogger/TestDirectReader.java @@ -392,8 +392,7 @@ public void testPartialRead() throws Exception { NativeIOImpl nativeIO = new NativeIOImpl() { @Override public long pread(int fd, long buf, long size, long offset) throws NativeIOException { - long read = super.pread(fd, buf, size, offset); - return Math.min(read, Buffer.ALIGNMENT); // force only less than a buffer read + return super.pread(fd, buf, Math.min(size, Buffer.ALIGNMENT), offset); } @Override @@ -463,8 +462,7 @@ public long pread(int fd, long buf, long size, long offset) throws NativeIOExcep } calls++; - long read = super.pread(fd, buf, size, offset); - return Math.min(read, Buffer.ALIGNMENT * 2L); + return super.pread(fd, buf, Math.min(size, Buffer.ALIGNMENT * 2L), offset); } } @@ -488,6 +486,43 @@ public long pread(int fd, long buf, long size, long offset) throws NativeIOExcep assertThat(nativeIO.calls, equalTo(2)); } + @Test + public void testFailedPartialReadInvalidatesCachedBlock() throws Exception { + File ledgerDir = tmpDirs.createNew("partialReadCache", "logs"); + + writeFileWithPattern(ledgerDir, 1234, 0xbeefcafe, 1, 1 << 20); + + class FailingShortReadNativeIO extends NativeIOImpl { + int calls; + + @Override + public long pread(int fd, long buf, long size, long offset) throws NativeIOException { + calls++; + if (calls == 2) { + return super.pread(fd, buf, Buffer.ALIGNMENT * 2L, offset); + } else if (calls == 3) { + return 0; + } + return super.pread(fd, buf, size, offset); + } + } + + FailingShortReadNativeIO nativeIO = new FailingShortReadNativeIO(); + try (LogReader reader = new DirectReader(1234, logFilename(ledgerDir, 1234), + ByteBufAllocator.DEFAULT, + nativeIO, Buffer.ALIGNMENT * 4, + 1 << 20, opLogger)) { + assertThat(reader.readIntAt(0), equalTo(0xbeefcafe)); + + IOException exception = Assertions.assertThrows( + IOException.class, () -> reader.readIntAt(Buffer.ALIGNMENT * 4L)); + Assertions.assertFalse(exception instanceof EOFException); + + assertThat(reader.readIntAt(0), equalTo(0xbeefcafe)); + } + assertThat(nativeIO.calls, equalTo(4)); + } + @Test public void testLargeEntry() throws Exception { File ledgerDir = tmpDirs.createNew("largeEntries", "logs");