This is an automated email from the ASF dual-hosted git repository.

szetszwo pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/master by this push:
     new 644bb55e6ba HDDS-16365. Pipeline.getReplicaIndexes() unnecessarily 
create a new map. (#11194)
644bb55e6ba is described below

commit 644bb55e6ba96f8a24731224a304d8ab443413eb
Author: Tsz-Wo Nicholas Sze <[email protected]>
AuthorDate: Sun Sep 13 09:52:18 2026 -0700

    HDDS-16365. Pipeline.getReplicaIndexes() unnecessarily create a new map. 
(#11194)
---
 .../hadoop/hdds/scm/storage/BlockInputStream.java  | 29 +++-----
 .../hadoop/hdds/scm/storage/BlockOutputStream.java | 13 +---
 .../hadoop/hdds/scm/storage/ChunkInputStream.java  |  6 +-
 .../org/apache/hadoop/hdds/client/BlockID.java     | 22 +++---
 .../apache/hadoop/hdds/scm/pipeline/Pipeline.java  | 23 +++++-
 .../hdds/scm/storage/ContainerProtocolCalls.java   | 86 ++++------------------
 ...stContainerReconciliationWithMockDatanodes.java |  3 +-
 .../debug/replicas/BlockExistenceVerifier.java     |  2 +-
 .../debug/replicas/chunk/ChunkKeyHandler.java      |  2 +-
 .../client/checksum/ECFileChecksumHelper.java      |  2 +-
 .../checksum/ReplicatedFileChecksumHelper.java     |  2 +-
 .../hadoop/hdds/scm/TestXceiverClientGrpc.java     |  2 +-
 .../hdds/scm/storage/TestContainerCommandsEC.java  |  2 +-
 .../ozone/client/rpc/TestECKeyOutputStream.java    |  8 +-
 .../java/org/apache/hadoop/ozone/om/ScmClient.java | 44 ++++++-----
 15 files changed, 92 insertions(+), 154 deletions(-)

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 92b03bd2acc..6d5ac9c8cf2 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 b960e753744..438425a3ceb 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 34ef7a71bd3..ee70a67dd85 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 7141a65306d..68da2f5192c 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) {
+      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 f6065a578e2..bff5c1043fe 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 if this pipeline contains all replica indexes within the 
given range.
    */
+  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<DatanodeDetails, Integer> getReplicaIndexes() {
-    return 
this.getNodes().stream().collect(Collectors.toMap(Function.identity(), 
this::getReplicaIndex));
+    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 d2898c395d9..b8d832cf45a 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> T tryEachDatanode(Pipeline pipeline,
    */
   public static GetBlockResponseProto getBlock(XceiverClientSpi xceiverClient,
       List<Validator> validators, BlockID blockID, Token<? extends 
TokenIdentifier> token,
-      Map<DatanodeDetails, Integer> 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<? extends TokenIdentifier> token, Map<DatanodeDetails, Integer> 
replicaIndexes) throws IOException {
-    return getBlock(xceiverClient, getValidatorList(), datanodeBlockID, token, 
replicaIndexes);
+      Token<? extends TokenIdentifier> 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<? extends TokenIdentifier> token,
       DatanodeDetails datanode,
-      Map<DatanodeDetails, Integer> 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<Validator> validators,
       ContainerCommandRequestProto.Builder builder, BlockID blockID,
-      DatanodeDetails datanode, Map<DatanodeDetails, Integer> 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<Validator> toValidatorList(Validator 
validator) {
     return Collections.unmodifiableList(validators);
   }
 
-  public static HashMap<DatanodeDetails, GetBlockResponseProto>
-      getBlockFromAllNodes(
-      XceiverClientSpi xceiverClient,
-      DatanodeBlockID datanodeBlockID,
-      Token<OzoneBlockTokenIdentifier> token)
-      throws IOException, InterruptedException {
-    GetBlockRequestProto.Builder readBlockRequest = GetBlockRequestProto
-            .newBuilder()
-            .setBlockID(datanodeBlockID);
-    HashMap<DatanodeDetails, GetBlockResponseProto> 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<DatanodeDetails, ContainerCommandResponseProto> responses =
-            xceiverClient.sendCommandOnAllNodes(request);
-    for (Map.Entry<DatanodeDetails, ContainerCommandResponseProto> entry:
-           responses.entrySet()) {
-      datanodeToResponseMap.put(entry.getKey(), 
entry.getValue().getGetBlock());
-    }
-    return datanodeToResponseMap;
-  }
-
   public static HashMap<DatanodeDetails, ReadContainerResponseProto>
       readContainerFromAllNodes(XceiverClientSpi client, long containerID,
       String encodedToken) throws IOException, InterruptedException {
@@ -1000,12 +954,12 @@ public static ContainerCommandRequestProto 
buildReadBlockCommandProto(
       Token<? extends TokenIdentifier> 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<DatanodeDetails, Integer> 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 c051b3478b4..a97a3418cc8 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 5cae3453321..0a5c89bf203 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 14635411922..b1ae8d9c8aa 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 1f6e18426aa..faaa1e7725c 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<ContainerProtos.ChunkInfo> 
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 14d6d0b05a3..42a07690a3f 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<ContainerProtos.ChunkInfo> 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 ca346a6bc98..8cfcac04334 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 a70333d4b21..3034adec9a3 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(
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 5a7f2ad68f2..097a1cf84e9 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.getReplicaIndexes().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 868a248e909..a61595faf06 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<Long, Pipeline> getContainerLocations(Iterable<Long> containerIds,
                                                   boolean forceRefresh)
       throws IOException {
@@ -160,29 +174,13 @@ public Map<Long, Pipeline> 
getContainerLocations(Iterable<Long> containerIds,
     }
     try {
       Map<Long, Pipeline> result = containerLocationCache.getAll(containerIds);
-      // Don't keep empty pipelines or insufficient EC pipelines in the cache.
-      List<Long> 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<Long, Pipeline> entry : result.entrySet()) {
+        if (shouldInvalidate(entry.getValue())) {
+          containerLocationCache.invalidate(entry.getKey());
+        }
+      }
+
       return result;
     } catch (ExecutionException e) {
       return handleCacheExecutionException(e);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to