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]