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 d773701300e5..472d3038fd4b 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 @@ -54,7 +54,6 @@ import org.apache.hadoop.hdds.utils.ConnectionFailureUtils; 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.security.token.Token; import org.apache.ratis.protocol.exceptions.TimeoutIOException; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; @@ -232,7 +231,7 @@ private int preadOnce(XceiverClientGrpc client, StreamingReader reader, Pipeline throw new IOException("Uninitialized StreamingReadResponse: " + blockID); } client.streamRead(ContainerProtocolCalls.buildReadBlockCommandProto( - blockID, blockOffset, length, responseDataSize, tokenRef.get(), pipeline), response); + blockID, blockOffset, length, responseDataSize, verifyChecksum, tokenRef.get(), pipeline), response); int copied = 0; while (copied < length) { @@ -475,7 +474,7 @@ synchronized void readBlockImpl(long length) throws IOException { throw new IOException("Uninitialized StreamingReadResponse: " + blockID); } xceiverClient.streamRead(ContainerProtocolCalls.buildReadBlockCommandProto( - blockID, requestedLength, length, responseDataSize, tokenRef.get(), pipelineRef.get()), r); + blockID, requestedLength, length, responseDataSize, verifyChecksum, tokenRef.get(), pipelineRef.get()), r); } private void handleExceptions(IOException cause) throws IOException { @@ -696,8 +695,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman try { ByteBuffer data = readBlock.getData().asReadOnlyByteBuffer(); if (verifyChecksum) { - ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData()); - Checksum.verifyChecksum(data, checksumData, 0); + Checksum.validateChecksums(data, readBlock.getOffset(), 0, readBlock.getChunkInfoListList()); } offerToQueue(readBlock); } catch (Exception e) { @@ -709,7 +707,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman LOG.warn("Failed to process block {} response at offset={}, size={}: {}, {}", getBlockID().getContainerBlockID(), offset, data.size(), StringUtils.bytes2Hex(data.asReadOnlyByteBuffer(), 10), - readBlock.getChecksumData(), e); + readBlock.getChunkInfoListList(), e); if (r != null) { r.getRequestObserver().onError(e); } @@ -783,10 +781,8 @@ private void setCompleted() { private void offerToQueue(ReadBlockResponseProto item) { if (LOG.isTraceEnabled()) { - final ContainerProtos.ChecksumData checksumData = item.getChecksumData(); - LOG.trace("{}: enqueue response offset {}, length {}, numChecksums {}, bytesPerChecksum={}", - name, item.getOffset(), item.getData().size(), - checksumData.getChecksumsList().size(), checksumData.getBytesPerChecksum()); + LOG.trace("{}: enqueue response offset {}, length {}, chunks {}", + name, item.getOffset(), item.getData().size(), item.getChunkInfoListCount()); } final boolean offered = responseQueue.offer(item); Preconditions.assertTrue(offered, () -> "Failed to offer " + item); 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 d162b80ad94b..ef960291f513 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 @@ -68,6 +68,7 @@ import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.security.token.OzoneBlockTokenIdentifier; +import org.apache.hadoop.ozone.common.Checksum; import org.apache.hadoop.ozone.common.OzoneChecksumException; import org.apache.hadoop.security.token.Token; import org.apache.ratis.protocol.exceptions.TimeoutIOException; @@ -76,6 +77,7 @@ import org.apache.ratis.thirdparty.io.grpc.StatusRuntimeException; import org.apache.ratis.thirdparty.io.grpc.stub.ClientCallStreamObserver; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; /** * Tests for StreamBlockInputStream custom configuration behavior. @@ -1138,10 +1140,12 @@ private ReadBlockResponseProto buildReadBlockResponse(byte[] data) { return ReadBlockResponseProto.newBuilder() .setOffset(0) .setData(ByteString.copyFrom(data)) - .setChecksumData(ChecksumData.newBuilder() - .setType(ContainerProtos.ChecksumType.NONE) - .setBytesPerChecksum(data.length) - .build()) + .addChunkInfoList(ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk").setOffset(0).setLen(data.length) + .setChecksumData(ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.NONE) + .setBytesPerChecksum(data.length) + .build())) .build(); } @@ -1152,10 +1156,12 @@ private ContainerCommandResponseProto buildResponseProto(byte[] data, long offse .setReadBlock(ReadBlockResponseProto.newBuilder() .setOffset(offset) .setData(ByteString.copyFrom(data)) - .setChecksumData(ChecksumData.newBuilder() - .setType(ContainerProtos.ChecksumType.NONE) - .setBytesPerChecksum(data.length) - .build()) + .addChunkInfoList(ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk").setOffset(offset).setLen(data.length) + .setChecksumData(ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.NONE) + .setBytesPerChecksum(data.length) + .build())) .build()) .build(); } @@ -1167,12 +1173,96 @@ private ContainerCommandResponseProto buildCorruptResponseProto(byte[] data, lon .setReadBlock(ReadBlockResponseProto.newBuilder() .setOffset(offset) .setData(ByteString.copyFrom(data)) - .setChecksumData(ChecksumData.newBuilder() - .setType(ContainerProtos.ChecksumType.CRC32) - .setBytesPerChecksum(data.length) - .addChecksums(ByteString.copyFrom(new byte[4])) - .build()) + .addChunkInfoList(ContainerProtos.ChunkInfo.newBuilder() + .setChunkName("chunk").setOffset(offset).setLen(data.length) + .setChecksumData(ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32) + .setBytesPerChecksum(data.length) + .addChecksums(ByteString.copyFrom(new byte[4])) + .build())) .build()) .build(); } + + private static ContainerProtos.ChunkInfo checksummedChunk(byte[] data, int offset, int length) throws Exception { + return ContainerProtos.ChunkInfo.newBuilder().setChunkName("chunk-" + offset).setOffset(offset).setLen(length) + .setChecksumData(new Checksum(ContainerProtos.ChecksumType.CRC32, 4) + .computeChecksum(ByteBuffer.wrap(data, offset, length)).getProtoBufMessage()).build(); + } + + private static ContainerCommandResponseProto response(byte[] data, long offset, + List chunks) { + return ContainerCommandResponseProto.newBuilder().setCmdType(Type.ReadBlock) + .setResult(ContainerProtos.Result.SUCCESS).setReadBlock(ReadBlockResponseProto.newBuilder() + .setOffset(offset).setData(ByteString.copyFrom(data)).addAllChunkInfoList(chunks)).build(); + } + + @Test + void testChunkChecksumsAndSeekPrefix() throws Exception { + byte[] data = new byte[] {0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11}; + List chunks = Arrays.asList( + checksummedChunk(data, 0, 3), checksummedChunk(data, 3, 9)); + for (int seek : new int[] {0, 4, 8, 11}) { + int offset = seek == 0 ? 0 : 3 + (seek - 3) / 4 * 4; + ContainerCommandResponseProto response = response(Arrays.copyOfRange(data, offset, data.length), offset, + offset == 0 ? chunks : chunks.subList(1, 2)); + assertStreamResponse(response, data, seek, true, false); + } + } + + @Test + void testMissingMetadataAndVerificationDisabled() throws Exception { + byte[] data = new byte[] {0, 1, 2, 3}; + ContainerCommandResponseProto missing = response(data, 0, Collections.emptyList()); + assertStreamResponse(missing, data, 0, true, true); + assertStreamResponse(missing, data, 0, false, false); + assertStreamResponse(buildResponseProto(data, 0), data, 0, true, false); + byte[] corruptData = data.clone(); + corruptData[0]++; + ContainerCommandResponseProto corrupt = response(corruptData, 0, + Collections.singletonList(checksummedChunk(data, 0, data.length))); + assertStreamResponse(corrupt, data, 0, true, true); + assertStreamResponse(corrupt, corruptData, 0, false, false); + } + + private void assertStreamResponse(ContainerCommandResponseProto response, byte[] expected, int seek, + boolean verifyChecksum, boolean fails) throws Exception { + assertStreamResponse(response, expected, seek, verifyChecksum, fails, false); + assertStreamResponse(response, expected, seek, verifyChecksum, fails, true); + } + + private void assertStreamResponse(ContainerCommandResponseProto response, byte[] expected, int seek, + boolean verifyChecksum, boolean fails, boolean positioned) throws Exception { + OzoneClientConfig config = newStreamReadConfig(); + config.setChecksumVerify(verifyChecksum); + config.setMaxReadRetryCount(0); + ClientCallStreamObserver observer = mock(ClientCallStreamObserver.class); + XceiverClientGrpc client = mockCapturingStreamingReadClient(observer, reader -> reader.onNext(response)); + XceiverClientFactory factory = mock(XceiverClientFactory.class); + when(factory.acquireClientForReadData(any(Pipeline.class))).thenReturn(client); + try (StreamBlockInputStream stream = new StreamBlockInputStream(new BlockID(1, 16258), expected.length, + mockStandalonePipeline(), null, factory, NO_REFRESH, config)) { + stream.seek(seek); + byte[] actual = new byte[expected.length - seek]; + if (fails) { + IOException error = assertThrows(IOException.class, () -> { + if (positioned) { + stream.readPositioned(seek, ByteBuffer.wrap(actual)); + } else { + stream.read(actual); + } + }); + assertThat(hasCause(error, OzoneChecksumException.class)).isTrue(); + assertArrayEquals(new byte[actual.length], actual); + } else { + assertEquals(actual.length, positioned + ? stream.readPositioned(seek, ByteBuffer.wrap(actual)) : stream.read(actual)); + assertArrayEquals(Arrays.copyOfRange(expected, seek, expected.length), actual); + } + } + ArgumentCaptor request = ArgumentCaptor.forClass(ContainerCommandRequestProto.class); + verify(client).streamRead(request.capture(), any()); + assertEquals(verifyChecksum, request.getValue().getReadBlock().getIncludeChecksums()); + verify(client, times(1)).completeStreamRead(); + } } 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 831c899bde57..1ff6009d0c43 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 @@ -44,7 +44,6 @@ 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.datanode.proto.ContainerProtos.WriteChunkResponseProto; -import org.apache.hadoop.ozone.common.ChecksumData; import org.apache.hadoop.ozone.common.ChunkBufferToByteString; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.com.google.protobuf.UnsafeByteOperations; @@ -344,13 +343,14 @@ public static ContainerCommandResponseProto getReadChunkResponse( } public static ContainerCommandResponseProto getReadBlockResponse( - ContainerCommandRequestProto request, ChecksumData checksumData, ByteBuffer data, long offset) { + ContainerCommandRequestProto request, List chunks, ByteBuffer data, long offset) { - ContainerProtos.ReadBlockResponseProto response = ContainerProtos.ReadBlockResponseProto.newBuilder() - .setChecksumData(checksumData.getProtoBufMessage()) + ContainerProtos.ReadBlockResponseProto.Builder response = ContainerProtos.ReadBlockResponseProto.newBuilder() .setData(ByteString.copyFrom(data)) - .setOffset(offset) - .build(); + .setOffset(offset); + if (request.getReadBlock().getIncludeChecksums()) { + response.addAllChunkInfoList(chunks); + } return getSuccessResponseBuilder(request) .setReadBlock(response) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java index b8d832cf45a4..5e0f1b3bc9e6 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java @@ -950,7 +950,7 @@ public static List toValidatorList(Validator validator) { } public static ContainerCommandRequestProto buildReadBlockCommandProto( - BlockID blockID, long offset, long length, int responseDataSize, + BlockID blockID, long offset, long length, int responseDataSize, boolean includeChecksums, Token token, Pipeline pipeline) throws IOException { final DatanodeDetails datanode = pipeline.getClosestNode(); @@ -959,6 +959,7 @@ public static ContainerCommandRequestProto buildReadBlockCommandProto( .setOffset(offset) .setLength(length) .setResponseDataSize(responseDataSize) + .setIncludeChecksums(includeChecksums) .setBlockID(blockID.getDatanodeBlockIDProtobufBuilder(replicaIndex)); final ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto.newBuilder().setCmdType(Type.ReadBlock) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java index 128666784eb4..a1b50a87d793 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java @@ -17,8 +17,6 @@ package org.apache.hadoop.ozone.common; -import static org.apache.ratis.util.Preconditions.assertSame; - import com.google.common.annotations.VisibleForTesting; import java.nio.ByteBuffer; import java.security.MessageDigest; @@ -441,37 +439,37 @@ public static void verifyChecksum(List bufferList, int startIndex, C public static void validateChecksums(ByteBuffer data, long blockOffset, int startIndex, final List chunks) throws OzoneChecksumException { - long firstChunkOffset = blockOffset - chunks.get(startIndex).getOffset(); - int bytesPerChecksum = chunks.get(startIndex).getChecksumData().getBytesPerChecksum(); - long readLength = data.remaining(); - int dataOffset = data.position(); - - if (readLength <= 0) { + if (!data.hasRemaining()) { return; } - - assertSame(0, firstChunkOffset % bytesPerChecksum, "blockOffset % bytesPerChecksum"); - ContainerProtos.ChunkInfo firstChunk = chunks.get(startIndex); - - // verify first chunk, only the first chunk is not start at zero. - int firstChunkIndex = (int) (firstChunkOffset / bytesPerChecksum); - int dataLimit = (int) Math.min(firstChunk.getLen() - firstChunkOffset, readLength); - Checksum.verifySingleChunk(data, dataOffset, dataLimit, firstChunk, firstChunkIndex); - dataOffset += dataLimit; - readLength -= dataLimit; - startIndex++; - - while (readLength > 0) { - ContainerProtos.ChunkInfo chunkInfo = chunks.get(startIndex); - dataLimit = (int) Math.min(chunkInfo.getLen(), readLength); - - Checksum.verifySingleChunk(data, dataOffset, dataLimit, chunkInfo, 0); - - dataOffset += dataLimit; - readLength -= dataLimit; - startIndex++; + int dataOffset = data.position(); + int remaining = data.remaining(); + long offset = blockOffset; + while (remaining > 0) { + if (startIndex >= chunks.size()) { + throw new OzoneChecksumException("Missing chunk metadata at offset " + offset); + } + ChunkInfo chunk = chunks.get(startIndex++); + long chunkOffset = chunk.getOffset(); + long length = chunk.getLen(); + if (offset < chunkOffset || offset - chunkOffset >= length + || (offset != blockOffset && offset != chunkOffset)) { + throw new OzoneChecksumException("Invalid chunk coverage at offset " + offset); + } + long relativeOffset = offset - chunkOffset; + int size = (int) Math.min(remaining, length - relativeOffset); + if (chunk.getChecksumData().getType() != ChecksumType.NONE) { + int bytesPerChecksum = chunk.getChecksumData().getBytesPerChecksum(); + if (bytesPerChecksum <= 0 || relativeOffset % bytesPerChecksum != 0 + || (relativeOffset + size != length && (relativeOffset + size) % bytesPerChecksum != 0)) { + throw new OzoneChecksumException("Invalid checksum boundary at offset " + offset); + } + verifySingleChunk(data, dataOffset, size, chunk, (int) (relativeOffset / bytesPerChecksum)); + } + dataOffset += size; + offset += size; + remaining -= size; } - } private static void verifySingleChunk(ByteBuffer data, int dataOffset, int dataLimit, diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java index 7ccc5d8547e2..57a35c21e570 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java @@ -25,7 +25,9 @@ import static org.mockito.Mockito.when; import java.nio.ByteBuffer; +import java.util.Arrays; import java.util.Collections; +import java.util.List; import org.apache.commons.lang3.RandomStringUtils; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.junit.jupiter.api.Test; @@ -133,4 +135,65 @@ public void testRejectsIncrementalBufferInWriteMode(boolean useChecksumCache) { exception.getMessage()); } } + + private static ContainerProtos.ChunkInfo chunk(byte[] data, int offset, int length, int interval) throws Exception { + return ContainerProtos.ChunkInfo.newBuilder().setChunkName("chunk-" + offset).setOffset(offset).setLen(length) + .setChecksumData(new Checksum(ContainerProtos.ChecksumType.CRC32, interval) + .computeChecksum(ByteBuffer.wrap(data, offset, length)).getProtoBufMessage()).build(); + } + + @Test + void testChunkRelativeRangesPreserveBuffer() throws Exception { + byte[] data = "ABCDEFGHIJKL".getBytes(UTF_8); + List chunks = Arrays.asList(chunk(data, 0, 3, 4), chunk(data, 3, 9, 4)); + for (int start : new int[] {0, 3, 7, 11}) { + ByteBuffer buffer = ByteBuffer.wrap(data, start, data.length - start); + buffer.mark(); + Checksum.validateChecksums(buffer, start, start == 0 ? 0 : 1, chunks); + assertEquals(start, buffer.position()); + assertEquals(data.length, buffer.limit()); + buffer.reset(); + } + Checksum.validateChecksums(ByteBuffer.allocate(0), 0, 0, Collections.emptyList()); + data[8]++; + assertThrows(OzoneChecksumException.class, + () -> Checksum.validateChecksums(ByteBuffer.wrap(data), 0, 0, chunks)); + } + + @Test + void testMalformedChunkCoverageAndAlignment() throws Exception { + byte[] data = "ABCDEFGHIJKL".getBytes(UTF_8); + ContainerProtos.ChunkInfo first = chunk(data, 0, 3, 4); + ContainerProtos.ChunkInfo second = chunk(data, 3, 9, 4); + for (List chunks : Arrays.asList( + Collections.emptyList(), Collections.singletonList(first), + Arrays.asList(second, first), + Arrays.asList(first, second.toBuilder().setOffset(2).build()), + Arrays.asList(first, second.toBuilder().setOffset(4).build()))) { + assertThrows(OzoneChecksumException.class, + () -> Checksum.validateChecksums(ByteBuffer.wrap(data), 0, 0, chunks)); + } + assertThrows(OzoneChecksumException.class, () -> Checksum.validateChecksums( + ByteBuffer.wrap(data), -1, 0, Arrays.asList(first, second))); + for (int[] range : new int[][] {{4, 4}, {3, 3}}) { + assertThrows(OzoneChecksumException.class, () -> Checksum.validateChecksums( + ByteBuffer.wrap(data, range[0], range[1]), range[0], 0, Collections.singletonList(second))); + } + ContainerProtos.ChunkInfo invalid = second.toBuilder().setChecksumData( + second.getChecksumData().toBuilder().setBytesPerChecksum(0)).build(); + assertThrows(OzoneChecksumException.class, () -> Checksum.validateChecksums( + ByteBuffer.wrap(data, 3, 9), 3, 0, Collections.singletonList(invalid))); + } + + @Test + void testNoneAndDifferentChunkParameters() throws Exception { + byte[] data = "ABCDEFGHIJKL".getBytes(UTF_8); + ContainerProtos.ChunkInfo first = chunk(data, 0, 3, 2).toBuilder() + .setChecksumData(Checksum.getNoChecksumDataProto()).build(); + List chunks = Arrays.asList(first, chunk(data, 3, 9, 3)); + Checksum.validateChecksums(ByteBuffer.wrap(data, 1, 11), 1, 0, chunks); + Checksum.validateChecksums(ByteBuffer.wrap(data), 0, 0, + Arrays.asList(chunk(data, 0, 3, 2), chunk(data, 3, 9, 3))); + } + } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java index 34298c7263fb..1564e3ddbf5f 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java @@ -877,7 +877,7 @@ public void streamDataReadOnly(ContainerCommandRequestProto msg, ContainerProtos.Result.CONTAINER_INTERNAL_ERROR); } perf.appendPreOpLatencyMs(Time.monotonicNow() - startTime); - responseProto = handler.readBlock(msg, container, blockFile, streamObserver, false); + responseProto = handler.readBlock(msg, container, blockFile, streamObserver); long oPLatencyMS = Time.monotonicNow() - startTime; metrics.incContainerOpsLatencies(cmdType, oPLatencyMS); if (responseProto == null) { diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java index 53dd03627e13..d6828423a36f 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java @@ -302,5 +302,5 @@ public abstract void copyContainer( public abstract ContainerCommandResponseProto readBlock( ContainerCommandRequestProto msg, Container container, RandomAccessFileChannel blockFile, - StreamObserver streamObserver, boolean testVariableChunks); + StreamObserver streamObserver); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/BlockReadCursor.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/BlockReadCursor.java new file mode 100644 index 000000000000..e877039829ab --- /dev/null +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/BlockReadCursor.java @@ -0,0 +1,125 @@ +/* + * 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.hadoop.ozone.container.keyvalue; + +import java.io.IOException; +import java.util.Arrays; +import java.util.List; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; + +/** Tracks chunk-relative checksum boundaries for a streaming block read. */ +class BlockReadCursor { + private static final int STREAMING_BYTES_PER_CHUNK = 1024 * 64; + + private final List chunks; + // [c1 offset, c2 offset, ..., cn offset, cn offset + cn len] + private final long[] chunkOffsets; + private final int responseDataSize; + private final long start; + private final long end; + private long offset; + private int chunkIndex; + + BlockReadCursor(long requestedOffset, long length, int responseSize, List chunks) throws IOException { + this.chunks = chunks; + chunkOffsets = new long[chunks.size() + 1]; + int bufferSize = responseSize; + for (int i = 0; i < chunks.size(); i++) { + ChunkInfo chunk = chunks.get(i); + chunkOffsets[i] = chunk.getOffset(); + // One checksum interval (or the short chunk containing it) must always fit in the buffer. + bufferSize = Math.max(bufferSize, (int) Math.min(chunk.getLen(), interval(chunk))); + } + ChunkInfo lastChunk = chunks.get(chunks.size() - 1); + long blockEnd = lastChunk.getOffset() + lastChunk.getLen(); + chunkOffsets[chunks.size()] = blockEnd; + this.responseDataSize = bufferSize; + chunkIndex = findChunk(requestedOffset); + start = length == 0 ? requestedOffset : alignDown(requestedOffset, chunkIndex); + offset = start; + if (length == 0) { + end = start; + } else { + long requestedEnd = requestedOffset + Math.min(length, blockEnd - requestedOffset); + int last = findChunk(requestedEnd - 1); + long floor = alignDown(requestedEnd - 1, last); + end = floor + Math.min(interval(chunks.get(last)), chunkOffsets[last + 1] - floor); + } + } + + private static int interval(ChunkInfo chunk) throws IOException { + // Retain the existing read alignment for chunks without checksums. + if (chunk.getChecksumData().getType() == ChecksumType.NONE) { + return STREAMING_BYTES_PER_CHUNK; + } + int size = chunk.getChecksumData().getBytesPerChecksum(); + if (size <= 0) { + throw new IOException("Invalid bytes per checksum: " + size); + } + return size; + } + + private long alignDown(long position, int index) throws IOException { + return position - (position - chunkOffsets[index]) % interval(chunks.get(index)); + } + + private int findChunk(long position) { + int index = Arrays.binarySearch(chunkOffsets, chunkIndex, chunks.size(), position); + if (index >= 0) { + return index; + } + int insertionPoint = -index - 1; + return insertionPoint - 1; + } + + long offset() { + return offset; + } + + boolean hasRemaining() { + return offset < end; + } + + int responseDataSize() { + return responseDataSize; + } + + int nextReadLength() throws IOException { + long limit = offset + Math.min(responseDataSize, end - offset); + if (limit < end) { + limit = alignDown(limit, findChunk(limit)); + } + return Math.toIntExact(limit - offset); + } + + List chunksForRead(int length) { + return chunks.subList(chunkIndex, findChunk(offset + length - 1) + 1); + } + + void advance(int bytesRead) { + offset += bytesRead; + if (hasRemaining()) { + chunkIndex = findChunk(offset); + } + } + + long bytesRead() { + return offset - start; + } +} 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 bab8b03bbf8c..b4342fcdc915 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 @@ -62,12 +62,12 @@ import static org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion.DEFAULT_LAYOUT; import static org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion.FILE_PER_BLOCK; import static org.apache.hadoop.ozone.container.keyvalue.helpers.BlockUtils.getBlockMapKey; -import static org.apache.ratis.util.Preconditions.assertSame; import static org.apache.ratis.util.Preconditions.assertTrue; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; import com.google.common.util.concurrent.Striped; +import java.io.EOFException; import java.io.File; import java.io.FileNotFoundException; import java.io.FilenameFilter; @@ -104,7 +104,6 @@ import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.protocol.DatanodeDetails; 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.ContainerCommandRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State; @@ -132,7 +131,6 @@ import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.client.io.BlockInputStreamFactoryImpl; import org.apache.hadoop.ozone.common.Checksum; -import org.apache.hadoop.ozone.common.ChecksumData; import org.apache.hadoop.ozone.common.ChunkBuffer; import org.apache.hadoop.ozone.common.ChunkBufferToByteString; import org.apache.hadoop.ozone.common.OzoneChecksumException; @@ -188,7 +186,6 @@ public class KeyValueHandler extends Handler { private static final Logger LOG = LoggerFactory.getLogger( KeyValueHandler.class); - private static final int STREAMING_BYTES_PER_CHUNK = 1024 * 64; private final BlockManager blockManager; private final ChunkManager chunkManager; @@ -2284,7 +2281,7 @@ boolean deleteUnreferencedFile(File file) { public ContainerCommandResponseProto readBlock( ContainerCommandRequestProto request, Container kvContainer, RandomAccessFileChannel blockFile, - StreamObserver streamObserver, boolean testVariableChunks) { + StreamObserver streamObserver) { if (kvContainer.getContainerData().getLayoutVersion() != FILE_PER_BLOCK) { return ContainerUtils.logAndReturnError(LOG, @@ -2300,7 +2297,7 @@ public ContainerCommandResponseProto readBlock( } try { final long startTime = Time.monotonicNow(); - final long bytesRead = readBlockImpl(request, blockFile, kvContainer, streamObserver, testVariableChunks); + final long bytesRead = readBlockImpl(request, blockFile, kvContainer, streamObserver); KeyValueContainerData containerData = (KeyValueContainerData) kvContainer .getContainerData(); HddsVolume volume = containerData.getVolume(); @@ -2324,7 +2321,7 @@ public ContainerCommandResponseProto readBlock( } private long readBlockImpl(ContainerCommandRequestProto request, RandomAccessFileChannel blockFile, - Container kvContainer, StreamObserver streamObserver, boolean testVariableChunks) + Container kvContainer, StreamObserver streamObserver) throws IOException { final ReadBlockRequestProto readBlock = request.getReadBlock(); int responseDataSize = readBlock.getResponseDataSize(); @@ -2348,76 +2345,26 @@ private long readBlockImpl(ContainerCommandRequestProto request, RandomAccessFil "Requested offset " + readBlock.getOffset() + " is beyond the end of block " + blockID + " with size " + blockData.getSize())); } - final List chunkInfos = blockData.getChunks(); - final ChecksumType checksumType = chunkInfos.get(0).getChecksumData().getType(); - int bytesPerChecksum = STREAMING_BYTES_PER_CHUNK; - if (checksumType != ContainerProtos.ChecksumType.NONE) { - bytesPerChecksum = chunkInfos.get(0).getChecksumData().getBytesPerChecksum(); - } - - // TODO: Support client-side flag to toggle checksum verification. - // If checksum is disabled, chunk offset adjustment can be skipped. - int chunkIndex = ReadBlockComputation.searchChunk(readBlock.getOffset(), chunkInfos); - ReadBlockComputation readBlockComputation = - new ReadBlockComputation(responseDataSize, bytesPerChecksum, chunkInfos, chunkIndex); - long adjustedOffset = readBlockComputation.computeAdjustedOffset(readBlock.getOffset()); - - long adjustLength = readBlockComputation.computeAdjustedLength( - readBlock.getOffset(), readBlock.getLength(), adjustedOffset); - - ChecksumData checksumData = new ChecksumData(checksumType, bytesPerChecksum); - final ByteBuffer buffer = ByteBuffer.allocate(responseDataSize); - blockFile.position(adjustedOffset); - long totalDataLength = 0; - int numResponses = 0; - Preconditions.checkState(adjustLength <= blockData.getSize() - adjustedOffset); - LOG.debug("adjustedOffset {}, requiredLength {}, blockSize {}", - adjustedOffset, adjustLength, blockData.getSize()); - for (boolean shouldRead = true; totalDataLength < adjustLength && shouldRead;) { - - int bufferLimit = readBlockComputation.computeBufferLimit(adjustedOffset, adjustLength - totalDataLength); - - buffer.limit(bufferLimit); - - shouldRead = blockFile.read(buffer); + final BlockReadCursor cursor = new BlockReadCursor(readBlock.getOffset(), readBlock.getLength(), + responseDataSize, blockData.getChunks()); + final ByteBuffer buffer = ByteBuffer.allocate(cursor.responseDataSize()); + blockFile.position(cursor.offset()); + while (cursor.hasRemaining()) { + buffer.clear().limit(cursor.nextReadLength()); + blockFile.read(buffer); + if (buffer.hasRemaining()) { + throw new EOFException("Unexpected end of block " + blockID + " at " + cursor.offset()); + } buffer.flip(); final int readLength = buffer.remaining(); - if (readLength == 0) { - assertTrue(!shouldRead); - break; - } - assertTrue(readLength > 0, () -> "readLength = " + readLength + " <= 0"); - - if (checksumType != ContainerProtos.ChecksumType.NONE) { - - if (!testVariableChunks) { - // This is the old approach. - // It has some bugs, but we are keeping it for now due to backward compatibility. - checksumData = new ChecksumData(checksumType, bytesPerChecksum, getChecksums(adjustedOffset, readLength, - bytesPerChecksum, chunkInfos)); - } - - - if (validateChunkChecksumData) { - Checksum.validateChecksums(buffer.duplicate(), adjustedOffset, chunkIndex, chunkInfos); - } - LOG.debug("Read {} at adjustedOffset {}, readLength {}, bytesPerChecksum {}", - readBlock, adjustedOffset, readLength, bytesPerChecksum); + final List chunks = cursor.chunksForRead(readLength); + if (validateChunkChecksumData) { + Checksum.validateChecksums(buffer, cursor.offset(), 0, chunks); } - final ContainerCommandResponseProto response = getReadBlockResponse( - request, checksumData, buffer, adjustedOffset); - final int dataLength = response.getReadBlock().getData().size(); - LOG.debug("server onNext response {}: dataLength={}, numChecksums={}", - numResponses, dataLength, response.getReadBlock().getChecksumData().getChecksumsList().size()); - streamObserver.onNext(response); - buffer.clear(); - - adjustedOffset += readLength; - totalDataLength += dataLength; - numResponses++; - chunkIndex = readBlockComputation.findChunk(adjustedOffset); + streamObserver.onNext(getReadBlockResponse(request, chunks, buffer, cursor.offset())); + cursor.advance(readLength); } - return totalDataLength; + return cursor.bytesRead(); } /** @@ -2432,37 +2379,6 @@ private static long rejectReadBlock(RandomAccessFileChannel blockFile, return 0; } - static List getChecksums(long blockOffset, int readLength, int bytesPerChecksum, - final List chunks) { - final int bytesPerChunk = Math.toIntExact(chunks.get(0).getLen()); - 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)); - } - return checksums; - } - @Override public void addFinalizedBlock(Container container, long localID) { KeyValueContainer keyValueContainer = (KeyValueContainer)container; diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/ReadBlockComputation.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/ReadBlockComputation.java deleted file mode 100644 index de5cb5e46fd7..000000000000 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/ReadBlockComputation.java +++ /dev/null @@ -1,159 +0,0 @@ -/* - * 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.hadoop.ozone.container.keyvalue; - -import java.util.List; -import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; -import org.apache.ratis.util.Preconditions; - -/** - * Utility class to compute checksum-aligned offsets, lengths, and buffer - * limits for streaming readBlock operations. - */ -class ReadBlockComputation { - - private final int responseDataSize; - private final int bitMask; - private final int bytesPerChecksum; - private final List chunks; - private int chunkIndex; - private long lastPosition; - - ReadBlockComputation(int responseDataSize, int bytesPerChecksum, - List chunks, int firstChunkIndex) { - this.responseDataSize = responseDataSize; - - Preconditions.assertSame(1, Long.bitCount(bytesPerChecksum), "bitCount"); // check power of 2. - // Suppose bytesPerChecksum = 00001000 (a power of 2), then bitMask = 11111000. - // The following two computations are the same: - // We will use (n & bitMask) to compute ((n / bytesPerChecksum) * bytesPerChecksum). - this.bitMask = -bytesPerChecksum; - this.bytesPerChecksum = bytesPerChecksum; - - this.chunks = chunks; - this.chunkIndex = firstChunkIndex; - this.lastPosition = 0; - } - - /** - * Binary-search for the index of the chunk whose start offset is the - * largest value ≤ {@code targetOffset}. - */ - static int searchChunk(long targetOffset, List chunkInfoList) { - int low = 0; - int high = chunkInfoList.size() - 1; - while (low <= high) { - int mid = (low + high) >>> 1; - long midVal = chunkInfoList.get(mid).getOffset(); - if (midVal <= targetOffset) { - low = mid + 1; - } else { - high = mid - 1; - } - } - return high; - } - - /** - * Return the read offset aligned back to the previous checksum boundary - * within the chunk that contains {@code blockOffset}. - */ - long computeAdjustedOffset(long blockOffset) { - long offsetAlignment = (blockOffset - chunks.get(chunkIndex).getOffset()) % bytesPerChecksum; - return blockOffset - offsetAlignment; - } - - /** - * Return the length of data that must be read so that the range - * {@code [blockOffset, blockOffset + blockLength)} is fully covered - * with checksum-aligned boundaries. - */ - long computeAdjustedLength(long blockOffset, long blockLength, long adjustedOffset) { - // We use an inclusive blockEnd to straightforwardly identify which chunk and checksum - // cover the final byte. - // For example, if bytesPerChecksum = 16 and we read up to byte 36 (exclusive readEnd = 36), - // it correctly locates the chunk. However, calculating the checksum index from an - // exclusive boundary requires extra +1/-1 edge-case handling, specifically when the - // end aligns exactly on a checksum boundary. - // Using the inclusive index makes ((blockEnd - chunkOffset) / bytesPerChecksum) straightforward. - long blockEnd = blockOffset + blockLength - 1; - ChunkInfo lastChunk = chunks.get(searchChunk(blockEnd, chunks)); - long chunkOffset = lastChunk.getOffset(); - - long roundDown = (blockEnd - chunkOffset) & bitMask; - long roundUp = roundDown + bytesPerChecksum; - long chunkLength = Math.min(roundUp, lastChunk.getLen()); - return chunkOffset + chunkLength - adjustedOffset; - } - - /** - * Compute the buffer limit for the next read iteration, aligned to - * checksum boundaries so that partial checksums are not sent across - * response messages. - */ - public int computeBufferLimit(long offset, long remainingLength) { - if (responseDataSize >= remainingLength) { - return Math.toIntExact(remainingLength); - } - - /* - * Use an exclusive boundary here. If we used an inclusive boundary - * that perfectly aligned with the start of the next chunk, it would - * incorrectly truncate the buffer limit. - * - * Example Scenario: - * - bytesPerChecksum = 16 (bitMask = -16) - * - Chunk 1: offset = 0, length = 20 - * - Chunk 2: offset = 20, length = 20 - * - Current offset = 0, responseDataSize = 20 - * - * If using an INCLUSIVE boundary (end = 20 - 1 = 19): - * - endChunk = searchChunk(19) -> Chunk 1 (offset = 0) - * - lengthExcludingEndChunk = 0 - 0 = 0 - * - Result = ((20 - 0) & -16) + 0 = 16 (Incorrectly truncates 4 bytes) - * - * If using an EXCLUSIVE boundary (end = 20): - * - endChunk = searchChunk(20) -> Chunk 2 (offset = 20) - * - lengthExcludingEndChunk = 20 - 0 = 20 - * - Result = ((20 - 20) & -16) + 20 = 20 (Correctly reads the entire Chunk 1) - * - * Use searchChunk, not findChunk: the read may stop before this end, so the findChunk cursor must not move here. - */ - ChunkInfo endChunk = chunks.get(searchChunk(offset + responseDataSize, chunks)); - final int lengthExcludingEndChunk = Math.toIntExact(endChunk.getOffset() - offset); - // bytesPerChecksum must be a power of 2. - return ((responseDataSize - lengthExcludingEndChunk) & bitMask) + lengthExcludingEndChunk; - } - - /** - * @param position must be increasing in subsequent calls to this method. - * @return the chunk containing the given position. - */ - int findChunk(long position) { - Preconditions.assertTrue(position >= lastPosition); - lastPosition = position; - - for (; chunkIndex < chunks.size(); chunkIndex++) { - final ChunkInfo chunk = chunks.get(chunkIndex); - if (position >= chunk.getOffset() && position < chunk.getOffset() + chunk.getLen()) { - return chunkIndex; - } - } - return chunkIndex; - } -} diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestBlockReadCursor.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestBlockReadCursor.java new file mode 100644 index 000000000000..1dda101b8c7f --- /dev/null +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestBlockReadCursor.java @@ -0,0 +1,118 @@ +/* + * 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.hadoop.ozone.container.keyvalue; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; +import org.junit.jupiter.api.Test; + +/** Tests the response ranges selected by {@link BlockReadCursor}. */ +class TestBlockReadCursor { + private static ChunkInfo chunk(long offset, long length, int interval) { + return ChunkInfo.newBuilder().setChunkName("chunk-" + offset).setOffset(offset).setLen(length) + .setChecksumData(ContainerProtos.ChecksumData.newBuilder() + .setType(ContainerProtos.ChecksumType.CRC32).setBytesPerChecksum(interval)).build(); + } + + @Test + void testShiftedRangesAndLookahead() throws Exception { + List chunks = Arrays.asList(chunk(0, 3, 4), chunk(3, 9, 4), chunk(12, 7, 3)); + BlockReadCursor cursor = new BlockReadCursor(4, 14, 6, chunks); + assertEquals(3, cursor.offset()); + assertEquals(4, cursor.nextReadLength()); + assertEquals(4, cursor.nextReadLength()); + assertEquals(chunks.subList(1, 2), cursor.chunksForRead(4)); + cursor.advance(4); + assertEquals(7, cursor.offset()); + assertEquals(5, cursor.nextReadLength()); + cursor.advance(5); + assertEquals(12, cursor.offset()); + assertEquals(6, cursor.nextReadLength()); + cursor.advance(6); + assertFalse(cursor.hasRemaining()); + assertEquals(15, cursor.bytesRead()); + } + + @Test + void testSpanningChunksAndShortFinalChecksum() throws Exception { + List chunks = Arrays.asList(chunk(0, 3, 4), chunk(3, 9, 4)); + BlockReadCursor cursor = new BlockReadCursor(0, Long.MAX_VALUE, 8, chunks); + assertEquals(7, cursor.nextReadLength()); + assertEquals(chunks, cursor.chunksForRead(7)); + cursor.advance(7); + assertEquals(5, cursor.nextReadLength()); + cursor.advance(5); + assertFalse(cursor.hasRemaining()); + assertEquals(12, cursor.bytesRead()); + } + + @Test + void testReadsAtChunkAndChecksumBoundaries() throws Exception { + List chunks = Arrays.asList(chunk(0, 3, 4), chunk(3, 9, 4), chunk(12, 7, 3)); + for (int[] range : new int[][] {{0, 3, 0}, {3, 7, 1}, {7, 11, 1}, {11, 12, 1}, + {12, 15, 2}, {15, 18, 2}, {18, 19, 2}}) { + for (int requestedOffset : new int[] {range[0], range[1] - 1}) { + BlockReadCursor cursor = new BlockReadCursor(requestedOffset, 1, 8, chunks); + int length = range[1] - range[0]; + assertEquals(range[0], cursor.offset()); + assertEquals(length, cursor.nextReadLength()); + assertEquals(Collections.singletonList(chunks.get(range[2])), cursor.chunksForRead(length)); + cursor.advance(length); + assertFalse(cursor.hasRemaining()); + } + } + } + + @Test + void testMinimumBufferAndExactChunkEnd() throws Exception { + List chunks = Arrays.asList(chunk(0, 3, 4), chunk(3, 9, 4)); + BlockReadCursor cursor = new BlockReadCursor(3, 9, 1, chunks); + for (int expected : new int[] {4, 4, 1}) { + assertTrue(cursor.hasRemaining()); + assertEquals(expected, cursor.nextReadLength()); + cursor.advance(expected); + } + assertFalse(cursor.hasRemaining()); + assertFalse(new BlockReadCursor(3, 0, 1, chunks).hasRemaining()); + BlockReadCursor shortChunk = new BlockReadCursor(0, 3, 1, + Collections.singletonList(chunk(0, 3, Integer.MAX_VALUE))); + assertEquals(3, shortChunk.responseDataSize()); + assertEquals(3, shortChunk.nextReadLength()); + } + + @Test + void testRangeNearLongLimit() throws Exception { + List chunks = Arrays.asList(chunk(0, Long.MAX_VALUE - 9, 4), chunk(Long.MAX_VALUE - 9, 9, 4)); + BlockReadCursor cursor = new BlockReadCursor(Long.MAX_VALUE - 8, Long.MAX_VALUE, 4, chunks); + assertEquals(Long.MAX_VALUE - 9, cursor.offset()); + for (int length : new int[] {4, 4, 1}) { + assertEquals(length, cursor.nextReadLength()); + cursor.advance(length); + } + assertEquals(Long.MAX_VALUE, cursor.offset()); + assertFalse(cursor.hasRemaining()); + } + +} 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 c35c003141cb..6a5bf67a66de 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 @@ -71,7 +71,6 @@ import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Semaphore; -import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.atomic.AtomicInteger; import org.apache.commons.io.FileUtils; import org.apache.hadoop.conf.StorageUnit; @@ -133,6 +132,8 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -1125,7 +1126,7 @@ public void onCompleted() { blockFile = new RandomAccessFileChannel(); ContainerCommandResponseProto response = kvHandler.readBlock( - readBlockRequest, container, blockFile, streamObserver, false); + readBlockRequest, container, blockFile, streamObserver); assertNull(response, "ReadBlock should return null on success"); assertTrue(responseCount.get() > 0, "Should receive at least one response"); @@ -1210,6 +1211,64 @@ public void testReadBlockErrorOnTheStreamClosesTheBlockFile() throws Exception { } } + @Test + void testReadBlockResponseSize() throws Exception { + try (StreamFixture fixture = new StreamFixture()) { + fixture.appendChunk("chunk1", 0, BLOCK_SIZE); + for (int responseSize : new int[] {0, 64 * 1024 * 1024}) { + assertResponses(fixture.read(0, BLOCK_SIZE, responseSize), 0, BLOCK_SIZE); + } + } + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testReadBlockWithoutResponseChecksums(boolean validateChecksums) throws Exception { + DatanodeConfiguration datanodeConfig = conf.getObject(DatanodeConfiguration.class); + datanodeConfig.setChunkDataValidationCheck(validateChecksums); + conf.setFromObject(datanodeConfig); + try (StreamFixture fixture = new StreamFixture()) { + long offset = 0; + for (int length : new int[] {3, 9}) { + ChunkInfo chunk = new ChunkInfo("chunk-" + offset, offset, length); + chunk.setChecksumData(new Checksum(ContainerProtos.ChecksumType.CRC32, 4) + .computeChecksum(ByteBuffer.wrap(writtenBytes(offset, length)))); + fixture.appendChunk(chunk); + offset += length; + } + ContainerProtos.ReadBlockRequestProto.Builder request = ContainerProtos.ReadBlockRequestProto.newBuilder() + .setOffset(4).setLength(8).setResponseDataSize(4); + ReadBlockResult withChecksums = fixture.read(request); + assertResponses(withChecksums, 3, 9); + assertResponses(fixture.read(request.setIncludeChecksums(true)), 3, 9); + ReadBlockResult withoutChecksums = fixture.read(request.setIncludeChecksums(false)); + assertNull(withoutChecksums.getResponse()); + assertThat(withoutChecksums.getErrors()).isEmpty(); + assertEquals(withChecksums.getDataResponses().size(), withoutChecksums.getDataResponses().size()); + for (int i = 0; i < withChecksums.getDataResponses().size(); i++) { + ContainerProtos.ReadBlockResponseProto expected = withChecksums.getDataResponses().get(i).getReadBlock(); + assertEquals(expected.toBuilder().clearChunkInfoList().build(), + withoutChecksums.getDataResponses().get(i).getReadBlock()); + } + + File file = ContainerLayoutVersion.FILE_PER_BLOCK.getChunkFile( + fixture.container.getContainerData(), fixture.blockID, "unused"); + try (java.io.RandomAccessFile corrupt = new java.io.RandomAccessFile(file, "rw")) { + corrupt.seek(3); + corrupt.write(255); + } + ReadBlockResult corrupt = fixture.read(request); + if (validateChecksums) { + assertEquals(ContainerProtos.Result.IO_EXCEPTION, corrupt.getResponse().getResult()); + assertThat(corrupt.getDataResponses()).isEmpty(); + } else { + assertNull(corrupt.getResponse()); + assertThat(corrupt.getErrors()).isEmpty(); + assertEquals((byte) 255, corrupt.getDataResponses().get(0).getReadBlock().getData().byteAt(0)); + } + } + } + /** * A chunk boundary inside the response window that is not checksum aligned makes a response shorter than * responseDataSize; the next response must continue from where the short one stopped. @@ -1238,7 +1297,7 @@ private static byte[] writtenBytes(long offset, int length) { } /** Asserts a successful read of {@code totalLength} contiguous bytes from {@code firstOffset}. */ - private static void assertResponses(ReadBlockResult result, long firstOffset, long totalLength) { + private static void assertResponses(ReadBlockResult result, long firstOffset, long totalLength) throws Exception { assertNull(result.getResponse()); assertThat(result.getErrors()).isEmpty(); long offset = firstOffset; @@ -1247,6 +1306,8 @@ private static void assertResponses(ReadBlockResult result, long firstOffset, lo assertEquals(offset, response.getReadBlock().getOffset()); final ByteString data = response.getReadBlock().getData(); assertArrayEquals(writtenBytes(offset, data.size()), data.toByteArray()); + Checksum.validateChecksums(data.asReadOnlyByteBuffer(), offset, 0, + response.getReadBlock().getChunkInfoListList()); offset += data.size(); } assertEquals(totalLength, offset - firstOffset, "total bytes delivered"); @@ -1291,9 +1352,13 @@ private final class StreamFixture implements AutoCloseable { /** Write one more chunk and re-put the block with it. */ void appendChunk(String name, long offset, int length) throws Exception { - ChunkInfo chunkInfo = new ChunkInfo(name, offset, length); + appendChunk(new ChunkInfo(name, offset, length)); + } + + void appendChunk(ChunkInfo chunkInfo) throws Exception { blockData.addChunk(chunkInfo.getProtoBufMessage()); - ChunkBuffer data = ChunkBuffer.wrap(ByteBuffer.wrap(writtenBytes(offset, length))); + ChunkBuffer data = ChunkBuffer.wrap(ByteBuffer.wrap( + writtenBytes(chunkInfo.getOffset(), (int) chunkInfo.getLen()))); kvHandler.getChunkManager().writeChunk(container, blockID, chunkInfo, data, DispatcherContext.getHandleWriteChunk()); kvHandler.getBlockManager().putBlock(container, blockData); @@ -1305,19 +1370,19 @@ ReadBlockResult read(long offset, long length) { } ReadBlockResult read(long offset, long length, int responseDataSize) { + return read(ContainerProtos.ReadBlockRequestProto.newBuilder() + .setOffset(offset).setLength(length).setResponseDataSize(responseDataSize)); + } + + ReadBlockResult read(ContainerProtos.ReadBlockRequestProto.Builder readBlock) { ContainerCommandRequestProto request = ContainerCommandRequestProto.newBuilder() .setCmdType(ContainerProtos.Type.ReadBlock) .setContainerID(container.getContainerData().getContainerID()) .setDatanodeUuid(DATANODE_UUID) - .setReadBlock(ContainerProtos.ReadBlockRequestProto.newBuilder() - .setBlockID(blockID.getDatanodeBlockIDProtobuf()) - .setOffset(offset) - .setLength(length) - .setResponseDataSize(responseDataSize) - .build()) + .setReadBlock(readBlock.setBlockID(blockID.getDatanodeBlockIDProtobuf())) .build(); ReadBlockResult result = new ReadBlockResult(blockID); - result.setResponse(kvHandler.readBlock(request, container, blockFile, result, false)); + result.setResponse(kvHandler.readBlock(request, container, blockFile, result)); return result; } @@ -1443,7 +1508,7 @@ public void testReadBlockWithSmallChunks() throws Exception { } byte[] rawData = new byte[totalLen]; - ThreadLocalRandom.current().nextBytes(rawData); + new java.util.Random(16258).nextBytes(rawData); writeBlock(handlerWithVolume, container, blockID, chunkLens, rawData); readBlockAndVerify(handlerWithVolume, container, blockID, rawData, 2, 2, 2, 2); } finally { @@ -1471,7 +1536,7 @@ public void testReadBlockWithChunksNotMultipleOfBytesPerChecksum() throws Except } byte[] rawData = new byte[totalLen]; - ThreadLocalRandom.current().nextBytes(rawData); + new java.util.Random(16258).nextBytes(rawData); writeBlock(handlerWithVolume, container, blockID, chunkLens, rawData); readBlockAndVerify(handlerWithVolume, container, blockID, rawData, 2048, 2048, 1044, 3072); } finally { @@ -1480,8 +1545,12 @@ public void testReadBlockWithChunksNotMultipleOfBytesPerChecksum() throws Except } } - @Test - public void testReadBlockMultipleResponse() throws Exception { + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void testReadBlockMultipleResponse(boolean validateChecksums) throws Exception { + DatanodeConfiguration datanodeConfig = conf.getObject(DatanodeConfiguration.class); + datanodeConfig.setChunkDataValidationCheck(validateChecksums); + conf.setFromObject(datanodeConfig); int[] chunkLens = {(1 << 20) + 10, 20, 4096}; Path testDir = Files.createTempDirectory("testReadBlock"); try { @@ -1499,7 +1568,7 @@ public void testReadBlockMultipleResponse() throws Exception { } byte[] rawData = new byte[totalLen]; - ThreadLocalRandom.current().nextBytes(rawData); + new java.util.Random(16258).nextBytes(rawData); writeBlock(handlerWithVolume, container, blockID, chunkLens, rawData); readBlockAndVerify(handlerWithVolume, container, blockID, rawData, 20, (1 << 20) + 20, 0, (1 << 20) + 30 + 1024); } finally { @@ -1556,7 +1625,8 @@ private void writeBlock(HandlerWithVolumeSet handlerWithVolume, KeyValueContaine */ @SuppressWarnings("checkstyle:ParameterNumber") private void readBlockAndVerify(HandlerWithVolumeSet handlerWithVolume, KeyValueContainer container, - BlockID blockID, byte[] rawData, long readOffset, long length, long readBackOffset, long readBackLength) { + BlockID blockID, byte[] rawData, long readOffset, long length, long readBackOffset, long readBackLength) + throws Exception { ContainerCommandRequestProto readBlockRequest = ContainerCommandRequestProto.newBuilder() .setCmdType(ContainerProtos.Type.ReadBlock) @@ -1579,7 +1649,7 @@ private void readBlockAndVerify(HandlerWithVolumeSet handlerWithVolume, KeyValue try (RandomAccessFileChannel blockFile = new RandomAccessFileChannel()) { ContainerCommandResponseProto response = handlerWithVolume.getHandler().readBlock( - readBlockRequest, container, blockFile, streamObserver, true); + readBlockRequest, container, blockFile, streamObserver); assertNull(response, "ReadBlock should return null on success"); } @@ -1591,12 +1661,22 @@ private void readBlockAndVerify(HandlerWithVolumeSet handlerWithVolume, KeyValue for (ContainerCommandResponseProto resp : capturedResponses) { returnedDataLen += resp.getReadBlock().getData().size(); } + List storedChunks = handlerWithVolume.getHandler().getBlockManager() + .getBlock(container, blockID).getChunks(); ByteBuffer allData = ByteBuffer.allocate(returnedDataLen); long firstResponseOffset = capturedResponses.get(0).getReadBlock().getOffset(); for (ContainerCommandResponseProto resp : capturedResponses) { assertEquals(ContainerProtos.Result.SUCCESS, resp.getResult()); assertTrue(resp.hasReadBlock()); + long responseOffset = resp.getReadBlock().getOffset(); + long responseEnd = responseOffset + resp.getReadBlock().getData().size(); + List overlapping = storedChunks.stream() + .filter(chunk -> chunk.getOffset() < responseEnd && chunk.getOffset() + chunk.getLen() > responseOffset) + .collect(java.util.stream.Collectors.toList()); + assertEquals(overlapping, resp.getReadBlock().getChunkInfoListList()); + Checksum.validateChecksums(resp.getReadBlock().getData().asReadOnlyByteBuffer(), + resp.getReadBlock().getOffset(), 0, resp.getReadBlock().getChunkInfoListList()); allData.put(resp.getReadBlock().getData().asReadOnlyByteBuffer()); } @@ -1613,4 +1693,28 @@ private void readBlockAndVerify(HandlerWithVolumeSet handlerWithVolume, KeyValue } } + @Test + void testReadBlockZeroLengthAndOversizedRange() throws Exception { + try (StreamFixture fixture = new StreamFixture()) { + fixture.appendChunk("chunk1", 0, BLOCK_SIZE); + assertResponses(fixture.read(1, 0), 1, 0); + assertResponses(fixture.read(1, Long.MAX_VALUE), 0, BLOCK_SIZE); + } + } + + @Test + void testReadBlockUnexpectedEof() throws Exception { + try (StreamFixture fixture = new StreamFixture()) { + fixture.appendChunk("chunk1", 0, BLOCK_SIZE); + File file = ContainerLayoutVersion.FILE_PER_BLOCK.getChunkFile( + fixture.container.getContainerData(), fixture.blockID, "unused"); + try (java.io.RandomAccessFile truncated = new java.io.RandomAccessFile(file, "rw")) { + truncated.setLength(BLOCK_SIZE - 1); + } + ReadBlockResult result = fixture.read(0, BLOCK_SIZE); + assertEquals(ContainerProtos.Result.IO_EXCEPTION, result.getResponse().getResult()); + assertThat(result.getDataResponses()).isEmpty(); + } + } + } diff --git a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto index 72d50c975a52..98c47c07c3ab 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -395,12 +395,15 @@ message ReadBlockRequestProto { required uint64 offset = 2; optional uint64 length = 3; optional uint32 responseDataSize = 4; + optional bool includeChecksums = 5 [default = true]; } message ReadBlockResponseProto { - required ChecksumData checksumData = 1; + reserved 1; + reserved "checksumData"; required uint64 offset = 2; required bytes data = 3; + repeated ChunkInfo 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 4adb0567e134..3ce0af445277 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 @@ -19,6 +19,7 @@ import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.ONE; import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_CLIENT_BYTES_PER_CHECKSUM_MIN_SIZE; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -38,6 +39,7 @@ import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; import org.apache.hadoop.hdds.scm.OzoneClientConfig; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.scm.storage.StreamBlockInputStream; @@ -51,8 +53,12 @@ import org.apache.hadoop.ozone.client.OzoneClientFactory; import org.apache.hadoop.ozone.client.io.KeyInputStream; import org.apache.hadoop.ozone.client.protocol.ClientProtocol; +import org.apache.hadoop.ozone.container.common.helpers.BlockData; import org.apache.hadoop.ozone.container.common.impl.ContainerData; import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion; +import org.apache.hadoop.ozone.container.common.interfaces.DBHandle; +import org.apache.hadoop.ozone.container.keyvalue.KeyValueContainerData; +import org.apache.hadoop.ozone.container.keyvalue.helpers.BlockUtils; import org.apache.hadoop.ozone.om.BucketForTesting; import org.apache.hadoop.ozone.om.helpers.OmKeyInfo; import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo; @@ -145,6 +151,63 @@ void testReadKey256k() throws Exception { runTestReadKey(KEY_SIZE, bytesPerChecksum); } + @Test + void testUnevenFlushedChunks() throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(cluster.getConf()); + OzoneClientConfig config = conf.getObject(OzoneClientConfig.class); + config.setStreamReadBlock(true); + config.setChecksumVerify(true); + config.setBytesPerChecksum(16 * 1024); + config.setStreamReadResponseDataSize(32 * 1024); + config.setStreamBufferFlushDelay(false); + conf.setFromObject(config); + byte[] expected = new byte[8193 + 32771 + 65543]; + new java.util.Random(16258).nextBytes(expected); + try (OzoneClient client = OzoneClientFactory.getRpcClient(conf)) { + BucketForTesting bucket = BucketForTesting.newBuilder(client).build(); + String key = "uneven-flushed-chunks"; + try (OutputStream out = bucket.delegate().createKey(key, expected.length, + RatisReplicationConfig.getInstance(ONE), Collections.emptyMap())) { + int offset = 0; + for (int length : new int[] {8193, 32771, 65543}) { + out.write(expected, offset, length); + out.flush(); + offset += length; + } + } + OmKeyInfo info = client.getProxy().getKeyInfo(bucket.delegate().getVolumeName(), + bucket.delegate().getName(), key, false); + List locations = info.getLatestVersionLocations().createLocationList(); + assertEquals(1, locations.size()); + BlockID blockID = locations.get(0).getBlockID(); + KeyValueContainerData container = (KeyValueContainerData) cluster.getHddsDatanodes().get(0) + .getDatanodeStateMachine().getContainer().getContainerSet() + .getContainer(blockID.getContainerID()).getContainerData(); + List chunks; + try (DBHandle db = BlockUtils.getDB(container, conf)) { + BlockData block = db.getStore().getBlockByID(blockID, container.getBlockKey(blockID.getLocalID())); + chunks = block.getChunks(); + } + assertTrue(chunks.size() >= 3, "flush must persist separate chunks"); + assertTrue(chunks.stream().map(ChunkInfo::getLen).distinct().count() > 1, "stored chunks must be nonuniform"); + long shiftedOffset = chunks.get(1).getOffset(); + assertTrue(shiftedOffset % config.getBytesPerChecksum() != 0, "second chunk must shift the checksum grid"); + try (KeyInputStream in = bucket.getKeyInputStream(key)) { + assertTrue(in.isStreamBlockInputStream()); + byte[] actual = new byte[expected.length]; + org.apache.hadoop.io.IOUtils.readFully(in, actual, 0, actual.length); + assertArrayEquals(expected, actual); + for (long position : new long[] {shiftedOffset + 1, shiftedOffset + 16385, expected.length - 17, 0}) { + in.seek(position); + int length = Math.min(20000, expected.length - (int) position); + byte[] range = new byte[length]; + org.apache.hadoop.io.IOUtils.readFully(in, range, 0, range.length); + assertArrayEquals(Arrays.copyOfRange(expected, (int) position, (int) position + length), range); + } + } + } + } + void runTestReadKey(SizeInBytes keySize, SizeInBytes bytesPerChecksum) throws Exception { final List datanodes = cluster.getHddsDatanodes(); assertEquals(1, datanodes.size());