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 c15cd338908d..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; @@ -532,10 +534,14 @@ 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 (readBlock.hasChunkInfoList()) { + verifyChecksumForReadBlock(data, checksumData, readBlock); + } else { + throw new IOException("Checksum data is missing for block " + getBlockID()); + } } 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(); @@ -628,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-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..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 @@ -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,14 +338,17 @@ 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 response = ReadBlockResponseProto.newBuilder() .setChecksumData(checksumData.getProtoBufMessage()) .setData(ByteString.copyFrom(data)) .setOffset(offset) + .setChunkInfoList(ChunkInfoList.newBuilder().addAllChunks(chunkInfoList).build()) .build(); + return getSuccessResponseBuilder(request) .setReadBlock(response) .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..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; @@ -1162,4 +1163,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 { 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..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 @@ -301,4 +301,44 @@ 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); + clientConfig.setStreamBufferFlushDelay(false); + 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); + } + } + } + } }