diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRackAware.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRackAware.java index 7bcc5fb8926d..9f3bd557fa79 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRackAware.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRackAware.java @@ -24,14 +24,12 @@ import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.stream.Collectors; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.scm.SCMCommonPlacementPolicy; import org.apache.hadoop.hdds.scm.exceptions.SCMException; import org.apache.hadoop.hdds.scm.net.NetConstants; import org.apache.hadoop.hdds.scm.net.NetworkTopology; -import org.apache.hadoop.hdds.scm.net.Node; import org.apache.hadoop.hdds.scm.node.NodeManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -442,95 +440,82 @@ private DatanodeDetails chooseNode(List excludedNodes, List affinityNodes, List usedNodes, long metadataSizeRequired, long dataSizeRequired) throws SCMException { - int ancestorGen = RACK_LEVEL; - int maxRetry = MAX_RETRY; - List excludedNodesForCapacity = null; - - // When affinity node is null, in this case new node to be selected - // should be in different rack than used nodes rack. - // Exclude nodes should be just excluded from topology node selection, - // which is filled in excludedNodesForCapacity - // Used node rack should not be part of rack selection - // which is filled in excludedNodes - if (affinityNodes == null && excludedNodes != null) { - excludedNodesForCapacity = excludedNodes.stream() - .map(DatanodeDetails::getNetworkFullPath) - .collect(Collectors.toList()); - excludedNodes = usedNodes; - } - - boolean isFallbacked = false; - while (true) { - metrics.incrDatanodeChooseAttemptCount(); - DatanodeDetails node = null; - if (affinityNodes != null) { - for (Node affinityNode : affinityNodes) { - node = (DatanodeDetails)networkTopology.chooseRandom( - NetConstants.ROOT, excludedNodesForCapacity, excludedNodes, - affinityNode, ancestorGen); - if (node != null) { - break; - } - } - } else { - node = (DatanodeDetails)networkTopology.chooseRandom(NetConstants.ROOT, - excludedNodesForCapacity, excludedNodes, null, ancestorGen); + List excludedNodesForCapacity = new ArrayList<>(); + if (excludedNodes != null) { + for (DatanodeDetails node : excludedNodes) { + excludedNodesForCapacity.add(node.getNetworkFullPath()); } + } - if (node == null) { - // cannot find the node which meets all constrains - LOG.warn("Failed to find the datanode for container. excludedNodes:" + - (excludedNodes == null ? "" : excludedNodes.toString()) + - ", affinityNode:" + - (affinityNodes == null ? "" : affinityNodes.stream() - .map(Object::toString).collect(Collectors.joining(", ")))); - if (fallback) { - isFallbacked = true; - // fallback, don't consider the affinity node - if (affinityNodes != null) { - affinityNodes = null; - continue; - } - // fallback, don't consider cross rack - if (ancestorGen == RACK_LEVEL) { - ancestorGen--; - continue; - } - } - // there is no constrains to reduce or fallback is true - throw new SCMException("No satisfied datanode to meet the" + - " excludedNodes and affinityNode constrains.", null); + // The places to look for a node, in order. Each one gets MAX_RETRY tries: + // 1. the rack of each affinity node; + // 2. any rack that none of usedNodes is on (right away if there are no affinity nodes, otherwise only as a + // fallback); + // 3. any rack at all, as a fallback. + List steps = new ArrayList<>(); + if (affinityNodes != null) { + for (DatanodeDetails affinityNode : affinityNodes) { + steps.add(new Step(affinityNode, RACK_LEVEL, false)); } + } + if (affinityNodes == null || fallback) { + steps.add(new Step(null, RACK_LEVEL, affinityNodes != null)); + } + if (fallback) { + steps.add(new Step(null, 0, true)); + } - if (usedNodes != null && usedNodes.contains(node)) { - if (excludedNodesForCapacity == null) { - excludedNodesForCapacity = new ArrayList<>(); + for (Step step : steps) { + int retries = MAX_RETRY; + while (retries > 0) { + metrics.incrDatanodeChooseAttemptCount(); + DatanodeDetails node = (DatanodeDetails) networkTopology.chooseRandom(NetConstants.ROOT, + excludedNodesForCapacity, step.affinityNode == null ? usedNodes : null, step.affinityNode, + step.ancestorGen); + if (node == null) { + LOG.warn("Failed to find the datanode for container. excludedNodes: {}, affinityNode: {}", + excludedNodes, step.affinityNode); + break; } excludedNodesForCapacity.add(node.getNetworkFullPath()); - continue; - } - - if (isValidNode(node, metadataSizeRequired, dataSizeRequired)) { - metrics.incrDatanodeChooseSuccessCount(); - if (isFallbacked) { - metrics.incrDatanodeChooseFallbackCount(); + if (usedNodes != null && usedNodes.contains(node)) { + continue; + } + if (isValidNode(node, metadataSizeRequired, dataSizeRequired)) { + metrics.incrDatanodeChooseSuccessCount(); + if (step.fallback) { + metrics.incrDatanodeChooseFallbackCount(); + } + return node; } - return node; + retries--; } - - maxRetry--; - if (maxRetry == 0) { - // avoid the infinite loop - String errMsg = "No satisfied datanode to meet the space constrains. " + // Decommissioned and maintenance nodes stay in the topology, so a rack can use up all its tries on them. + // If that happens on an affinity node's rack, move on to the next place. Anywhere else, give up. + if (retries == 0 && step.affinityNode == null) { + String errMsg = "No writable datanode with enough space was found. " + "metadata size required: " + metadataSizeRequired + " data size required: " + dataSizeRequired; LOG.info(errMsg); throw new SCMException(errMsg, null); } - if (excludedNodesForCapacity == null) { - excludedNodesForCapacity = new ArrayList<>(); - } - excludedNodesForCapacity.add(node.getNetworkFullPath()); + } + throw new SCMException("No satisfied datanode to meet the excludedNodes and affinityNode constrains.", null); + } + + /** One place where chooseNode looks for a node. */ + private static final class Step { + /** Look near this node, or anywhere when null. */ + private final DatanodeDetails affinityNode; + /** RACK_LEVEL: stay on affinityNode's rack, or, without one, avoid the racks of usedNodes. 0: any rack. */ + private final int ancestorGen; + /** Whether a node found here counts as a fallback in the metrics. */ + private final boolean fallback; + + Step(DatanodeDetails affinityNode, int ancestorGen, boolean fallback) { + this.affinityNode = affinityNode; + this.ancestorGen = ancestorGen; + this.fallback = fallback; } } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java index 5b1961df1071..8886dbc09aad 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java @@ -18,6 +18,7 @@ package org.apache.hadoop.hdds.scm.container.placement.algorithms; import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONING; import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.HEALTHY; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_DATANODE_RATIS_VOLUME_FREE_SPACE_MIN; @@ -678,6 +679,65 @@ public void chooseNodeWithUsedNodesMultipleRack(int datanodeCount) cluster.isSameParent(datanodes.get(5), datanodeDetails.get(0))); } + @Test + public void chooseNodeWhenUsedNodeRackHasOnlyDecommissionedNodes() throws SCMException { + setup(3 * NODE_PER_RACK); + // rack1 has no usable node left: node5 already has a replica, and nodes 6-9 are decommissioned. + for (int i = 6; i < 10; i++) { + dnInfos.get(i).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONED, HEALTHY)); + } + // The container has replicas on node5 (rack1) and node0 (rack0). rack1 is tried first. + List usedNodes = new ArrayList<>(Arrays.asList(datanodes.get(5), datanodes.get(0))); + + List datanodeDetails = policy.chooseDatanodes(usedNodes, new ArrayList<>(), null, 1, 0, 5); + + assertEquals(1, datanodeDetails.size()); + // Instead of giving up after rack1, placement tries rack0 and finds a node there. + assertTrue(cluster.isSameParent(datanodes.get(0), datanodeDetails.get(0))); + assertNotEquals(datanodes.get(0), datanodeDetails.get(0)); + } + + @Test + public void fallbackWhenUsedNodeRacksHaveOnlyDecommissionedNodes() throws SCMException { + setup(3 * NODE_PER_RACK); + // The replicas are on node0 (rack0) and node5 (rack1), and neither rack has another usable node: + // node1 and node6 are excluded because their replicas are being decommissioned, and the rest are decommissioned. + dnInfos.get(1).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONING, HEALTHY)); + for (int i : new int[] {2, 3, 4, 6, 7, 8, 9}) { + dnInfos.get(i).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONED, HEALTHY)); + } + List usedNodes = new ArrayList<>(Arrays.asList(datanodes.get(0), datanodes.get(5))); + List excludedNodes = new ArrayList<>(Arrays.asList(datanodes.get(1), datanodes.get(6))); + + List datanodeDetails = policy.chooseDatanodes(usedNodes, excludedNodes, null, 1, 0, 5); + + assertEquals(1, datanodeDetails.size()); + // After both racks fail, placement falls back to rack2, the only rack with usable nodes. + assertTrue(cluster.isSameParent(datanodes.get(10), datanodeDetails.get(0))); + } + + @Test + public void fallbackDoesNotSkipRackOfExcludedNode() throws SCMException { + setup(4 * NODE_PER_RACK); + // The replicas are on node0 (rack0), node5 (rack1) and node15 (rack3), and those racks have no other usable + // node. node10 on rack2 has a replica that is being decommissioned, so it is excluded. + for (int i : new int[] {1, 2, 3, 4, 6, 7, 8, 9, 16, 17, 18, 19}) { + dnInfos.get(i).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONED, HEALTHY)); + } + dnInfos.get(10).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONING, HEALTHY)); + List usedNodes = + new ArrayList<>(Arrays.asList(datanodes.get(0), datanodes.get(5), datanodes.get(15))); + List excludedNodes = new ArrayList<>(Arrays.asList(datanodes.get(10))); + + List datanodeDetails = policy.chooseDatanodes(usedNodes, excludedNodes, null, 1, 0, 5); + + assertEquals(1, datanodeDetails.size()); + // The fallback looks for a rack with no replica on it. That is rack2, and the excluded node there must not + // rule out the whole rack. + assertTrue(cluster.isSameParent(datanodes.get(10), datanodeDetails.get(0))); + assertNotEquals(datanodes.get(10), datanodeDetails.get(0)); + } + @Test public void chooseSingleNodeRackWithUsedAndExcludeNodes() throws SCMException {