Skip to content
Merged
Show file tree
Hide file tree
Changes from 14 commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
f6b0d03
HDDS-7492. Placement Policy Interface changes to handle misreplicatio…
swamirishi Nov 28, 2022
daeddb4
HDDS-7492. Address Review Comments
swamirishi Nov 28, 2022
b4b891b
HDDS-7492. Address Review Comments
swamirishi Nov 29, 2022
7d36041
HDDS-7492. Revert replication config addition to function
swamirishi Nov 29, 2022
7efb3ee
HDDS-7492. Remove Function ReplicasToRemove
swamirishi Nov 29, 2022
fd3ed13
HDDS-7492. Simplify toCopyFunction
swamirishi Nov 29, 2022
209429e
HDDS-7492. Fix Checkstyle Issues
swamirishi Nov 29, 2022
ce6654d
HDDS-7492. Address review comments
swamirishi Dec 2, 2022
045fbb8
HDDS-7492. Address review comments
swamirishi Dec 2, 2022
0e107f7
HDDS-7492. Adding test cases
swamirishi Dec 2, 2022
b261b17
HDDS-7492. Fix Checkstyle Issues
swamirishi Dec 2, 2022
96ece79
HDDS-7492. Add testcases
swamirishi Dec 2, 2022
923d940
HDDS-7492. Update Javadoc
swamirishi Dec 5, 2022
d3f052e
HDDS-7492. Add Max Replicas per Rack in SCM CommonPlacement Policy fo…
swamirishi Dec 5, 2022
8f023aa
HDDS-7492. Change algorithm to support max replicas per rack simplify…
swamirishi Dec 6, 2022
0d231eb
HDDS-7492. Fix bug to get non zero denominator for number of required…
swamirishi Dec 7, 2022
15472d0
HDDS-7492. Fix testcases
swamirishi Dec 7, 2022
26244fa
HDDS-7492. Fix Algorithm to remove max number of replicas
swamirishi Dec 7, 2022
20620e7
HDDS-7492. Fix Pair Import Issue
swamirishi Dec 7, 2022
f62044d
HDDS-7492. Fix number of racks required while validating container pl…
swamirishi Dec 7, 2022
75963bd
HDDS-7492. Fix Rack Scatter Policy to support max replicas per rack
swamirishi Dec 8, 2022
50648fa
HDDS-7492. Fix issue in rack scatter
swamirishi Dec 8, 2022
7eb77ce
HDDS-7492. Fix max replica per rack for pipeline placement policy & r…
swamirishi Dec 8, 2022
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 @@ -22,12 +22,13 @@
import java.io.IOException;
import java.util.Collections;
import java.util.List;
import java.util.Set;

