diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockInputStream.java index 92b03bd2acc5..6d5ac9c8cf21 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockInputStream.java @@ -35,6 +35,7 @@ import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChunkInfo; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.GetBlockRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.GetBlockResponseProto; import org.apache.hadoop.hdds.scm.OzoneClientConfig; import org.apache.hadoop.hdds.scm.XceiverClientFactory; @@ -268,24 +269,15 @@ protected BlockData getBlockDataUsingSCClient() throws IOException { blockID.getContainerID()); } - DatanodeBlockID.Builder blkIDBuilder = - DatanodeBlockID.newBuilder().setContainerID(blockID.getContainerID()) - .setLocalID(blockID.getLocalID()) - .setBlockCommitSequenceId(blockID.getBlockCommitSequenceId()); - int replicaIndex = pipeline.getReplicaIndex(xceiverClientShortCircuit.getDn()); - if (replicaIndex > 0) { - blkIDBuilder.setReplicaIndex(replicaIndex); - } - DatanodeBlockID datanodeBlockID = blkIDBuilder.build(); - ContainerProtos.GetBlockRequestProto.Builder readBlockRequest = - ContainerProtos.GetBlockRequestProto.newBuilder().setBlockID(datanodeBlockID) - .setRequestShortCircuitAccess(true); + final GetBlockRequestProto.Builder getBlockRequest = GetBlockRequestProto.newBuilder() + .setBlockID(blockID.getDatanodeBlockIDProtobufBuilder(replicaIndex)) + .setRequestShortCircuitAccess(true); ContainerProtos.ContainerCommandRequestProto.Builder builder = ContainerProtos.ContainerCommandRequestProto.newBuilder() .setCmdType(ContainerProtos.Type.GetBlock) - .setContainerID(datanodeBlockID.getContainerID()) - .setGetBlock(readBlockRequest) + .setContainerID(blockID.getContainerID()) + .setGetBlock(getBlockRequest) .setClientId(xceiverClientShortCircuit.getClientId()) .setCallId(xceiverClientShortCircuit.getCallId()); if (tokenRef.get() != null) { @@ -294,13 +286,12 @@ protected BlockData getBlockDataUsingSCClient() throws IOException { GetBlockResponseProto response = ContainerProtocolCalls.getBlock(xceiverClientShortCircuit, VALIDATORS, builder, xceiverClientShortCircuit.getDn()); - blockFileInputStream = xceiverClientShortCircuit.getFileInputStream( - builder.getCallId(), datanodeBlockID.getLocalID()); + blockFileInputStream = xceiverClientShortCircuit.getFileInputStream(builder.getCallId(), blockID.getLocalID()); if (blockFileInputStream == null) { - throw new IOException("Failed to get file InputStream for block " + datanodeBlockID); + throw new IOException("Failed to get file InputStream for block " + blockID); } else { if (LOG.isDebugEnabled()) { - LOG.debug("Get the FileInputStream of block {}", datanodeBlockID); + LOG.debug("Get the FileInputStream of block {}", blockID); } } return response.getBlockData(); @@ -316,7 +307,7 @@ protected BlockData getBlockDataUsingGRPCClient() throws IOException { } GetBlockResponseProto response = ContainerProtocolCalls.getBlock( - xceiverClientGrpc, VALIDATORS, blockID, tokenRef.get(), pipeline.getReplicaIndexes()); + xceiverClientGrpc, VALIDATORS, blockID, tokenRef.get(), pipeline); return response.getBlockData(); } diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java index b960e7537443..438425a3ceb3 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java @@ -187,16 +187,9 @@ public BlockOutputStream( KeyValue keyValue = KeyValue.newBuilder().setKey("TYPE").setValue("KEY").build(); - ContainerProtos.DatanodeBlockID.Builder blkIDBuilder = - ContainerProtos.DatanodeBlockID.newBuilder() - .setContainerID(blockID.getContainerID()) - .setLocalID(blockID.getLocalID()) - .setBlockCommitSequenceId(blockID.getBlockCommitSequenceId()); - if (replicationIndex > 0) { - blkIDBuilder.setReplicaIndex(replicationIndex); - } - this.containerBlockData = BlockData.newBuilder().setBlockID( - blkIDBuilder.build()).addMetadata(keyValue); + this.containerBlockData = BlockData.newBuilder() + .setBlockID(blockID.getDatanodeBlockIDProtobufBuilder(replicationIndex)) + .addMetadata(keyValue); this.pipeline = pipeline; // tell DataNode I will send incremental chunk list this.supportIncrementalChunkList = canEnableIncrementalChunkList(); diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ChunkInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ChunkInputStream.java index 34ef7a71bd3f..ee70a67dd858 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ChunkInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ChunkInputStream.java @@ -297,11 +297,7 @@ protected synchronized void releaseClient() { private void updateDatanodeBlockId(Pipeline pipeline) throws IOException { DatanodeDetails closestNode = pipeline.getClosestNode(); int replicaIdx = pipeline.getReplicaIndex(closestNode); - ContainerProtos.DatanodeBlockID.Builder builder = blockID.getDatanodeBlockIDProtobufBuilder(); - if (replicaIdx > 0) { - builder.setReplicaIndex(replicaIdx); - } - datanodeBlockID = builder.build(); + datanodeBlockID = blockID.getDatanodeBlockIDProtobufBuilder(replicaIdx).build(); } /** diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java index 7141a65306dd..81aefb54a0cf 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java @@ -19,7 +19,7 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import java.util.Objects; -import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; /** @@ -98,24 +98,24 @@ public void appendTo(StringBuilder sb) { } @JsonIgnore - public ContainerProtos.DatanodeBlockID getDatanodeBlockIDProtobuf() { - ContainerProtos.DatanodeBlockID.Builder blockID = getDatanodeBlockIDProtobufBuilder(); - if (replicaIndex != null) { - blockID.setReplicaIndex(replicaIndex); - } - return blockID.build(); + public DatanodeBlockID getDatanodeBlockIDProtobuf() { + return getDatanodeBlockIDProtobufBuilder(replicaIndex).build(); } @JsonIgnore - public ContainerProtos.DatanodeBlockID.Builder getDatanodeBlockIDProtobufBuilder() { - return ContainerProtos.DatanodeBlockID.newBuilder(). - setContainerID(containerBlockID.getContainerID()) + public DatanodeBlockID.Builder getDatanodeBlockIDProtobufBuilder(Integer replicaIdx) { + final DatanodeBlockID.Builder b = DatanodeBlockID.newBuilder() + .setContainerID(containerBlockID.getContainerID()) .setLocalID(containerBlockID.getLocalID()) .setBlockCommitSequenceId(blockCommitSequenceId); + if (replicaIdx != null && replicaIdx > 0) { + b.setReplicaIndex(replicaIdx); + } + return b; } @JsonIgnore - public static BlockID getFromProtobuf(ContainerProtos.DatanodeBlockID blockID) { + public static BlockID getFromProtobuf(DatanodeBlockID blockID) { return new BlockID(blockID.getContainerID(), blockID.getLocalID(), blockID.getBlockCommitSequenceId(), diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java index f6065a578e29..8712473ee85c 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java @@ -33,8 +33,6 @@ import java.util.Objects; import java.util.Optional; import java.util.Set; -import java.util.function.Function; -import java.util.stream.Collectors; import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicatedReplicationConfig; @@ -240,10 +238,27 @@ public int getReplicaIndex(DatanodeDetails dn) { } /** - * Get the replicaIndex Map. + * @param fromIndex the replica index starting from (inclusive) + * @param toIndex the replica index starting to (exclusive) + * @return true iff this pipeline contain all replica indexes within the given range. */ - public Map getReplicaIndexes() { - return this.getNodes().stream().collect(Collectors.toMap(Function.identity(), this::getReplicaIndex)); + public boolean containsAllReplicaIndexes(int fromIndex, int toIndex) { + final boolean[] contains = new boolean[toIndex - fromIndex]; + for (int r : replicaIndexes.values()) { + if (r >= fromIndex && r < toIndex) { + contains[r - fromIndex] = true; + } + } + for (boolean contain : contains) { + if (!contain) { + return false; + } + } + return true; + } + + public Map getReplicaIndexesForTesting() { + return Collections.unmodifiableMap(replicaIndexes); } /** 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 d2898c395d94..b8d832cf45a4 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 @@ -188,7 +188,7 @@ static T tryEachDatanode(Pipeline pipeline, */ public static GetBlockResponseProto getBlock(XceiverClientSpi xceiverClient, List validators, BlockID blockID, Token token, - Map replicaIndexes) throws IOException { + Pipeline pipeline) throws IOException { ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto .newBuilder() .setCmdType(Type.GetBlock) @@ -198,7 +198,7 @@ public static GetBlockResponseProto getBlock(XceiverClientSpi xceiverClient, } return tryEachDatanode(xceiverClient.getPipeline(), - d -> getBlock(xceiverClient, validators, builder, blockID, d, replicaIndexes), + d -> getBlock(xceiverClient, validators, builder, blockID, d, pipeline), d -> toErrorMessage(blockID, d)); } @@ -209,8 +209,8 @@ static String toErrorMessage(BlockID blockId, DatanodeDetails d) { public static GetBlockResponseProto getBlock(XceiverClientSpi xceiverClient, BlockID datanodeBlockID, - Token token, Map replicaIndexes) throws IOException { - return getBlock(xceiverClient, getValidatorList(), datanodeBlockID, token, replicaIndexes); + Token token, Pipeline pipeline) throws IOException { + return getBlock(xceiverClient, getValidatorList(), datanodeBlockID, token, pipeline); } /** @@ -221,7 +221,6 @@ public static GetBlockResponseProto getBlock(XceiverClientSpi xceiverClient, * @param blockID blockID to identify container * @param token a token for this block (may be null) * @param datanode datanode to query - * @param replicaIndexes replica indexes for EC pipelines * @return container protocol get block response * @throws IOException if there is an I/O error while performing the call */ @@ -230,7 +229,7 @@ public static GetBlockResponseProto getBlockFromDatanode( BlockID blockID, Token token, DatanodeDetails datanode, - Map replicaIndexes) throws IOException { + Pipeline pipeline) throws IOException { ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto .newBuilder() .setCmdType(Type.GetBlock) @@ -238,23 +237,19 @@ public static GetBlockResponseProto getBlockFromDatanode( if (token != null) { builder.setEncodedToken(token.encodeToUrlString()); } - return getBlock(xceiverClient, getValidatorList(), builder, blockID, datanode, - replicaIndexes); + return getBlock(xceiverClient, getValidatorList(), builder, blockID, datanode, pipeline); } private static GetBlockResponseProto getBlock(XceiverClientSpi xceiverClient, List validators, ContainerCommandRequestProto.Builder builder, BlockID blockID, - DatanodeDetails datanode, Map replicaIndexes) throws IOException { + DatanodeDetails datanode, Pipeline pipeline) throws IOException { String traceId = TracingUtil.exportCurrentSpan(); if (traceId != null) { builder.setTraceID(traceId); } - final DatanodeBlockID.Builder datanodeBlockID = blockID.getDatanodeBlockIDProtobufBuilder(); - int replicaIndex = replicaIndexes.getOrDefault(datanode, 0); - if (replicaIndex > 0) { - datanodeBlockID.setReplicaIndex(replicaIndex); - } + final int replicaIndex = pipeline.getReplicaIndex(datanode); + final DatanodeBlockID.Builder datanodeBlockID = blockID.getDatanodeBlockIDProtobufBuilder(replicaIndex); final GetBlockRequestProto.Builder readBlockRequest = GetBlockRequestProto.newBuilder() .setBlockID(datanodeBlockID.build()); final ContainerCommandRequestProto request = builder @@ -510,16 +505,10 @@ public static XceiverClientReply writeChunkAsync( boolean containerAutoCreate) throws IOException, ExecutionException, InterruptedException { - WriteChunkRequestProto.Builder writeChunkRequest = - WriteChunkRequestProto.newBuilder() - .setBlockID(DatanodeBlockID.newBuilder() - .setContainerID(blockID.getContainerID()) - .setLocalID(blockID.getLocalID()) - .setBlockCommitSequenceId(blockID.getBlockCommitSequenceId()) - .setReplicaIndex(replicationIndex) - .build()) - .setChunkData(chunk) - .setData(data); + final WriteChunkRequestProto.Builder writeChunkRequest = WriteChunkRequestProto.newBuilder() + .setBlockID(blockID.getDatanodeBlockIDProtobufBuilder(replicationIndex)) + .setChunkData(chunk) + .setData(data); if (blockData != null) { PutBlockRequestProto.Builder createBlockRequest = PutBlockRequestProto.newBuilder() @@ -931,41 +920,6 @@ public static List toValidatorList(Validator validator) { return Collections.unmodifiableList(validators); } - public static HashMap - getBlockFromAllNodes( - XceiverClientSpi xceiverClient, - DatanodeBlockID datanodeBlockID, - Token token) - throws IOException, InterruptedException { - GetBlockRequestProto.Builder readBlockRequest = GetBlockRequestProto - .newBuilder() - .setBlockID(datanodeBlockID); - HashMap datanodeToResponseMap - = new HashMap<>(); - String id = xceiverClient.getPipeline().getFirstNode().getUuidString(); - ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto - .newBuilder() - .setCmdType(Type.GetBlock) - .setContainerID(datanodeBlockID.getContainerID()) - .setDatanodeUuid(id) - .setGetBlock(readBlockRequest); - if (token != null) { - builder.setEncodedToken(token.encodeToUrlString()); - } - String traceId = TracingUtil.exportCurrentSpan(); - if (traceId != null) { - builder.setTraceID(traceId); - } - ContainerCommandRequestProto request = builder.build(); - Map responses = - xceiverClient.sendCommandOnAllNodes(request); - for (Map.Entry entry: - responses.entrySet()) { - datanodeToResponseMap.put(entry.getKey(), entry.getValue().getGetBlock()); - } - return datanodeToResponseMap; - } - public static HashMap readContainerFromAllNodes(XceiverClientSpi client, long containerID, String encodedToken) throws IOException, InterruptedException { @@ -1000,12 +954,12 @@ public static ContainerCommandRequestProto buildReadBlockCommandProto( Token token, Pipeline pipeline) throws IOException { final DatanodeDetails datanode = pipeline.getClosestNode(); - final DatanodeBlockID datanodeBlockID = getDatanodeBlockID(blockID, datanode, pipeline.getReplicaIndexes()); + final int replicaIndex = pipeline.getReplicaIndex(datanode); final ReadBlockRequestProto.Builder readBlockRequest = ReadBlockRequestProto.newBuilder() .setOffset(offset) .setLength(length) .setResponseDataSize(responseDataSize) - .setBlockID(datanodeBlockID); + .setBlockID(blockID.getDatanodeBlockIDProtobufBuilder(replicaIndex)); final ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto.newBuilder().setCmdType(Type.ReadBlock) .setContainerID(blockID.getContainerID()); @@ -1017,14 +971,4 @@ public static ContainerCommandRequestProto buildReadBlockCommandProto( .setReadBlock(readBlockRequest) .build(); } - - static DatanodeBlockID getDatanodeBlockID(BlockID blockID, DatanodeDetails datanode, - Map replicaIndexes) { - final DatanodeBlockID.Builder b = blockID.getDatanodeBlockIDProtobufBuilder(); - final int replicaIndex = replicaIndexes.getOrDefault(datanode, 0); - if (replicaIndex > 0) { - b.setReplicaIndex(replicaIndex); - } - return b.build(); - } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestContainerReconciliationWithMockDatanodes.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestContainerReconciliationWithMockDatanodes.java index c051b3478b4c..a97a3418cc87 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestContainerReconciliationWithMockDatanodes.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/TestContainerReconciliationWithMockDatanodes.java @@ -34,7 +34,6 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; -import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.spy; @@ -390,7 +389,7 @@ private static void mockContainerProtocolCalls(FailureLocation failureLocation, }); // Mock getBlock - containerProtocolMock.when(() -> ContainerProtocolCalls.getBlock(any(), any(), any(), any(), anyMap())) + containerProtocolMock.when(() -> ContainerProtocolCalls.getBlock(any(), any(), any(), any(), any())) .thenAnswer(inv -> { XceiverClientSpi xceiverClientSpi = inv.getArgument(0); BlockID blockID = inv.getArgument(2); diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java index 5cae34533219..0a5c89bf2032 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java @@ -58,7 +58,7 @@ public BlockVerificationResult verifyBlock(DatanodeDetails datanode, OmKeyLocati client, keyLocation.getBlockID(), keyLocation.getToken(), - pipeline.getReplicaIndexes() + pipeline ); boolean hasBlock = response != null && response.hasBlockData(); diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/chunk/ChunkKeyHandler.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/chunk/ChunkKeyHandler.java index 146354119220..b1ae8d9c8aa1 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/chunk/ChunkKeyHandler.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/chunk/ChunkKeyHandler.java @@ -131,7 +131,7 @@ protected void execute(OzoneClient client, OzoneAddress address) keyLocation.getBlockID(), keyLocation.getToken(), datanodeDetails, - pipeline.getReplicaIndexes()); + pipeline); if (blockResponse == null || !blockResponse.hasBlockData()) { System.err.printf("GetBlock call failed on %s datanode and %s block.%n", diff --git a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ECFileChecksumHelper.java b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ECFileChecksumHelper.java index 1f6e18426aac..faaa1e7725cb 100644 --- a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ECFileChecksumHelper.java +++ b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ECFileChecksumHelper.java @@ -123,7 +123,7 @@ protected List getChunkInfos(OmKeyLocationInfo xceiverClientSpi = getXceiverClientFactory().acquireClientForReadData(pipeline); ContainerProtos.GetBlockResponseProto response = ContainerProtocolCalls - .getBlock(xceiverClientSpi, blockID, token, pipeline.getReplicaIndexes()); + .getBlock(xceiverClientSpi, blockID, token, pipeline); chunks = response.getBlockData().getChunksList(); } finally { diff --git a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ReplicatedFileChecksumHelper.java b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ReplicatedFileChecksumHelper.java index 14d6d0b05a32..42a07690a3ff 100644 --- a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ReplicatedFileChecksumHelper.java +++ b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/checksum/ReplicatedFileChecksumHelper.java @@ -75,7 +75,7 @@ protected List getChunkInfos( } xceiverClientSpi = getXceiverClientFactory().acquireClientForReadData(pipeline); ContainerProtos.GetBlockResponseProto response = ContainerProtocolCalls - .getBlock(xceiverClientSpi, blockID, token, pipeline.getReplicaIndexes()); + .getBlock(xceiverClientSpi, blockID, token, pipeline); chunks = response.getBlockData().getChunksList(); } finally { diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestXceiverClientGrpc.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestXceiverClientGrpc.java index ca346a6bc986..8cfcac043343 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestXceiverClientGrpc.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestXceiverClientGrpc.java @@ -269,7 +269,7 @@ private void invokeXceiverClientGetBlock(XceiverClientSpi client) .setContainerID(1) .setLocalID(1) .setBlockCommitSequenceId(1) - .build()), null, client.getPipeline().getReplicaIndexes()); + .build()), null, client.getPipeline()); } private void invokeXceiverClientReadChunk(XceiverClientSpi client) diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java index a70333d4b217..1e82a544f17f 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java @@ -516,7 +516,7 @@ public void testCreateRecoveryContainer() throws Exception { ContainerProtos.ReadChunkResponseProto readChunkResponseProto = ContainerProtocolCalls.readChunk(dnClient, writeChunkRequest.getWriteChunk().getChunkData(), - blockID.getDatanodeBlockIDProtobufBuilder().setReplicaIndex(replicaIndex).build(), null, + blockID.getDatanodeBlockIDProtobufBuilder(replicaIndex).build(), null, blockToken); ByteBuffer[] readOnlyByteBuffersArray = BufferUtils .getReadOnlyByteBuffersArray( @@ -834,7 +834,7 @@ private void triggerRetryByCloseContainer(OzoneOutputStream out) { BlockID entryBlockID = blockOutputStreamEntry.getBlockID(); long entryContainerID = entryBlockID.getContainerID(); Pipeline entryPipeline = blockOutputStreamEntry.getPipeline(); - Map replicaIndexes = entryPipeline.getReplicaIndexes(); + Map replicaIndexes = entryPipeline.getReplicaIndexesForTesting(); try { for (Map.Entry entry : replicaIndexes.entrySet()) { DatanodeDetails key = entry.getKey(); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestECKeyOutputStream.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestECKeyOutputStream.java index 5a7f2ad68f22..3dd63bbd05e5 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestECKeyOutputStream.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestECKeyOutputStream.java @@ -225,9 +225,11 @@ public void testECKeyCreatetWithDatanodeIdChange() OmKeyLocationInfo omKeyLocationInfo = locationInfoList.get(0); long containerId = omKeyLocationInfo.getContainerID(); Pipeline pipeline = omKeyLocationInfo.getPipeline(); - DatanodeDetails dnWithReplicaIndex1 = - pipeline.getReplicaIndexes().entrySet().stream().filter(e -> e.getValue() == 1).map(Map.Entry::getKey) - .findFirst().get(); + DatanodeDetails dnWithReplicaIndex1 = pipeline.getReplicaIndexesForTesting().entrySet().stream() + .filter(e -> e.getValue() == 1) + .map(Map.Entry::getKey) + .findFirst() + .get(); Mockito.when(handlers.get(dnWithReplicaIndex1.getUuidString()).getDatanodeId()) .thenAnswer(i -> { if (!failed.get()) { diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java index 868a248e909e..a61595faf069 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java @@ -152,6 +152,20 @@ public StorageContainerLocationProtocol getContainerClient() { return this.containerClient; } + /** Don't keep empty pipelines or insufficient EC pipelines in the cache. */ + static boolean shouldInvalidate(Pipeline pipeline) { + if (pipeline.isEmpty()) { + return true; + } + + final ReplicationConfig repConfig = pipeline.getReplicationConfig(); + if (!(repConfig instanceof ECReplicationConfig)) { + return false; + } + final int d = ((ECReplicationConfig) repConfig).getData(); + return !pipeline.containsAllReplicaIndexes(1, d + 1); // 1 <= i <= d are the data strips + } + public Map getContainerLocations(Iterable containerIds, boolean forceRefresh) throws IOException { @@ -160,29 +174,13 @@ public Map getContainerLocations(Iterable containerIds, } try { Map result = containerLocationCache.getAll(containerIds); - // Don't keep empty pipelines or insufficient EC pipelines in the cache. - List uncachePipelines = result.entrySet().stream() - .filter(e -> { - Pipeline pipeline = e.getValue(); - // filter empty pipelines - if (pipeline.isEmpty()) { - return true; - } - // filter insufficient EC pipelines which missing any data index - ReplicationConfig repConfig = pipeline.getReplicationConfig(); - if (repConfig instanceof ECReplicationConfig) { - int d = ((ECReplicationConfig) repConfig).getData(); - for (int i = 1; i <= d; i++) { - if (!pipeline.getReplicaIndexes().containsValue(i)) { - return true; - } - } - } - return false; - }) - .map(Map.Entry::getKey) - .collect(Collectors.toList()); - containerLocationCache.invalidateAll(uncachePipelines); + + for (Map.Entry entry : result.entrySet()) { + if (shouldInvalidate(entry.getValue())) { + containerLocationCache.invalidate(entry.getKey()); + } + } + return result; } catch (ExecutionException e) { return handleCacheExecutionException(e);