From 6f964d6e02a64a469b19aec43f32ccd7a31b1e5a Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 21:05:21 +0800 Subject: [PATCH 1/9] fix server side bug --- .../scm/storage/StreamBlockInputStream.java | 2 +- .../ContainerCommandResponseBuilders.java | 16 ++- .../container/keyvalue/KeyValueHandler.java | 59 +++++----- .../keyvalue/TestKeyValueHandler.java | 102 ++++++++++++++++++ .../main/proto/DatanodeClientProtocol.proto | 1 + 5 files changed, 150 insertions(+), 30 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index b008abdb56c4..1ce3e5438ef4 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -541,7 +541,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman final StreamingReadResponse r = getResponse(); LOG.warn("Failed to process block {} response at offset={}, size={}: {}, {}", getBlockID().getContainerBlockID(), - offset, data.size(), StringUtils.bytes2Hex(data.substring(0, 10).asReadOnlyByteBuffer()), + offset, data.size(), StringUtils.bytes2Hex(data.asReadOnlyByteBuffer()), readBlock.getChecksumData(), e); setFailed(e); r.getRequestObserver().onError(e); diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java index c60f3d7449a0..e2bcc4a1667c 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java @@ -27,6 +27,7 @@ import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.BlockData; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfoList; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto.Builder; @@ -39,6 +40,7 @@ import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ListBlockResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.PutBlockResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.PutSmallFileResponseProto; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadBlockResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadChunkResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadContainerResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result; @@ -336,16 +338,20 @@ public static ContainerCommandResponseProto getReadChunkResponse( } public static ContainerCommandResponseProto getReadBlockResponse( - ContainerCommandRequestProto request, ChecksumData checksumData, ByteBuffer data, long offset) { + ContainerCommandRequestProto request, ChecksumData checksumData, + ByteBuffer data, long offset, List chunkInfoList, boolean verifyChecksum) { - ContainerProtos.ReadBlockResponseProto response = ContainerProtos.ReadBlockResponseProto.newBuilder() + ContainerProtos.ReadBlockResponseProto.Builder builder = ReadBlockResponseProto.newBuilder() .setChecksumData(checksumData.getProtoBufMessage()) .setData(ByteString.copyFrom(data)) - .setOffset(offset) - .build(); + .setOffset(offset); + + if (verifyChecksum) { + builder.setChunkInfoList(ChunkInfoList.newBuilder().addAllChunks(chunkInfoList)); + } return getSuccessResponseBuilder(request) - .setReadBlock(response) + .setReadBlock(builder) .build(); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java index 1b3f399b46fa..34999aec1e56 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java @@ -2208,7 +2208,7 @@ private long readBlockImpl(ContainerCommandRequestProto request, RandomAccessFil } } final ContainerCommandResponseProto response = getReadBlockResponse( - request, checksumData, buffer, adjustedOffset); + request, checksumData, buffer, adjustedOffset, chunkInfos, verifyChecksum); final int dataLength = response.getReadBlock().getData().size(); LOG.debug("server onNext response {}: dataLength={}, numChecksums={}", numResponses, dataLength, response.getReadBlock().getChecksumData().getChecksumsList().size()); @@ -2225,30 +2225,41 @@ private long readBlockImpl(ContainerCommandRequestProto request, RandomAccessFil static List getChecksums(long blockOffset, int readLength, int bytesPerChunk, int bytesPerChecksum, final List chunks) { assertSame(0, blockOffset % bytesPerChecksum, "blockOffset % bytesPerChecksum"); - final int numChecksums = 1 + (readLength - 1) / bytesPerChecksum; - final List checksums = new ArrayList<>(numChecksums); - for (int i = 0; i < numChecksums; i++) { - // As the checksums are stored "chunk by chunk", we need to figure out which chunk we start reading from, - // and its offset to pull out the correct checksum bytes for each read. - final int n = i * bytesPerChecksum; - final long offset = blockOffset + n; - final int c = Math.toIntExact(offset / bytesPerChunk); - final int chunkOffset = Math.toIntExact(offset % bytesPerChunk); - final int csi = chunkOffset / bytesPerChecksum; - - assertTrue(c < chunks.size(), - () -> "chunkIndex = " + c + " >= chunk.size()" + chunks.size()); - final ContainerProtos.ChunkInfo chunk = chunks.get(c); - if (c < chunks.size() - 1) { - assertSame(bytesPerChunk, chunk.getLen(), "bytesPerChunk"); - } - final ContainerProtos.ChecksumData checksumDataProto = chunks.get(c).getChecksumData(); - assertSame(bytesPerChecksum, checksumDataProto.getBytesPerChecksum(), "bytesPerChecksum"); - final List checksumsList = checksumDataProto.getChecksumsList(); - assertTrue(csi < checksumsList.size(), - () -> "checksumIndex = " + csi + " >= checksumsList.size()" + checksumsList.size()); - checksums.add(checksumsList.get(csi)); + final List checksums = new ArrayList<>(); + + long currentChunkOffset = 0; + for (ContainerProtos.ChunkInfo chunk : chunks) { + long chunkStart = currentChunkOffset; + long chunkEnd = chunkStart + chunk.getLen(); + + long overlapStart = Math.max(blockOffset, chunkStart); + long overlapEnd = Math.min(blockOffset + readLength, chunkEnd); + + if (overlapStart < overlapEnd) { + long offsetInChunk = overlapStart - chunkStart; + long endOffsetInChunk = overlapEnd - chunkStart; + + int firstChecksumIndex = Math.toIntExact(offsetInChunk / bytesPerChecksum); + int lastChecksumIndex = Math.toIntExact((endOffsetInChunk - 1) / bytesPerChecksum); + + ContainerProtos.ChecksumData checksumDataProto = chunk.getChecksumData(); + assertSame(bytesPerChecksum, checksumDataProto.getBytesPerChecksum(), "bytesPerChecksum"); + List checksumsList = checksumDataProto.getChecksumsList(); + + for (int csi = firstChecksumIndex; csi <= lastChecksumIndex; csi++) { + final int finalCsi = csi; + assertTrue(finalCsi < checksumsList.size(), + () -> "checksumIndex = " + finalCsi + " >= checksumsList.size()" + checksumsList.size()); + checksums.add(checksumsList.get(finalCsi)); + } + } + + currentChunkOffset = chunkEnd; + if (currentChunkOffset >= blockOffset + readLength) { + break; + } } + return checksums; } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java index f77d6fec2cbc..dc0cec0ddc49 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java @@ -1162,4 +1162,106 @@ boolean deleteUnreferencedFile(File file) { return false; } } + + @Test + public void testGetChecksumsWithVaryingChunkSizes() { + int bytesPerChecksum = 1024; + int bytesPerChunk = 16 * 1024; + + ContainerProtos.ChunkInfo chunk1 = ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk1") + .setOffset(0) + .setLen(1024) + .setChecksumData(ContainerProtos.ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32) + .setBytesPerChecksum(bytesPerChecksum) + .addChecksums(org.apache.ratis.thirdparty.com.google.protobuf.ByteString.copyFromUtf8("chk1")) + .build()) + .build(); + + ContainerProtos.ChunkInfo chunk2 = ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk2") + .setOffset(1024) + .setLen(10) + .setChecksumData(ContainerProtos.ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32) + .setBytesPerChecksum(bytesPerChecksum) + .addChecksums(org.apache.ratis.thirdparty.com.google.protobuf.ByteString.copyFromUtf8("chk2")) + .build()) + .build(); + + ContainerProtos.ChunkInfo chunk3 = ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk3") + .setOffset(1034) + .setLen(2048) + .setChecksumData(ContainerProtos.ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32) + .setBytesPerChecksum(bytesPerChecksum) + .addChecksums(org.apache.ratis.thirdparty.com.google.protobuf.ByteString.copyFromUtf8("chk3-1")) + .addChecksums(org.apache.ratis.thirdparty.com.google.protobuf.ByteString.copyFromUtf8("chk3-2")) + .build()) + .build(); + + List chunks = java.util.Arrays.asList(chunk1, chunk2, chunk3); + + // Read full block (1024 + 10 + 2048) + List checksums = KeyValueHandler.getChecksums(0, 3082, bytesPerChunk, bytesPerChecksum, chunks); + assertEquals(4, checksums.size()); + assertEquals("chk1", checksums.get(0).toStringUtf8()); + assertEquals("chk2", checksums.get(1).toStringUtf8()); + assertEquals("chk3-1", checksums.get(2).toStringUtf8()); + assertEquals("chk3-2", checksums.get(3).toStringUtf8()); + + // Read from offset 1024 + checksums = KeyValueHandler.getChecksums(1024, 2058, bytesPerChunk, bytesPerChecksum, chunks); + assertEquals(3, checksums.size()); + assertEquals("chk2", checksums.get(0).toStringUtf8()); + assertEquals("chk3-1", checksums.get(1).toStringUtf8()); + assertEquals("chk3-2", checksums.get(2).toStringUtf8()); + + // Read from offset 2048 + checksums = KeyValueHandler.getChecksums(2048, 1034, bytesPerChunk, bytesPerChecksum, chunks); + assertEquals(2, checksums.size()); + assertEquals("chk3-1", checksums.get(0).toStringUtf8()); + assertEquals("chk3-2", checksums.get(1).toStringUtf8()); + } + + @Test + public void testGetChecksumsWithSmallChunks() { + int bytesPerChecksum = 1024; + int bytesPerChunk = 16 * 1024; + + ContainerProtos.ChunkInfo chunk1 = ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk1") + .setOffset(0) + .setLen(1) + .setChecksumData(ContainerProtos.ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32) + .setBytesPerChecksum(bytesPerChecksum) + .addChecksums(org.apache.ratis.thirdparty.com.google.protobuf.ByteString.copyFromUtf8("chk1")) + .build()) + .build(); + + ContainerProtos.ChunkInfo chunk2 = ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk2") + .setOffset(1) + .setLen(1) + .setChecksumData(ContainerProtos.ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32) + .setBytesPerChecksum(bytesPerChecksum) + .addChecksums(org.apache.ratis.thirdparty.com.google.protobuf.ByteString.copyFromUtf8("chk2")) + .build()) + .build(); + + List chunks = java.util.Arrays.asList(chunk1, chunk2); + + // Read full block (1 + 1 = 2 bytes) + List checksums = + KeyValueHandler.getChecksums(0, 2, bytesPerChunk, bytesPerChecksum, chunks); + + // According to the bug report, this should return 2 checksums. + assertEquals(2, checksums.size()); + assertEquals("chk1", checksums.get(0).toStringUtf8()); + assertEquals("chk2", checksums.get(1).toStringUtf8()); + } } diff --git a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto index e33b07b50aa4..881571a51e26 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -386,6 +386,7 @@ message ReadBlockResponseProto { required ChecksumData checksumData = 1; required uint64 offset = 2; required bytes data = 3; + optional ChunkInfoList chunkInfoList = 4; } message EchoRequestProto { From a65be9321d27296e0b0a342bdf735f675f061f8d Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 22:16:31 +0800 Subject: [PATCH 2/9] HDDS-15857. Fix checksum calculation in KeyValueHandler for variable-sized chunks --- .../scm/storage/StreamBlockInputStream.java | 41 ++++++++++++++++++- .../ContainerCommandResponseBuilders.java | 11 +++-- .../ozone/client/rpc/read/TestStreamRead.java | 38 +++++++++++++++++ 3 files changed, 83 insertions(+), 7 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index 1ce3e5438ef4..9ae79fcaab4f 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -532,7 +532,46 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman ByteBuffer data = readBlock.getData().asReadOnlyByteBuffer(); if (verifyChecksum) { ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData()); - Checksum.verifyChecksum(data, checksumData, 0); + if (checksumData.getChecksumType() == ContainerProtos.ChecksumType.NONE) { + // Checksum is set to NONE. No further verification is required. + } else if (readBlock.hasChunkInfoList()) { + int bytesPerChecksum = checksumData.getBytesPerChecksum(); + long blockOffset = readBlock.getOffset(); + long readLength = data.remaining(); + long currentChunkOffset = 0; + int checksumIndex = 0; + int dataOffset = 0; + + for (ContainerProtos.ChunkInfo chunk : readBlock.getChunkInfoList().getChunksList()) { + long chunkStart = currentChunkOffset; + long chunkEnd = chunkStart + chunk.getLen(); + + long overlapStart = Math.max(blockOffset, chunkStart); + long overlapEnd = Math.min(blockOffset + readLength, chunkEnd); + + if (overlapStart < overlapEnd) { + int overlapLen = Math.toIntExact(overlapEnd - overlapStart); + ByteBuffer chunkData = data.duplicate(); + chunkData.position(data.position() + dataOffset); + chunkData.limit(data.position() + dataOffset + overlapLen); + + Checksum.verifyChecksum(chunkData, checksumData, checksumIndex); + + dataOffset += overlapLen; + + long offsetInChunk = overlapStart - chunkStart; + long endOffsetInChunk = overlapEnd - chunkStart; + + int firstChecksumIndex = Math.toIntExact(offsetInChunk / bytesPerChecksum); + int lastChecksumIndex = Math.toIntExact((endOffsetInChunk - 1) / bytesPerChecksum); + + checksumIndex += (lastChecksumIndex - firstChecksumIndex + 1); + } + currentChunkOffset += chunk.getLen(); + } + } else { + Checksum.verifyChecksum(data, checksumData, 0); + } } offerToQueue(readBlock); } catch (Exception e) { diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java index e2bcc4a1667c..f0af4432b7de 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ContainerCommandResponseBuilders.java @@ -341,17 +341,16 @@ public static ContainerCommandResponseProto getReadBlockResponse( ContainerCommandRequestProto request, ChecksumData checksumData, ByteBuffer data, long offset, List chunkInfoList, boolean verifyChecksum) { - ContainerProtos.ReadBlockResponseProto.Builder builder = ReadBlockResponseProto.newBuilder() + ContainerProtos.ReadBlockResponseProto response = ReadBlockResponseProto.newBuilder() .setChecksumData(checksumData.getProtoBufMessage()) .setData(ByteString.copyFrom(data)) - .setOffset(offset); + .setOffset(offset) + .setChunkInfoList(ChunkInfoList.newBuilder().addAllChunks(chunkInfoList).build()) + .build(); - if (verifyChecksum) { - builder.setChunkInfoList(ChunkInfoList.newBuilder().addAllChunks(chunkInfoList)); - } return getSuccessResponseBuilder(request) - .setReadBlock(builder) + .setReadBlock(response) .build(); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java index 9fc217b6df3a..8d14f7558d89 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java @@ -301,4 +301,42 @@ static String runTestReadKey(String name, SizeInBytes keySize, SizeInBytes buffe print(name, keySizeByte, elapsedNanos, bufferSize, computedMD5); return computedMD5; } + + @Test + void testSmallChunksWithLargeChecksum() throws Exception { + final SizeInBytes bytesPerChecksum = SizeInBytes.valueOf(512); + System.out.println("cluster starting ..."); + try (MiniOzoneCluster cluster = newCluster(bytesPerChecksum.getSizeInt())) { + cluster.waitForClusterToBeReady(); + System.out.println("cluster ready"); + + OzoneConfiguration conf = cluster.getConf(); + OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class); + clientConfig.setStreamReadBlock(true); + final OzoneConfiguration steamReadConf = new OzoneConfiguration(conf); + steamReadConf.setFromObject(clientConfig); + + try (OzoneClient streamReadClient = OzoneClientFactory.getRpcClient(steamReadConf)) { + final TestBucket testBucket = TestBucket.newBuilder(streamReadClient).build(); + final String keyName = "keySmallChunks"; + + // Write a key with multiple small chunks + byte[] data = new byte[]{1, 2, 3, 4, 5}; + try (OutputStream out = testBucket.delegate().createKey(keyName, data.length, RatisReplicationConfig.getInstance(ONE), Collections.emptyMap())) { + for (byte b : data) { + out.write(b); + out.flush(); // Forces a chunk to be created + } + } + + // Read it back using stream read + try (InputStream in = testBucket.delegate().readKey(keyName)) { + byte[] readData = new byte[data.length]; + int bytesRead = in.read(readData); + assertEquals(data.length, bytesRead); + org.junit.jupiter.api.Assertions.assertArrayEquals(data, readData); + } + } + } + } } From 3779fad6d6f858b4d5470e482a227e6642b18796 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 22:36:22 +0800 Subject: [PATCH 3/9] update test --- .../org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java | 1 + 1 file changed, 1 insertion(+) diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java index 8d14f7558d89..c8817b5f08d0 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java @@ -313,6 +313,7 @@ void testSmallChunksWithLargeChecksum() throws Exception { OzoneConfiguration conf = cluster.getConf(); OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class); clientConfig.setStreamReadBlock(true); + clientConfig.setStreamBufferFlushDelay(false); final OzoneConfiguration steamReadConf = new OzoneConfiguration(conf); steamReadConf.setFromObject(clientConfig); From e9b31838cedcfab55d7557dd9a7b102aac4244ca Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 22:53:08 +0800 Subject: [PATCH 4/9] throw exception if checksum data is missing --- .../apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index e81e15fd0a50..1776f1c3e458 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -570,7 +570,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman currentChunkOffset += chunk.getLen(); } } else { - Checksum.verifyChecksum(data, checksumData, 0); + throw new IOException("Checksum data is missing for block " + getBlockID()); } } offerToQueue(readBlock); From e0ae9473844524b859407b3d5b2aeeaa47d8743d Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 22:58:31 +0800 Subject: [PATCH 5/9] fix checkstyle --- .../hadoop/hdds/scm/storage/StreamBlockInputStream.java | 4 +--- .../hadoop/ozone/container/keyvalue/TestKeyValueHandler.java | 3 ++- .../apache/hadoop/ozone/client/rpc/read/TestStreamRead.java | 3 ++- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index 1776f1c3e458..8efecc62325d 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -532,9 +532,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman ByteBuffer data = readBlock.getData().asReadOnlyByteBuffer(); if (verifyChecksum) { ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData()); - if (checksumData.getChecksumType() == ContainerProtos.ChecksumType.NONE) { - // Checksum is set to NONE. No further verification is required. - } else if (readBlock.hasChunkInfoList()) { + if (readBlock.hasChunkInfoList()) { int bytesPerChecksum = checksumData.getBytesPerChecksum(); long blockOffset = readBlock.getOffset(); long readLength = data.remaining(); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java index dc0cec0ddc49..0e3323338af1 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestKeyValueHandler.java @@ -119,6 +119,7 @@ import org.apache.hadoop.util.Time; import org.apache.ozone.test.GenericTestUtils; import org.apache.ozone.test.GenericTestUtils.LogCapturer; +import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; @@ -1205,7 +1206,7 @@ public void testGetChecksumsWithVaryingChunkSizes() { List chunks = java.util.Arrays.asList(chunk1, chunk2, chunk3); // Read full block (1024 + 10 + 2048) - List checksums = KeyValueHandler.getChecksums(0, 3082, bytesPerChunk, bytesPerChecksum, chunks); + List checksums = KeyValueHandler.getChecksums(0, 3082, bytesPerChunk, bytesPerChecksum, chunks); assertEquals(4, checksums.size()); assertEquals("chk1", checksums.get(0).toStringUtf8()); assertEquals("chk2", checksums.get(1).toStringUtf8()); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java index c8817b5f08d0..22fd89b23ee2 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java @@ -323,7 +323,8 @@ void testSmallChunksWithLargeChecksum() throws Exception { // Write a key with multiple small chunks byte[] data = new byte[]{1, 2, 3, 4, 5}; - try (OutputStream out = testBucket.delegate().createKey(keyName, data.length, RatisReplicationConfig.getInstance(ONE), Collections.emptyMap())) { + try (OutputStream out = testBucket.delegate() + .createKey(keyName, data.length, RatisReplicationConfig.getInstance(ONE), Collections.emptyMap())) { for (byte b : data) { out.write(b); out.flush(); // Forces a chunk to be created From ec1d56b3a490dd3f6dcae4b6183c102df31252a4 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 23:24:42 +0800 Subject: [PATCH 6/9] fix findbug --- .../apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index 8efecc62325d..f39fe9ceefc8 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -55,6 +55,7 @@ import org.apache.hadoop.io.retry.RetryPolicy; import org.apache.hadoop.ozone.common.Checksum; import org.apache.hadoop.ozone.common.ChecksumData; +import org.apache.hadoop.ozone.common.OzoneChecksumException; import org.apache.hadoop.security.token.Token; import org.apache.ratis.protocol.exceptions.TimeoutIOException; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; @@ -572,7 +573,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman } } offerToQueue(readBlock); - } catch (Exception e) { + } catch (IOException | RuntimeException e) { // Record the failure first: the log and observer calls below must not mask it. setFailed(e); final ByteString data = readBlock.getData(); From 86cf0f3f87bb664742846ffb19aa60a941ae3d87 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 14 Jul 2026 23:43:05 +0800 Subject: [PATCH 7/9] fix checkstyle --- .../apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java | 1 - 1 file changed, 1 deletion(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index f39fe9ceefc8..282ab059a580 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -55,7 +55,6 @@ import org.apache.hadoop.io.retry.RetryPolicy; import org.apache.hadoop.ozone.common.Checksum; import org.apache.hadoop.ozone.common.ChecksumData; -import org.apache.hadoop.ozone.common.OzoneChecksumException; import org.apache.hadoop.security.token.Token; import org.apache.ratis.protocol.exceptions.TimeoutIOException; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; From 445579dec6f6edf78ebe18cf6bdc45dd981a6119 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Wed, 15 Jul 2026 07:37:43 +0800 Subject: [PATCH 8/9] fix bugs --- .../scm/storage/StreamBlockInputStream.java | 78 +++++++++++-------- .../storage/TestStreamBlockInputStream.java | 18 ++++- .../ozone/contract/TestOzoneContractFSO.java | 1 + 3 files changed, 61 insertions(+), 36 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java index 282ab059a580..6ad4735aec36 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java @@ -41,6 +41,7 @@ import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadBlockResponseProto; import org.apache.hadoop.hdds.scm.OzoneClientConfig; import org.apache.hadoop.hdds.scm.StreamingReadResponse; @@ -55,6 +56,7 @@ import org.apache.hadoop.io.retry.RetryPolicy; import org.apache.hadoop.ozone.common.Checksum; import org.apache.hadoop.ozone.common.ChecksumData; +import org.apache.hadoop.ozone.common.OzoneChecksumException; import org.apache.hadoop.security.token.Token; import org.apache.ratis.protocol.exceptions.TimeoutIOException; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; @@ -533,40 +535,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman if (verifyChecksum) { ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData()); if (readBlock.hasChunkInfoList()) { - int bytesPerChecksum = checksumData.getBytesPerChecksum(); - long blockOffset = readBlock.getOffset(); - long readLength = data.remaining(); - long currentChunkOffset = 0; - int checksumIndex = 0; - int dataOffset = 0; - - for (ContainerProtos.ChunkInfo chunk : readBlock.getChunkInfoList().getChunksList()) { - long chunkStart = currentChunkOffset; - long chunkEnd = chunkStart + chunk.getLen(); - - long overlapStart = Math.max(blockOffset, chunkStart); - long overlapEnd = Math.min(blockOffset + readLength, chunkEnd); - - if (overlapStart < overlapEnd) { - int overlapLen = Math.toIntExact(overlapEnd - overlapStart); - ByteBuffer chunkData = data.duplicate(); - chunkData.position(data.position() + dataOffset); - chunkData.limit(data.position() + dataOffset + overlapLen); - - Checksum.verifyChecksum(chunkData, checksumData, checksumIndex); - - dataOffset += overlapLen; - - long offsetInChunk = overlapStart - chunkStart; - long endOffsetInChunk = overlapEnd - chunkStart; - - int firstChecksumIndex = Math.toIntExact(offsetInChunk / bytesPerChecksum); - int lastChecksumIndex = Math.toIntExact((endOffsetInChunk - 1) / bytesPerChecksum); - - checksumIndex += (lastChecksumIndex - firstChecksumIndex + 1); - } - currentChunkOffset += chunk.getLen(); - } + verifyChecksumForReadBlock(data, checksumData, readBlock); } else { throw new IOException("Checksum data is missing for block " + getBlockID()); } @@ -665,5 +634,46 @@ public void setStreamingReadResponse(StreamingReadResponse streamingReadResponse public String toString() { return name; } + + private void verifyChecksumForReadBlock( + ByteBuffer data, ChecksumData checksumData, ReadBlockResponseProto readBlock) + throws OzoneChecksumException { + if (!checksumData.getChecksumType().equals(ChecksumType.NONE)) { + int bytesPerChecksum = checksumData.getBytesPerChecksum(); + long blockOffset = readBlock.getOffset(); + long readLength = data.remaining(); + long currentChunkOffset = 0; + int checksumIndex = 0; + int dataOffset = 0; + + for (ContainerProtos.ChunkInfo chunk : readBlock.getChunkInfoList().getChunksList()) { + long chunkStart = currentChunkOffset; + long chunkEnd = chunkStart + chunk.getLen(); + + long overlapStart = Math.max(blockOffset, chunkStart); + long overlapEnd = Math.min(blockOffset + readLength, chunkEnd); + + if (overlapStart < overlapEnd) { + int overlapLen = Math.toIntExact(overlapEnd - overlapStart); + ByteBuffer chunkData = data.duplicate(); + chunkData.position(data.position() + dataOffset); + chunkData.limit(data.position() + dataOffset + overlapLen); + + Checksum.verifyChecksum(chunkData, checksumData, checksumIndex); + + dataOffset += overlapLen; + + long offsetInChunk = overlapStart - chunkStart; + long endOffsetInChunk = overlapEnd - chunkStart; + + int firstChecksumIndex = Math.toIntExact(offsetInChunk / bytesPerChecksum); + int lastChecksumIndex = Math.toIntExact((endOffsetInChunk - 1) / bytesPerChecksum); + + checksumIndex += (lastChecksumIndex - firstChecksumIndex + 1); + } + currentChunkOffset += chunk.getLen(); + } + } + } } } diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java index cded373e7c81..21c53e2209ba 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java @@ -46,9 +46,13 @@ import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumData; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfoList; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadBlockResponseProto; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.OzoneClientConfig; @@ -773,15 +777,25 @@ private ContainerCommandResponseProto buildResponseProto(byte[] data, long offse private ContainerCommandResponseProto buildCorruptResponseProto(byte[] data, long offset) { return ContainerCommandResponseProto.newBuilder() .setCmdType(Type.ReadBlock) - .setResult(ContainerProtos.Result.SUCCESS) + .setResult(Result.SUCCESS) .setReadBlock(ReadBlockResponseProto.newBuilder() .setOffset(offset) .setData(ByteString.copyFrom(data)) .setChecksumData(ChecksumData.newBuilder() - .setType(ContainerProtos.ChecksumType.CRC32) + .setType(ChecksumType.CRC32) .setBytesPerChecksum(data.length) .addChecksums(ByteString.copyFrom(new byte[4])) .build()) + .setChunkInfoList(ChunkInfoList.newBuilder() + .addChunks(ChunkInfo.newBuilder() + .setChunkName("chunk") + .setOffset(offset) + .setChecksumData(ChecksumData.newBuilder() + .setType(ChecksumType.CRC32) + .setBytesPerChecksum(data.length) + .addChecksums(ByteString.copyFrom(new byte[4])) + .build()) + .setLen(data.length).build()).build()) .build()) .build(); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java index 1164d53c162c..bbe5f03fe76e 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java @@ -32,6 +32,7 @@ class TestOzoneContractFSO extends AbstractOzoneContractTest { @Override protected OzoneConfiguration createOzoneConfig() { OzoneConfiguration conf = super.createOzoneConfig(); + conf.setBoolean("ozone.client.stream.readblock.enable", true); conf.set(OZONE_DEFAULT_BUCKET_LAYOUT, FILE_SYSTEM_OPTIMIZED.name()); return conf; } From 4745196a0e6a904bf4b817de202d6834b78fd7c6 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Wed, 15 Jul 2026 09:25:18 +0800 Subject: [PATCH 9/9] revert TestOzoneContractFSO --- .../apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java | 1 - 1 file changed, 1 deletion(-) diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java index bbe5f03fe76e..1164d53c162c 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/contract/TestOzoneContractFSO.java @@ -32,7 +32,6 @@ class TestOzoneContractFSO extends AbstractOzoneContractTest { @Override protected OzoneConfiguration createOzoneConfig() { OzoneConfiguration conf = super.createOzoneConfig(); - conf.setBoolean("ozone.client.stream.readblock.enable", true); conf.set(OZONE_DEFAULT_BUCKET_LAYOUT, FILE_SYSTEM_OPTIMIZED.name()); return conf; }