Skip to content
Draft
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,9 +440,7 @@ 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;
List<String> excludedNodesForCapacity = new ArrayList<>();

// When affinity node is null, in this case new node to be selected
// should be in different rack than used nodes rack.
Expand All @@ -453,84 +449,82 @@ private DatanodeDetails chooseNode(List<DatanodeDetails> excludedNodes,
// 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());
for (DatanodeDetails node : excludedNodes) {
excludedNodesForCapacity.add(node.getNetworkFullPath());
}
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);
}

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 except the racks of excludedNodes (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<>();
boolean outOfRetries = false;
for (Step step : steps) {
int retries = MAX_RETRY;
while (retries > 0) {
metrics.incrDatanodeChooseAttemptCount();
DatanodeDetails node = (DatanodeDetails) networkTopology.chooseRandom(NetConstants.ROOT,
excludedNodesForCapacity, excludedNodes, 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--;
}
// Decommissioned and maintenance nodes stay in the topology, so a place can use up all its tries on them,
// e.g. the only other rack while a whole rack is being decommissioned. Then move on to the next place.
outOfRetries = retries == 0;
}
if (outOfRetries) {
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);
}
throw new SCMException("No satisfied datanode to meet the excludedNodes and affinityNode constrains.", null);
}

maxRetry--;
if (maxRetry == 0) {
// avoid the infinite loop
String errMsg = "No satisfied datanode to meet the space constrains. "
+ "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());
/** 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 excludedNodes. 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 @@ -29,6 +29,7 @@
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos;
import org.apache.hadoop.hdds.scm.ContainerPlacementStatus;
import org.apache.hadoop.hdds.scm.PlacementPolicy;
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
Expand Down Expand Up @@ -133,8 +134,9 @@ public int processAndSendCommands(
List<DatanodeDetails> usedDns = replicas.stream()
.map(ContainerReplica::getDatanodeDetails)
.collect(Collectors.toList());
if (containerPlacement.validateContainerPlacement(usedDns,
usedDns.size()).isPolicySatisfied()) {
ContainerPlacementStatus placementStatus =
containerPlacement.validateContainerPlacement(usedDns, usedDns.size());
if (placementStatus.isPolicySatisfied()) {
LOG.info("Container {} is currently not misreplicated",
container.getContainerID());
return 0;
Expand All @@ -158,6 +160,21 @@ public int processAndSendCommands(
excludedAndUsedNodes.getUsedNodes(),
excludedAndUsedNodes.getExcludedNodes(), currentContainerSize,
container);
if (!targetDatanodes.isEmpty()) {
// When no other rack has a usable node, placement can fall back to a rack that already has a replica.
// Copying there doesn't help: the container is still mis-replicated, the extra copy gets deleted, and the
// next run copies it again, over and over.
List<DatanodeDetails> dnsAfterCopy = new ArrayList<>(usedDns);
dnsAfterCopy.addAll(targetDatanodes);
ContainerPlacementStatus placementAfterCopy =
containerPlacement.validateContainerPlacement(dnsAfterCopy, usedDns.size());
if (!placementAfterCopy.isPolicySatisfied()
&& placementAfterCopy.misReplicationCount() >= placementStatus.misReplicationCount()) {
LOG.info("Not copying container {} to {}, as that would not reduce its mis-replication.",
container.getContainerID(), targetDatanodes);
return 0;
}
}
List<DatanodeDetails> availableSources = sources.stream()
.map(ContainerReplica::getDatanodeDetails)
.collect(Collectors.toList());
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,67 @@ 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)));
}

@Test
public void fallbackToSameRackWhenOtherRackHasOnlyDecommissionedNodes() throws SCMException {
setup(2 * NODE_PER_RACK);
// All of rack1 is being decommissioned. node5 has a replica that is being decommissioned, so it is excluded,
// and nodes 6-9 are already decommissioned.
dnInfos.get(5).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONING, HEALTHY));
for (int i = 6; i < 10; i++) {
dnInfos.get(i).setNodeStatus(NodeStatus.valueOf(DECOMMISSIONED, HEALTHY));
}
// The other two replicas are on rack0.
List<DatanodeDetails> usedNodes = new ArrayList<>(Arrays.asList(datanodes.get(0), datanodes.get(1)));
List<DatanodeDetails> excludedNodes = new ArrayList<>(Arrays.asList(datanodes.get(5)));

List<DatanodeDetails> datanodeDetails = policy.chooseDatanodes(usedNodes, excludedNodes, null, 1, 0, 5);

assertEquals(1, datanodeDetails.size());
// rack1 has no usable node, so the new replica has to go to rack0 as well.
assertTrue(cluster.isSameParent(datanodes.get(0), datanodeDetails.get(0)));
assertThat(usedNodes).doesNotContain(datanodeDetails.get(0));
assertEquals(1, metrics.getDatanodeChooseFallbackCount());
// Without fallback, there is nowhere to put it.
assertThrows(SCMException.class, () -> policyNoFallback.chooseDatanodes(usedNodes, excludedNodes, null, 1, 0, 5));
}

@Test
public void chooseSingleNodeRackWithUsedAndExcludeNodes()
throws SCMException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -125,8 +125,12 @@ static PlacementPolicy mockPlacementPolicy() {
PlacementPolicy placementPolicy = mock(PlacementPolicy.class);
ContainerPlacementStatus mockedContainerPlacementStatus = mock(ContainerPlacementStatus.class);
when(mockedContainerPlacementStatus.isPolicySatisfied()).thenReturn(false);
when(placementPolicy.validateContainerPlacement(anyList(),
anyInt())).thenReturn(mockedContainerPlacementStatus);
// With the new targets added, the placement is fine.
ContainerPlacementStatus placementAfterCopy = mock(ContainerPlacementStatus.class);
when(placementAfterCopy.isPolicySatisfied()).thenReturn(true);
when(placementPolicy.validateContainerPlacement(anyList(), anyInt())).thenAnswer(invocation ->
invocation.<List<DatanodeDetails>>getArgument(0).size() > invocation.<Integer>getArgument(1)
? placementAfterCopy : mockedContainerPlacementStatus);
return placementPolicy;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,19 @@ public void testMisReplicationWithPendingOps()
pendingOp, 0, 1, 0);
}

@Test
public void testMisReplicationWithTargetsNotReducingIt() throws IOException {
Set<ContainerReplica> availableReplicas = ReplicationTestUtil
.createReplicas(Pair.of(IN_SERVICE, 0), Pair.of(IN_SERVICE, 0), Pair.of(IN_SERVICE, 0));
// The target doesn't help, e.g. because it is on the same rack as all the existing replicas.
PlacementPolicy placementPolicy = mock(PlacementPolicy.class);
ContainerPlacementStatus mockedContainerPlacementStatus = mock(ContainerPlacementStatus.class);
when(mockedContainerPlacementStatus.isPolicySatisfied()).thenReturn(false);
when(mockedContainerPlacementStatus.misReplicationCount()).thenReturn(1);
when(placementPolicy.validateContainerPlacement(anyList(), anyInt())).thenReturn(mockedContainerPlacementStatus);
testMisReplication(availableReplicas, placementPolicy, Collections.emptyList(), 0, 1, 1, 0);
}

@Test
public void testAllSourcesOverloaded() throws IOException {
ReplicationManager replicationManager = getReplicationManager();
Expand Down
Loading