Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -442,95 +440,82 @@ private DatanodeDetails chooseNode(List<DatanodeDetails> excludedNodes,
List<DatanodeDetails> affinityNodes, List<DatanodeDetails> usedNodes,
long metadataSizeRequired,
long dataSizeRequired) throws SCMException {
int ancestorGen = RACK_LEVEL;
int maxRetry = MAX_RETRY;
List<String> 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<String> 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<Step> 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;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<DatanodeDetails> usedNodes = new ArrayList<>(Arrays.asList(datanodes.get(5), datanodes.get(0)));

List<DatanodeDetails> 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<DatanodeDetails> usedNodes = new ArrayList<>(Arrays.asList(datanodes.get(0), datanodes.get(5)));
List<DatanodeDetails> excludedNodes = new ArrayList<>(Arrays.asList(datanodes.get(1), datanodes.get(6)));

List<DatanodeDetails> 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)));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we also assert metrics.getDatanodeChooseFallbackCount() is 1 here and 0 in the preceding test, since trying another affinity rack should not count as a fallback?

}

@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<DatanodeDetails> usedNodes =
new ArrayList<>(Arrays.asList(datanodes.get(0), datanodes.get(5), datanodes.get(15)));
List<DatanodeDetails> excludedNodes = new ArrayList<>(Arrays.asList(datanodes.get(10)));

List<DatanodeDetails> 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 {
Expand Down
Loading