/**
* A PlacementPolicy support choosing datanodes to build
* pipelines or containers with specified constraints.
*/
public interface PlacementPolicy {
public interface PlacementPolicy<Replica> {
Comment thread
swamirishi marked this conversation as resolved.

default List<DatanodeDetails> chooseDatanodes(
List<DatanodeDetails> excludedNodes,
Expand Down Expand Up @@ -60,9 +61,17 @@ List<DatanodeDetails> chooseDatanodes(List<DatanodeDetails> usedNodes,
* Given a list of datanode and the number of replicas required, return
* a PlacementPolicyStatus object indicating if the container meets the
* placement policy - ie is it on the correct number of racks, etc.
* @param dns List of datanodes holding a replica of the container
* @param dns List of replica holding a replica of the container
* @param replicas The expected number of replicas
*/
ContainerPlacementStatus validateContainerPlacement(
List<DatanodeDetails> dns, int replicas);
List<DatanodeDetails> dns, int replicas);

/**
* Given a set of replicas of a container which are
* neither over underreplicated nor overreplicated,
* return a set of replicas to copy to another node to fix misreplication.
* @param replicas
*/
Set<Replica> replicasToCopyToFixMisreplication(Set<Replica> replicas);
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,35 +18,43 @@
package org.apache.hadoop.hdds.scm;


import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
import java.util.Random;
import java.util.stream.Collectors;

import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
import com.google.common.collect.Sets;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.MetadataStorageReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.StorageReportProto;
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
import org.apache.hadoop.hdds.scm.container.placement.algorithms.ContainerPlacementStatusDefault;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.net.NetworkTopology;
import org.apache.hadoop.hdds.scm.net.Node;
import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.node.NodeStatus;

import com.google.common.annotations.VisibleForTesting;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.PriorityQueue;
import java.util.Queue;
import java.util.Random;
import java.util.Set;
import java.util.stream.Collectors;

/**
* This policy implements a set of invariants which are common
* for all basic placement policies, acts as the repository of helper
* functions which are common to placement policies.
*/
public abstract class SCMCommonPlacementPolicy implements PlacementPolicy {
public abstract class SCMCommonPlacementPolicy implements
PlacementPolicy<ContainerReplica> {
@VisibleForTesting
static final Logger LOG =
LoggerFactory.getLogger(SCMCommonPlacementPolicy.class);
Expand Down Expand Up @@ -346,6 +354,21 @@ protected int getRequiredRackCount(int numReplicas) {
return 1;
}

/**
* Default implementation to return the max number of replicas per rack.
* For simple policies that are not rack aware
* we return numReplicas, from this default implementation.
*
* @param numReplicas - The desired replica counts
* @return The max number of replicas per rack
*/
protected int getMaxReplicasPerRack(int numReplicas) {
return numReplicas / getRequiredRackCount(numReplicas)
+ Math.min(numReplicas % getRequiredRackCount(numReplicas), 1);
}



/**
* This default implementation handles rack aware policies and non rack
* aware policies. If a future placement policy needs to check more than racks
Expand All @@ -364,6 +387,7 @@ public ContainerPlacementStatus validateContainerPlacement(
List<DatanodeDetails> dns, int replicas) {
NetworkTopology topology = nodeManager.getClusterNetworkTopologyMap();
int requiredRacks = getRequiredRackCount(replicas);
int maxReplicasPerRack = getMaxReplicasPerRack(replicas);
if (topology == null || replicas == 1 || requiredRacks == 1) {
if (dns.size() > 0) {
// placement is always satisfied if there is at least one DN.
Expand All @@ -378,16 +402,17 @@ public ContainerPlacementStatus validateContainerPlacement(
// The leaf nodes are all at max level, so the number of nodes at
// leafLevel - 1 is the rack count
numRacks = topology.getNumOfNodes(maxLevel - 1);
final long currentRackCount = dns.stream()
.map(d -> topology.getAncestor(d, 1))
.distinct()
.count();
Map<Node, Long> currentRackCount = dns.stream()
.collect(Collectors.groupingBy(this::getPlacementGroup,
Collectors.counting()));

if (replicas < requiredRacks) {
requiredRacks = replicas;
}
return new ContainerPlacementStatusDefault(
(int)currentRackCount, requiredRacks, numRacks);
currentRackCount.size(), requiredRacks, numRacks, maxReplicasPerRack,
currentRackCount.values().stream().map(Long::intValue)
.collect(Collectors.toList()));
}

/**
Expand Down Expand Up @@ -426,4 +451,67 @@ public boolean isValidNode(DatanodeDetails datanodeDetails,
}
return false;
}

/**
* Given a set of replicas of a container which are
* neither over underreplicated nor overreplicated,
* return a set of replicas to copy to another node to fix misreplication.
* @param replicas
*/
@Override
public Set<ContainerReplica> replicasToCopyToFixMisreplication(
Set<ContainerReplica> replicas) {
Map<Node, List<ContainerReplica>> placementGroupReplicaIdMap
= replicas.stream().collect(Collectors.groupingBy(replica ->
this.getPlacementGroup(replica.getDatanodeDetails())));

int totalNumberOfReplicas = replicas.size();
int requiredNumberOfPlacementGroups =
getRequiredRackCount(totalNumberOfReplicas);
int additionalNumberOfRacksRequired = Math.max(
requiredNumberOfPlacementGroups - placementGroupReplicaIdMap.size(),
0);
int replicasPerPlacementGroup =
getMaxReplicasPerRack(totalNumberOfReplicas);
Set<ContainerReplica> copyReplicaSet = Sets.newHashSet();

for (List<ContainerReplica> replicaList: placementGroupReplicaIdMap
.values()) {
if (replicaList.size() > replicasPerPlacementGroup) {
List<ContainerReplica> replicasToBeCopied = replicaList.stream()
.limit(replicaList.size() - replicasPerPlacementGroup)
.collect(Collectors.toList());
copyReplicaSet.addAll(replicasToBeCopied);
replicaList.removeAll(replicasToBeCopied);
}
}
if (additionalNumberOfRacksRequired > copyReplicaSet.size()) {
Comment thread
swamirishi marked this conversation as resolved.
Outdated
additionalNumberOfRacksRequired -= copyReplicaSet.size();
Queue<List<ContainerReplica>> placementGroupReplicas =
new PriorityQueue<>((o1, o2) ->
Integer.compare(o2.size(), o1.size()));
placementGroupReplicas.addAll(placementGroupReplicaIdMap.values());
while (placementGroupReplicas.size() > 0
&& additionalNumberOfRacksRequired > 0) {
List<ContainerReplica> replicaList = placementGroupReplicas.poll();
int numberOfReplicasToBeCopied = Math.max(1,
Math.min(replicaList.size()
- Optional.ofNullable(placementGroupReplicas.peek())
.map(List::size).orElse(0),
additionalNumberOfRacksRequired));
List<ContainerReplica> replicasToBeCopied = replicaList.stream()
.limit(numberOfReplicasToBeCopied)
.collect(Collectors.toList());
copyReplicaSet.addAll(replicasToBeCopied);
replicaList.removeAll(replicasToBeCopied);
placementGroupReplicas.add(replicaList);
additionalNumberOfRacksRequired -= replicasToBeCopied.size();
}
}
return copyReplicaSet;
}

protected Node getPlacementGroup(DatanodeDetails dn) {
return nodeManager.getClusterNetworkTopologyMap().getAncestor(dn, 1);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,9 @@

import org.apache.hadoop.hdds.scm.ContainerPlacementStatus;

import java.util.Collections;
import java.util.List;

/**
* Simple Status object to check if a container is replicated across enough
* racks.
Expand All @@ -29,33 +32,56 @@ public class ContainerPlacementStatusDefault
private final int currentRacks;
private final int totalRacks;

private final int maxReplicasPerRack;
private final List<Integer> rackReplicaCnts;


public ContainerPlacementStatusDefault(int currentRacks, int requiredRacks,
int totalRacks) {
int totalRacks, int maxReplicasPerRack, List<Integer> rackReplicaCnts) {
this.requiredRacks = requiredRacks;
this.currentRacks = currentRacks;
this.totalRacks = totalRacks;
this.maxReplicasPerRack = maxReplicasPerRack;
this.rackReplicaCnts = rackReplicaCnts;
}

public ContainerPlacementStatusDefault(int requiredRacks, int currentRacks,
int totalRacks) {
this(requiredRacks, currentRacks, totalRacks, 1,
currentRacks == 0 ? Collections.emptyList()
: Collections.nCopies(currentRacks, 1));
}

@Override
public boolean isPolicySatisfied() {
return currentRacks >= totalRacks || currentRacks >= requiredRacks;
if (currentRacks < Math.min(totalRacks, requiredRacks)) {
return false;
}
return rackReplicaCnts.stream().allMatch(cnt -> cnt <= maxReplicasPerRack);
}

@Override
public String misReplicatedReason() {
if (isPolicySatisfied()) {
return null;
}
return "The container is mis-replicated as it is on " + currentRacks +
" racks but should be on " + requiredRacks + " racks.";
if (currentRacks < Math.min(requiredRacks, maxReplicasPerRack)) {
return "The container is mis-replicated as it is on " + currentRacks +
" racks but should be on " + requiredRacks + " racks.";
}
return "The container is mis-replicated as max number of replicas per rack "
+ "is " + maxReplicasPerRack + " but number of replicas per rack" +
" are " + rackReplicaCnts.toString();
}

@Override
public int misReplicationCount() {
if (isPolicySatisfied()) {
return 0;
}
return requiredRacks - currentRacks;
return Math.max(requiredRacks - currentRacks,
rackReplicaCnts.stream().mapToInt(
cnt -> Math.max(maxReplicasPerRack - cnt, 0)).sum());
Comment on lines +82 to +84

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.

Trying to figure out what this part is doing. Should it be cnt - maxReplicasPerRack instead? @swamirishi @sodonnel

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

yeah you are right

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks for pointing this out. Let me fix this.

}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.scm.PlacementPolicy;
import org.apache.hadoop.hdds.scm.SCMCommonPlacementPolicy;
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.net.NetworkTopology;
import org.apache.hadoop.hdds.scm.node.NodeManager;
Expand All @@ -40,7 +41,7 @@
* can be practically used.
*/
public final class SCMContainerPlacementRandom extends SCMCommonPlacementPolicy
implements PlacementPolicy {
implements PlacementPolicy<ContainerReplica> {
@VisibleForTesting
public static final Logger LOG =
LoggerFactory.getLogger(SCMContainerPlacementRandom.class);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,11 @@ private boolean isNonClosedRatisThreePipeline(Pipeline p) {
&& !p.isClosed();
}

@Override
protected int getMaxReplicasPerRack(int numReplicas) {
return numReplicas - 1;
}

/**
* Filter out viable nodes based on
* 1. nodes that are healthy
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.hadoop.hdds.scm;

import com.google.common.base.Preconditions;
import com.google.common.collect.Sets;
import org.apache.commons.lang3.RandomUtils;
import org.apache.hadoop.hdds.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
Expand Down Expand Up @@ -84,7 +85,6 @@
import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.UUID;
Expand Down Expand Up @@ -691,11 +691,20 @@ public static Set<ContainerReplica> getReplicas(
}

public static Set<ContainerReplica> getReplicas(
final ContainerID containerId,
final ContainerReplicaProto.State state,
final long sequenceId,
final DatanodeDetails... datanodeDetails) {
return Sets.newHashSet(getReplicas(containerId, state, sequenceId,
Arrays.asList(datanodeDetails)));
}

public static List<ContainerReplica> getReplicas(
final ContainerID containerId,
final ContainerReplicaProto.State state,
final long sequenceId,
final DatanodeDetails... datanodeDetails) {
Set<ContainerReplica> replicas = new HashSet<>();
final Iterable<DatanodeDetails> datanodeDetails) {
List<ContainerReplica> replicas = new ArrayList<>();
for (DatanodeDetails datanode : datanodeDetails) {
replicas.add(getReplicas(containerId, state,
sequenceId, datanode.getUuid(), datanode));
Expand Down Expand Up @@ -744,14 +753,14 @@ public static ContainerReplica getReplicas(
return builder.build();
}

public static Set<ContainerReplica> getReplicasWithReplicaIndex(
public static List<ContainerReplica> getReplicasWithReplicaIndex(
final ContainerID containerId,
final ContainerReplicaProto.State state,
final long usedBytes,
final long keyCount,
final long sequenceId,
final DatanodeDetails... datanodeDetails) {
Set<ContainerReplica> replicas = new HashSet<>();
final Iterable<DatanodeDetails> datanodeDetails) {
List<ContainerReplica> replicas = new ArrayList<>();
int replicaIndex = 1;
for (DatanodeDetails datanode : datanodeDetails) {
replicas.add(getReplicaBuilder(containerId, state,
Expand All @@ -762,6 +771,17 @@ public static Set<ContainerReplica> getReplicasWithReplicaIndex(
return replicas;
}

public static Set<ContainerReplica> getReplicasWithReplicaIndex(
final ContainerID containerId,
final ContainerReplicaProto.State state,
final long usedBytes,
final long keyCount,
final long sequenceId,
final DatanodeDetails... datanodeDetails) {
return Sets.newHashSet(getReplicasWithReplicaIndex(containerId, state,
usedBytes, keyCount, sequenceId, Arrays.asList(datanodeDetails)));
}



public static Pipeline getRandomPipeline() {
Expand Down
Loading