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 dc6494ac7ae HDDS-15059. Shift streaming write sortDatanodes logic to 
OM (#10633)
dc6494ac7ae is described below

commit dc6494ac7ae8d98071fc600723c161b884cf4954
Author: Chi-Hsuan Huang <[email protected]>
AuthorDate: Tue Jul 28 01:46:07 2026 +0800

    HDDS-15059. Shift streaming write sortDatanodes logic to OM (#10633)
    
    Generated-by: Claude Code (Claude Opus 4.8)
---
 .../hadoop/hdds/scm/client/ScmTopologyClient.java  |   4 +-
 .../java/org/apache/hadoop/ozone/om/OmConfig.java  |  16 ++
 .../apache/hadoop/ozone/TestOMSortDatanodes.java   |  52 ++++
 .../org/apache/hadoop/ozone/om/KeyManager.java     |  25 ++
 .../org/apache/hadoop/ozone/om/KeyManagerImpl.java |  94 +++++--
 .../hadoop/ozone/om/OMPerformanceMetrics.java      |   7 +
 .../org/apache/hadoop/ozone/om/OzoneManager.java   |   6 +
 .../hadoop/ozone/om/request/key/OMKeyRequest.java  |  56 ++++-
 .../om/request/key/TestOMAllocateBlockRequest.java | 279 +++++++++++++++++++++
 9 files changed, 514 insertions(+), 25 deletions(-)

diff --git 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmTopologyClient.java
 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmTopologyClient.java
index d595bd6e095..7abd2c9c016 100644
--- 
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmTopologyClient.java
+++ 
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmTopologyClient.java
@@ -17,7 +17,6 @@
 
 package org.apache.hadoop.hdds.scm.client;
 
-import static java.util.Objects.requireNonNull;
 import static org.apache.hadoop.hdds.scm.net.NetConstants.ROOT;
 import static 
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_OM_NETWORK_TOPOLOGY_REFRESH_DURATION;
 import static 
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_OM_NETWORK_TOPOLOGY_REFRESH_DURATION_DEFAULT;
@@ -60,8 +59,7 @@ public ScmTopologyClient(
   }
 
   public NetworkTopology getClusterMap() {
-    return requireNonNull(cache.get(),
-        "ScmBlockLocationClient must have been initialized already.");
+    return cache.get();
   }
 
   public void start(ConfigurationSource conf) throws IOException {
diff --git 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OmConfig.java 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OmConfig.java
index b8a60f9bcdb..c981afd414d 100644
--- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OmConfig.java
+++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OmConfig.java
@@ -184,6 +184,18 @@ public class OmConfig extends ReconfigurableConfig {
   )
   private long followerReadLocalLeaseTimeMs;
 
+  @Config(key = "ozone.om.block.write.sort.datanodes.enabled",
+      defaultValue = "false",
+      type = ConfigType.BOOLEAN,
+      tags = {ConfigTag.OM, ConfigTag.PERFORMANCE},
+      description = "If true, OM sorts the streaming-write pipeline (nearest " 
+
+          "datanode first) locally using its cached cluster topology, instead 
" +
+          "of asking SCM to sort on every allocateBlock. Defaults to false so 
" +
+          "SCM performs the sort. Enable this to offload the sort from SCM " +
+          "when multiple OM services share a single SCM service."
+  )
+  private boolean sortDatanodesForWriteEnabled;
+
   public long getRatisBasedFinalizationTimeout() {
     return ratisBasedFinalizationTimeout;
   }
@@ -224,6 +236,10 @@ public void setAllowLeaderSkipLinearizableRead(boolean 
newValue) {
     allowLeaderSkipLinearizableRead = newValue;
   }
 
+  public boolean isSortDatanodesForWriteEnabled() {
+    return sortDatanodesForWriteEnabled;
+  }
+
   public boolean isFollowerReadLocalLeaseEnabled() {
     return followerReadLocalLeaseEnabled;
   }
diff --git 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOMSortDatanodes.java
 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOMSortDatanodes.java
index cfce524537f..a206276ce8f 100644
--- 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOMSortDatanodes.java
+++ 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOMSortDatanodes.java
@@ -22,6 +22,7 @@
 import static org.apache.hadoop.hdds.scm.net.NetConstants.ROOT_LEVEL;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.mockito.Mockito.mock;
 
 import com.google.common.collect.ImmutableMap;
@@ -176,6 +177,57 @@ private static void assertRackOrder(String rack, List<? 
extends DatanodeDetails>
     }
   }
 
+  @Test
+  public void sortDatanodesForWriteSortsRpcDeserializedPipeline() {
+    // Pipeline nodes arrive from SCM over RPC as deserialized DatanodeDetails
+    // with no topology linkage; the write sort must resolve them to OM's 
cluster
+    // map, otherwise every node is equidistant (MAX) and the order is random.
+    List<DatanodeDetails> rpcNodes = new ArrayList<>();
+    for (DatanodeDetails dn : nodeManager.getAllNodes()) {
+      rpcNodes.add(DatanodeDetails.getFromProtoBuf(dn.getProtoBufMessage()));
+    }
+    for (DatanodeDetails dn : nodeManager.getAllNodes()) {
+      // The client address is normally an IP, but the sort must resolve a 
client
+      // by either IP or hostname, so cover both.
+      List<? extends DatanodeDetails> byIp =
+          keyManager.sortDatanodesForWrite(rpcNodes, dn.getIpAddress(), 
om.getClusterMap());
+      assertEquals(dn, byIp.get(0),
+          "Source node should be sorted first for writes (IP client)");
+      assertRackOrder(dn.getNetworkLocation(), byIp);
+
+      List<? extends DatanodeDetails> byHostname =
+          keyManager.sortDatanodesForWrite(rpcNodes, dn.getHostName(), 
om.getClusterMap());
+      assertEquals(dn, byHostname.get(0),
+          "Source node should be sorted first for writes (hostname client)");
+      assertRackOrder(dn.getNetworkLocation(), byHostname);
+    }
+  }
+
+  @Test
+  public void sortDatanodesForWriteKeepsOrderForStaleTopology() {
+    List<DatanodeDetails> nodes = new ArrayList<>();
+    nodes.add(randomDatanodeDetails());
+    nodes.addAll(nodeManager.getAllNodes());
+
+    List<? extends DatanodeDetails> sorted =
+        keyManager.sortDatanodesForWrite(nodes, "edge0", om.getClusterMap());
+
+    assertSame(nodes, sorted,
+        "Pipeline order should be preserved when a node is missing from the OM 
topology");
+  }
+
+  @Test
+  public void sortDatanodesForWriteKeepsOrderWhenClientUnresolved() {
+    List<? extends DatanodeDetails> nodes = nodeManager.getAllNodes();
+    List<DatanodeDetails> original = new ArrayList<>(nodes);
+    // A client that resolves to no known rack must NOT trigger a shuffle.
+    String unresolved = nodes.get(0).getIpAddress() + "X";
+    List<? extends DatanodeDetails> result =
+        keyManager.sortDatanodesForWrite(nodes, unresolved, 
om.getClusterMap());
+    assertEquals(original, result,
+        "Write pipeline order must be preserved when client is unresolved");
+  }
+
   private String nodeAddress(DatanodeDetails dn) {
     boolean useHostname = config.getBoolean(
         HddsConfigKeys.HDDS_DATANODE_USE_DN_HOSTNAME,
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManager.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManager.java
index 214d738b630..4077ee088ef 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManager.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManager.java
@@ -23,6 +23,8 @@
 import java.util.List;
 import java.util.Map;
 import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.scm.net.NetworkTopology;
 import org.apache.hadoop.hdds.utils.BackgroundService;
 import org.apache.hadoop.hdds.utils.db.Table;
 import org.apache.hadoop.hdds.utils.db.TableIterator;
@@ -370,4 +372,27 @@ DeleteKeysResult getPendingDeletionSubFiles(long volumeId, 
long bucketId, OmKeyI
    * @return Background service.
    */
   KeyLifecycleService getKeyLifecycleService();
+
+  /**
+   * Sort the datanodes of a write pipeline by network-topology distance to the
+   * client, using OM's locally cached cluster map. Unlike the read-path sort,
+   * the original order is preserved when the client cannot be resolved, 
because
+   * the first node is used as the streaming-write primary.
+   *
+   * @param nodes the pipeline nodes to sort
+   * @param clientMachine client address (IP or hostname)
+   * @param clusterMap OM's cached cluster map used to resolve topology 
distance
+   * @return nodes sorted nearest-first, or the original {@code nodes} list
+   *     instance unchanged when sorting is skipped (client unresolved or stale
+   *     topology); callers may use reference equality to detect a skipped sort
+   */
+  List<? extends DatanodeDetails> sortDatanodesForWrite(
+      List<? extends DatanodeDetails> nodes, String clientMachine, 
NetworkTopology clusterMap);
+
+  /**
+   * @return true if OM should sort the streaming-write pipeline locally
+   *     ({@code ozone.om.block.write.sort.datanodes.enabled}); false to leave
+   *     the sort to SCM.
+   */
+  boolean isSortDatanodesForWriteEnabled();
 }
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java
index 6ecfe653995..f24a5c470fb 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java
@@ -123,7 +123,6 @@
 import org.apache.hadoop.crypto.key.KeyProviderCryptoExtension;
 import 
org.apache.hadoop.crypto.key.KeyProviderCryptoExtension.EncryptedKeyVersion;
 import org.apache.hadoop.fs.FileEncryptionInfo;
-import org.apache.hadoop.hdds.HddsConfigKeys;
 import org.apache.hadoop.hdds.client.BlockID;
 import org.apache.hadoop.hdds.client.ReplicationConfig;
 import org.apache.hadoop.hdds.conf.OzoneConfiguration;
@@ -132,6 +131,7 @@
 import org.apache.hadoop.hdds.scm.ScmConfigKeys;
 import 
org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
 import org.apache.hadoop.hdds.scm.net.InnerNode;
+import org.apache.hadoop.hdds.scm.net.NetworkTopology;
 import org.apache.hadoop.hdds.scm.net.Node;
 import org.apache.hadoop.hdds.scm.net.NodeImpl;
 import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
@@ -2262,7 +2262,12 @@ private void sortDatanodes(String clientMachine, 
List<OmKeyInfo> keyInfos) {
           List<? extends DatanodeDetails> sortedNodes = 
sortedPipelines.get(uuidSet);
           if (sortedNodes == null) {
             sortedNodes = sortDatanodes(nodes, clientMachine);
-            if (sortedNodes != null) {
+            // Cache only a freshly sorted order, not an input list returned
+            // unchanged when no sort happens: that order is per-pipeline and 
must
+            // not be reused for another pipeline with the same node set. The 
read
+            // sort always returns a new list, so this never skips caching 
here; it
+            // keeps the pattern identical to the write path.
+            if (sortedNodes != null && sortedNodes != nodes) {
               sortedPipelines.put(uuidSet, sortedNodes);
             }
           } else if (LOG.isDebugEnabled()) {
@@ -2280,32 +2285,87 @@ private void sortDatanodes(String clientMachine, 
List<OmKeyInfo> keyInfos) {
   @VisibleForTesting
   public List<? extends DatanodeDetails> sortDatanodes(List<? extends 
DatanodeDetails> nodes,
                                              String clientMachine) {
-    final Node client = getClientNode(clientMachine, nodes);
-    return ozoneManager.getClusterMap()
-        .sortByDistanceCost(client, nodes, nodes.size());
+    final NetworkTopology clusterMap = ozoneManager.getClusterMap();
+    final Node client = getClientNode(clientMachine, nodes, clusterMap);
+    return clusterMap.sortByDistanceCost(client, nodes, nodes.size());
+  }
+
+  @Override
+  public List<? extends DatanodeDetails> sortDatanodesForWrite(
+      List<? extends DatanodeDetails> nodes, String clientMachine, 
NetworkTopology clusterMap) {
+    Preconditions.checkArgument(!StringUtils.isEmpty(clientMachine),
+        "clientMachine is empty");
+    Objects.requireNonNull(clusterMap, "clusterMap is null");
+    return captureLatencyNs(
+        metrics.getAllocateBlockSortDatanodesLatencyNs(), () -> {
+          final Node client = getClientNode(clientMachine, nodes, clusterMap);
+          if (client == null) {
+            // Preserve pipeline order for writes: the first node is the write
+            // primary, so do not shuffle when the client cannot be resolved.
+            return nodes;
+          }
+          return sortByClusterMapDistance(clusterMap, client, nodes);
+        });
+  }
+
+  @Override
+  public boolean isSortDatanodesForWriteEnabled() {
+    return ozoneManager.getConfig().isSortDatanodesForWriteEnabled();
+  }
+
+  /**
+   * Sort a pipeline's nodes by topology distance to the client. The nodes come
+   * from SCM over RPC, so they are deserialized {@link DatanodeDetails} with 
no
+   * parent/level: the topology treats them as unknown (distance
+   * {@link Integer#MAX_VALUE}) and the order comes out random. Look each node
+   * up in OM's cluster map to get the topology-linked instance, sort those,
+   * then map the order back to the original nodes.
+   */
+  private List<? extends DatanodeDetails> sortByClusterMapDistance(
+      NetworkTopology clusterMap, Node client,
+      List<? extends DatanodeDetails> nodes) {
+    final List<Node> topologyNodes = new ArrayList<>(nodes.size());
+    final Map<String, DatanodeDetails> nodeByPath = new HashMap<>();
+    for (DatanodeDetails node : nodes) {
+      final Node resolved = clusterMap.getNode(node.getNetworkFullPath());
+      if (resolved == null) {
+        return nodes;
+      }
+      topologyNodes.add(resolved);
+      nodeByPath.put(resolved.getNetworkFullPath(), node);
+    }
+    final List<Node> sorted =
+        clusterMap.sortByDistanceCost(client, topologyNodes, 
topologyNodes.size());
+    final List<DatanodeDetails> result = new ArrayList<>(sorted.size());
+    for (Node node : sorted) {
+      result.add(nodeByPath.get(node.getNetworkFullPath()));
+    }
+    return result;
   }
 
   private Node getClientNode(String clientMachine,
-                             List<? extends DatanodeDetails> nodes) {
-    List<DatanodeDetails> matchingNodes = new ArrayList<>();
-    boolean useHostname = ozoneManager.getConfiguration().getBoolean(
-        HddsConfigKeys.HDDS_DATANODE_USE_DN_HOSTNAME,
-        HddsConfigKeys.HDDS_DATANODE_USE_DN_HOSTNAME_DEFAULT);
+      List<? extends DatanodeDetails> nodes, NetworkTopology clusterMap) {
     for (DatanodeDetails node : nodes) {
-      if ((useHostname ? node.getHostName() : node.getIpAddress()).equals(
-          clientMachine)) {
-        matchingNodes.add(node);
+      // Match by either IP or hostname, like SCM's getNodesByAddress. 
clientMachine
+      // may be a hostname on the read path; the streaming-write remoteAddress 
is
+      // typically an IP. Matching both covers use.datanode.hostname either 
way.
+      if (clientMachine.equals(node.getIpAddress())
+          || clientMachine.equals(node.getHostName())) {
+        // The pipeline nodes are RPC-deserialized and not linked into OM's
+        // cluster map; prefer the map's instance so distance can be computed.
+        final Node resolved = clusterMap.getNode(node.getNetworkFullPath());
+        return resolved != null ? resolved : node;
       }
     }
-    return !matchingNodes.isEmpty() ? matchingNodes.get(0) :
-        getOtherNode(clientMachine);
+    return getOtherNode(clientMachine, clusterMap);
   }
 
-  private Node getOtherNode(String clientMachine) {
+  private Node getOtherNode(String clientMachine,
+                            NetworkTopology clusterMap) {
     try {
       String clientLocation = resolveNodeLocation(clientMachine);
       if (clientLocation != null) {
-        Node rack = ozoneManager.getClusterMap().getNode(clientLocation);
+        Node rack = clusterMap.getNode(clientLocation);
         if (rack instanceof InnerNode) {
           return new NodeImpl(clientMachine, clientLocation,
               (InnerNode) rack, rack.getLevel() + 1,
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OMPerformanceMetrics.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OMPerformanceMetrics.java
index e142e4cfa71..3f965ba025e 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OMPerformanceMetrics.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OMPerformanceMetrics.java
@@ -67,6 +67,9 @@ public class OMPerformanceMetrics {
   @Metric(about = "Sort datanodes latency in getKeyInfo")
   private MutableRate getKeyInfoSortDatanodesLatencyNs;
 
+  @Metric(about = "Sort datanodes latency in allocateBlock (streaming write)")
+  private MutableRate allocateBlockSortDatanodesLatencyNs;
+
   @Metric(about = "resolveBucketLink latency in getKeyInfo")
   private MutableRate getKeyInfoResolveBucketLatencyNs;
 
@@ -246,6 +249,10 @@ MutableRate getGetKeyInfoSortDatanodesLatencyNs() {
     return getKeyInfoSortDatanodesLatencyNs;
   }
 
+  MutableRate getAllocateBlockSortDatanodesLatencyNs() {
+    return allocateBlockSortDatanodesLatencyNs;
+  }
+
   public void setForceContainerCacheRefresh(boolean value) {
     forceContainerCacheRefresh.add(value ? 1L : 0L);
   }
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
index 709c59da16b..2db20c59f31 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
@@ -18,6 +18,7 @@
 package org.apache.hadoop.ozone.om;
 
 import static java.nio.charset.StandardCharsets.UTF_8;
+import static java.util.Objects.requireNonNull;
 import static 
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_TRASH_INTERVAL_DEFAULT;
 import static 
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_TRASH_INTERVAL_KEY;
 import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_BLOCK_TOKEN_ENABLED;
@@ -1350,6 +1351,11 @@ public void setScmTopologyClient(
   }
 
   public NetworkTopology getClusterMap() {
+    return requireNonNull(getClusterMapAllowNull(),
+        "OM topology cache has not been initialized yet.");
+  }
+
+  public NetworkTopology getClusterMapAllowNull() {
     return scmTopologyClient.getClusterMap();
   }
 
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java
index d12b8fa0525..4e2253315b8 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyRequest.java
@@ -57,11 +57,14 @@
 import org.apache.hadoop.hdds.client.ContainerBlockID;
 import org.apache.hadoop.hdds.client.ECReplicationConfig;
 import org.apache.hadoop.hdds.client.ReplicationConfig;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
 import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
 import 
org.apache.hadoop.hdds.protocol.proto.HddsProtos.BlockTokenSecretProto.AccessModeProto;
 import org.apache.hadoop.hdds.scm.container.common.helpers.AllocatedBlock;
 import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList;
 import org.apache.hadoop.hdds.scm.exceptions.SCMException;
+import org.apache.hadoop.hdds.scm.net.NetworkTopology;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
 import org.apache.hadoop.hdds.security.token.OzoneBlockTokenIdentifier;
 import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
 import org.apache.hadoop.hdds.utils.db.cache.CacheValue;
@@ -69,6 +72,7 @@
 import org.apache.hadoop.ozone.OmUtils;
 import org.apache.hadoop.ozone.OzoneAcl;
 import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.om.KeyManager;
 import org.apache.hadoop.ozone.om.OMMetadataManager;
 import org.apache.hadoop.ozone.om.OmConfig;
 import org.apache.hadoop.ozone.om.OzoneManager;
@@ -189,15 +193,38 @@ protected List<OmKeyLocationInfo> allocateBlock(
       UserInfo userInfo, OzoneManager ozoneManager)
       throws IOException {
     final long scmBlockSize = ozoneManager.getScmBlockSize();
+    final KeyManager keyManager = ozoneManager.getKeyManager();
 
     int dataGroupSize = replicationConfig instanceof ECReplicationConfig
         ? ((ECReplicationConfig) replicationConfig).getData() : 1;
     final int numBlocks = (int) 
Math.min(ozoneManager.getPreallocateBlocksMax(),
         (requestedSize - 1) / (scmBlockSize * dataGroupSize) + 1);
 
-    String clientMachine = "";
-    if (shouldSortDatanodes) {
-      clientMachine = userInfo.getRemoteAddress();
+    final String scmClientMachine;
+    final String omClientMachine;
+    // Sorted order cached by datanode set so blocks whose pipelines share the
+    // same datanodes are sorted once (mirrors the read path's caching). Keyed 
by
+    // the UUID set so it is order-insensitive and dedups across pipelines.
+    final Map<Set<String>, List<? extends DatanodeDetails>> sortedByNodes;
+    final String remoteAddress = userInfo.getRemoteAddress();
+    final NetworkTopology clusterMap = shouldSortDatanodes
+        && keyManager.isSortDatanodesForWriteEnabled()
+        ? ozoneManager.getClusterMapAllowNull() : null;
+    if (!shouldSortDatanodes) {
+      scmClientMachine = "";
+      omClientMachine = "";
+      sortedByNodes = null;
+    } else if (clusterMap != null && !remoteAddress.isEmpty()) {
+      // Sort in OM: SCM skips sorting (empty machine), OM sorts by 
remoteAddress.
+      scmClientMachine = "";
+      omClientMachine = remoteAddress;
+      sortedByNodes = new HashMap<>();
+    } else {
+      // Sort in SCM (or keep order when remoteAddress is empty, since SCM 
skips
+      // sorting for an empty client machine).
+      scmClientMachine = remoteAddress;
+      omClientMachine = "";
+      sortedByNodes = null;
     }
 
     List<OmKeyLocationInfo> locationInfos = new ArrayList<>(numBlocks);
@@ -205,7 +232,7 @@ protected List<OmKeyLocationInfo> allocateBlock(
     final List<AllocatedBlock> allocatedBlocks;
     try {
       allocatedBlocks = 
ozoneManager.getScmClient().getBlockClient().allocateBlock(
-          scmBlockSize, numBlocks, replicationConfig, 
ozoneManager.getOMServiceId(), excludeList, clientMachine);
+          scmBlockSize, numBlocks, replicationConfig, 
ozoneManager.getOMServiceId(), excludeList, scmClientMachine);
     } catch (SCMException ex) {
       ozoneManager.getMetrics().incNumBlockAllocateCallFails();
       if (ex.getResult() == SCMException.ResultCodes.SAFE_MODE_EXCEPTION) {
@@ -216,11 +243,30 @@ protected List<OmKeyLocationInfo> allocateBlock(
     }
     for (AllocatedBlock allocatedBlock : allocatedBlocks) {
       BlockID blockID = new BlockID(allocatedBlock.getBlockID());
+      Pipeline pipeline = allocatedBlock.getPipeline();
+      if (sortedByNodes != null) {
+        final List<DatanodeDetails> nodes = pipeline.getNodes();
+        final Set<String> uuidSet = nodes.stream()
+            .map(DatanodeDetails::getUuidString).collect(Collectors.toSet());
+        List<? extends DatanodeDetails> sorted = sortedByNodes.get(uuidSet);
+        if (sorted == null) {
+          sorted = keyManager.sortDatanodesForWrite(nodes, omClientMachine, 
clusterMap);
+          // Cache only a freshly sorted order, not an input list returned
+          // unchanged when the client is unresolved: that order is 
per-pipeline
+          // and must not be reused for another pipeline with the same node 
set.
+          if (sorted != nodes) {
+            sortedByNodes.put(uuidSet, sorted);
+          }
+        }
+        if (!Objects.equals(sorted, pipeline.getNodesInOrder())) {
+          pipeline = pipeline.copyWithNodesInOrder(sorted);
+        }
+      }
       OmKeyLocationInfo.Builder builder = new OmKeyLocationInfo.Builder()
           .setBlockID(blockID)
           .setLength(scmBlockSize)
           .setOffset(0)
-          .setPipeline(allocatedBlock.getPipeline());
+          .setPipeline(pipeline);
       if (ozoneManager.isGrpcBlockTokenEnabled()) {
         final Token<OzoneBlockTokenIdentifier> token = 
ozoneManager.getBlockTokenSecretManager().generateToken(
             remoteUser, blockID, READ_WRITE, scmBlockSize);
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMAllocateBlockRequest.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMAllocateBlockRequest.java
index 8f720deaf1b..58949fe0981 100644
--- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMAllocateBlockRequest.java
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMAllocateBlockRequest.java
@@ -20,13 +20,41 @@
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
 
 import jakarta.annotation.Nonnull;
+import java.net.InetAddress;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
 import java.util.List;
 import java.util.UUID;
+import org.apache.hadoop.hdds.client.ContainerBlockID;
 import org.apache.hadoop.hdds.client.RatisReplicationConfig;
+import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
+import org.apache.hadoop.hdds.scm.container.common.helpers.AllocatedBlock;
+import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList;
+import org.apache.hadoop.hdds.scm.net.NetworkTopology;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
+import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
+import org.apache.hadoop.ipc_.Server;
 import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.om.KeyManager;
 import org.apache.hadoop.ozone.om.helpers.BucketLayout;
 import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
 import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo;
@@ -36,7 +64,11 @@
 import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.AllocateBlockRequest;
 import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyArgs;
 import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
+import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.UserInfo;
+import org.apache.hadoop.security.UserGroupInformation;
 import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
 
 /**
  * Tests OMAllocateBlockRequest class.
@@ -225,6 +257,253 @@ protected OMRequest doPreExecute(OMRequest 
originalOMRequest)
     return modifiedOmRequest;
   }
 
+  @Test
+  public void testAllocateBlockSendsClientMachineToScmWhenFlagOff() throws 
Exception {
+    // Flag off (default): OM must NOT sort; SCM receives the real client 
address
+    // so it performs the sort.
+    KeyManager mockKeyManager = mock(KeyManager.class);
+    when(mockKeyManager.isSortDatanodesForWriteEnabled()).thenReturn(false);
+    when(ozoneManager.getKeyManager()).thenReturn(mockKeyManager);
+
+    OMAllocateBlockRequest request =
+        getOmAllocateBlockRequest(createAllocateBlockRequestWithSort());
+    preExecuteWithClient(request, "1.2.3.4");
+
+    ArgumentCaptor<String> clientMachine = 
ArgumentCaptor.forClass(String.class);
+    verify(scmBlockLocationProtocol).allocateBlock(anyLong(), anyInt(), any(),
+        any(), any(), clientMachine.capture());
+    assertEquals("1.2.3.4", clientMachine.getValue());
+    verify(mockKeyManager, never()).sortDatanodesForWrite(any(), anyString(), 
any());
+  }
+
+  @Test
+  public void testAllocateBlockDoesNotSendClientMachineToScm() throws 
Exception {
+    // OM now sorts the write pipeline locally, so SCM must receive an empty
+    // clientMachine even when the client requests sorted datanodes.
+    KeyManager mockKeyManager = mock(KeyManager.class);
+    when(mockKeyManager.isSortDatanodesForWriteEnabled()).thenReturn(true);
+    when(mockKeyManager.sortDatanodesForWrite(any(), any(), any()))
+        .thenAnswer(inv -> inv.getArgument(0));
+    when(ozoneManager.getKeyManager()).thenReturn(mockKeyManager);
+    
when(ozoneManager.getClusterMapAllowNull()).thenReturn(mock(NetworkTopology.class));
+
+    OMAllocateBlockRequest request =
+        getOmAllocateBlockRequest(createAllocateBlockRequestWithSort());
+    preExecuteWithClient(request, "1.2.3.4");
+
+    ArgumentCaptor<String> clientMachine = 
ArgumentCaptor.forClass(String.class);
+    verify(scmBlockLocationProtocol).allocateBlock(anyLong(), anyInt(), any(),
+        any(), any(), clientMachine.capture());
+    assertEquals("", clientMachine.getValue());
+  }
+
+  @Test
+  public void testAllocateBlockFallsBackToScmWhenTopologyUnavailable() throws 
Exception {
+    KeyManager mockKeyManager = mock(KeyManager.class);
+    when(mockKeyManager.isSortDatanodesForWriteEnabled()).thenReturn(true);
+    when(ozoneManager.getKeyManager()).thenReturn(mockKeyManager);
+
+    OMAllocateBlockRequest request =
+        getOmAllocateBlockRequest(createAllocateBlockRequestWithSort());
+    preExecuteWithClient(request, "1.2.3.4");
+
+    ArgumentCaptor<String> clientMachine = 
ArgumentCaptor.forClass(String.class);
+    verify(scmBlockLocationProtocol).allocateBlock(anyLong(), anyInt(), any(),
+        any(), any(), clientMachine.capture());
+    assertEquals("1.2.3.4", clientMachine.getValue());
+    verify(mockKeyManager, never()).sortDatanodesForWrite(any(), anyString(), 
any());
+  }
+
+  @Test
+  public void testAllocateBlockSortsSharedPipelineOnce() throws Exception {
+    // Two blocks on the same 3-node pipeline must be sorted once, and the
+    // sorted order must land in every block's pipeline.
+    List<DatanodeDetails> nodes = Arrays.asList(
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails());
+    Pipeline pipeline = Pipeline.newBuilder()
+        .setState(Pipeline.PipelineState.OPEN)
+        .setId(PipelineID.randomId())
+        .setReplicationConfig(
+            StandaloneReplicationConfig.getInstance(ReplicationFactor.THREE))
+        .setNodes(nodes)
+        .build();
+    AllocatedBlock.Builder blockBuilder =
+        new AllocatedBlock.Builder().setPipeline(pipeline);
+    when(scmBlockLocationProtocol.allocateBlock(anyLong(), anyInt(), any(),
+        anyString(), any(ExcludeList.class), anyString())).thenAnswer(inv -> {
+          int num = inv.getArgument(1);
+          List<AllocatedBlock> blocks = new ArrayList<>(num);
+          for (int i = 0; i < num; i++) {
+            blockBuilder.setContainerBlockID(
+                new ContainerBlockID(CONTAINER_ID + i, LOCAL_ID + i));
+            blocks.add(blockBuilder.build());
+          }
+          return blocks;
+        });
+
+    List<DatanodeDetails> sortedOrder = new ArrayList<>(nodes);
+    Collections.reverse(sortedOrder);
+    KeyManager mockKeyManager = mock(KeyManager.class);
+    when(mockKeyManager.isSortDatanodesForWriteEnabled()).thenReturn(true);
+    when(mockKeyManager.sortDatanodesForWrite(any(), any(), any()))
+        .thenAnswer(inv -> sortedOrder);
+    when(ozoneManager.getKeyManager()).thenReturn(mockKeyManager);
+    
when(ozoneManager.getClusterMapAllowNull()).thenReturn(mock(NetworkTopology.class));
+
+    OMAllocateBlockRequest request =
+        getOmAllocateBlockRequest(createAllocateBlockRequest());
+    // requestedSize spans two scmBlockSize blocks on the same pipeline.
+    List<OmKeyLocationInfo> locations = 
request.allocateBlock(replicationConfig,
+        new ExcludeList(), 2 * scmBlockSize, true,
+        UserInfo.newBuilder().setRemoteAddress("1.2.3.4").build(), 
ozoneManager);
+
+    // Sorted once for the shared pipeline...
+    verify(mockKeyManager, times(1)).sortDatanodesForWrite(any(), 
eq("1.2.3.4"), any());
+    // ...and the sorted order is applied to every block's pipeline.
+    assertEquals(2, locations.size());
+    for (OmKeyLocationInfo location : locations) {
+      assertEquals(sortedOrder, location.getPipeline().getNodesInOrder());
+    }
+  }
+
+  @Test
+  public void testAllocateBlockKeepsPerPipelineOrderWhenSortSkipped() throws 
Exception {
+    // Two pipelines share the same datanode set but in a different order. When
+    // the sort is skipped (sortDatanodesForWrite returns the input unchanged),
+    // each pipeline must keep its own order: the unsorted result must not be
+    // cached under the node set and reused for the other pipeline.
+    DatanodeDetails a = MockDatanodeDetails.randomDatanodeDetails();
+    DatanodeDetails b = MockDatanodeDetails.randomDatanodeDetails();
+    DatanodeDetails c = MockDatanodeDetails.randomDatanodeDetails();
+    List<DatanodeDetails> nodes1 = Arrays.asList(a, b, c);
+    List<DatanodeDetails> nodes2 = Arrays.asList(c, b, a);
+    Pipeline pipeline1 = Pipeline.newBuilder()
+        .setState(Pipeline.PipelineState.OPEN)
+        .setId(PipelineID.randomId())
+        .setReplicationConfig(
+            StandaloneReplicationConfig.getInstance(ReplicationFactor.THREE))
+        .setNodes(nodes1)
+        .build();
+    Pipeline pipeline2 = Pipeline.newBuilder()
+        .setState(Pipeline.PipelineState.OPEN)
+        .setId(PipelineID.randomId())
+        .setReplicationConfig(
+            StandaloneReplicationConfig.getInstance(ReplicationFactor.THREE))
+        .setNodes(nodes2)
+        .build();
+    AllocatedBlock block1 = new AllocatedBlock.Builder().setPipeline(pipeline1)
+        .setContainerBlockID(new ContainerBlockID(CONTAINER_ID, 
LOCAL_ID)).build();
+    AllocatedBlock block2 = new AllocatedBlock.Builder().setPipeline(pipeline2)
+        .setContainerBlockID(new ContainerBlockID(CONTAINER_ID + 1, LOCAL_ID + 
1)).build();
+    when(scmBlockLocationProtocol.allocateBlock(anyLong(), anyInt(), any(),
+        anyString(), any(ExcludeList.class), anyString()))
+        .thenReturn(Arrays.asList(block1, block2));
+
+    KeyManager mockKeyManager = mock(KeyManager.class);
+    when(mockKeyManager.isSortDatanodesForWriteEnabled()).thenReturn(true);
+    // Skip the sort: return the input list instance unchanged.
+    when(mockKeyManager.sortDatanodesForWrite(any(), any(), any()))
+        .thenAnswer(inv -> inv.getArgument(0));
+    when(ozoneManager.getKeyManager()).thenReturn(mockKeyManager);
+    
when(ozoneManager.getClusterMapAllowNull()).thenReturn(mock(NetworkTopology.class));
+
+    OMAllocateBlockRequest request =
+        getOmAllocateBlockRequest(createAllocateBlockRequest());
+    List<OmKeyLocationInfo> locations = 
request.allocateBlock(replicationConfig,
+        new ExcludeList(), 2 * scmBlockSize, true,
+        UserInfo.newBuilder().setRemoteAddress("1.2.3.4").build(), 
ozoneManager);
+
+    assertEquals(2, locations.size());
+    // Each pipeline keeps its own order; the skipped-sort result is not 
shared.
+    assertEquals(nodes1, locations.get(0).getPipeline().getNodesInOrder());
+    assertEquals(nodes2, locations.get(1).getPipeline().getNodesInOrder());
+    // Sorted per pipeline, since the unsorted result is not cached.
+    verify(mockKeyManager, times(2)).sortDatanodesForWrite(any(), 
eq("1.2.3.4"), any());
+  }
+
+  @Test
+  public void testAllocateBlockKeepsOrderWhenRemoteAddressEmpty() throws 
Exception {
+    // Sort enabled and topology available, but the client has no remote 
address:
+    // OM must not sort, SCM receives an empty clientMachine, and the pipeline
+    // order is preserved.
+    List<DatanodeDetails> nodes = Arrays.asList(
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails());
+    Pipeline pipeline = Pipeline.newBuilder()
+        .setState(Pipeline.PipelineState.OPEN)
+        .setId(PipelineID.randomId())
+        .setReplicationConfig(
+            StandaloneReplicationConfig.getInstance(ReplicationFactor.THREE))
+        .setNodes(nodes)
+        .build();
+    AllocatedBlock block = new AllocatedBlock.Builder().setPipeline(pipeline)
+        .setContainerBlockID(new ContainerBlockID(CONTAINER_ID, 
LOCAL_ID)).build();
+    ArgumentCaptor<String> clientMachine = 
ArgumentCaptor.forClass(String.class);
+    when(scmBlockLocationProtocol.allocateBlock(anyLong(), anyInt(), any(),
+        anyString(), any(ExcludeList.class), clientMachine.capture()))
+        .thenReturn(Collections.singletonList(block));
+
+    KeyManager mockKeyManager = mock(KeyManager.class);
+    when(mockKeyManager.isSortDatanodesForWriteEnabled()).thenReturn(true);
+    when(ozoneManager.getKeyManager()).thenReturn(mockKeyManager);
+    
when(ozoneManager.getClusterMapAllowNull()).thenReturn(mock(NetworkTopology.class));
+
+    OMAllocateBlockRequest request =
+        getOmAllocateBlockRequest(createAllocateBlockRequest());
+    List<OmKeyLocationInfo> locations = 
request.allocateBlock(replicationConfig,
+        new ExcludeList(), scmBlockSize, true,
+        UserInfo.newBuilder().setRemoteAddress("").build(), ozoneManager);
+
+    assertEquals("", clientMachine.getValue());
+    verify(mockKeyManager, never()).sortDatanodesForWrite(any(), anyString(), 
any());
+    assertEquals(1, locations.size());
+    // Assert the write order (nodesInOrder), which copyWithNodesInOrder would
+    // have changed had OM sorted; it must stay as the original pipeline order.
+    assertEquals(nodes, locations.get(0).getPipeline().getNodesInOrder());
+  }
+
+  @Test
+  public void sortDatanodesForWriteRequiresClientMachine() {
+    List<DatanodeDetails> nodes = Arrays.asList(
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails());
+    assertThrows(IllegalArgumentException.class,
+        () -> keyManager.sortDatanodesForWrite(nodes, "", 
mock(NetworkTopology.class)));
+  }
+
+  // Like createAllocateBlockRequest, but sets sortDatanodes so preExecute
+  // resolves the client address from the RPC context.
+  private OMRequest createAllocateBlockRequestWithSort() {
+    KeyArgs keyArgs = KeyArgs.newBuilder()
+        
.setVolumeName(volumeName).setBucketName(bucketName).setKeyName(keyName)
+        .setFactor(((RatisReplicationConfig) 
replicationConfig).getReplicationFactor())
+        .setType(replicationConfig.getReplicationType())
+        .setSortDatanodes(true)
+        .build();
+    AllocateBlockRequest allocateBlockRequest = 
AllocateBlockRequest.newBuilder()
+        .setClientID(clientID).setKeyArgs(keyArgs).build();
+    return OMRequest.newBuilder()
+        .setCmdType(OzoneManagerProtocolProtos.Type.AllocateBlock)
+        .setClientId(UUID.randomUUID().toString())
+        .setAllocateBlockRequest(allocateBlockRequest).build();
+  }
+
+  // Run preExecute with a mocked RPC context so UserInfo carries 
clientAddress,
+  // the way an OM RPC handler thread would see it.
+  private void preExecuteWithClient(OMAllocateBlockRequest request, String 
clientAddress) throws Exception {
+    InetAddress clientIp = InetAddress.getByAddress(clientAddress, 
InetAddress.getByName(clientAddress).getAddress());
+    UserGroupInformation ugi = UserGroupInformation.getCurrentUser();
+    try (MockedStatic<Server> mockedRpcServer = mockStatic(Server.class)) {
+      mockedRpcServer.when(Server::getRemoteUser).thenReturn(ugi);
+      mockedRpcServer.when(Server::getRemoteIp).thenReturn(clientIp);
+      request.preExecute(ozoneManager);
+    }
+  }
+
   protected OMRequest createAllocateBlockRequest() {
 
     KeyArgs keyArgs = KeyArgs.newBuilder()


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


Reply via email to