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 @@ -17,25 +17,31 @@

package org.apache.hadoop.hdds.scm.container.balancer;

import static org.apache.hadoop.util.StringUtils.byteDesc;

import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.conf.StorageUnit;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto;
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
import org.apache.hadoop.ozone.OzoneConsts;

/**
* Orchestrates cluster analysis, estimation and recommendation for container balancer.
*/
public final class ContainerBalancerAdvisor {

private static final long MIN_DELETE_PHASE_MILLIS = Duration.ofMinutes(9).toMillis();
private static final double PLANNING_ITERATION_BUFFER = 1.3d;
private static final String DATANODE_OFFSET_KEY =
"hdds.scm.replication.event.timeout.datanode.offset";

Expand Down Expand Up @@ -82,6 +88,187 @@ public static List<ContainerBalancerEstimation> estimate(OzoneConfiguration conf
return Collections.unmodifiableList(estimations);
}

/**
* Recommends balancer configuration for one or more balancer profiles.
*
* <p>If {@link AdvisorRequest#profile} is set, returns a result for that profile only.
* Otherwise returns SLOW, MEDIUM, and FAST. Per-profile validation failures are returned with
* {@link ContainerBalancerRecommendation#succeeded()} false instead of aborting other profiles.
*/
public static List<ContainerBalancerRecommendation> recommend(
OzoneConfiguration conf, AdvisorRequest request) {
Objects.requireNonNull(conf, "conf");
Objects.requireNonNull(request, "request");
List<DatanodeUsageInfoProto> nodes = Objects.requireNonNull(request.nodes, "nodes");

ContainerBalancerConfiguration balancerConfig = conf.getObject(ContainerBalancerConfiguration.class);

double thresholdPercent = request.thresholdPercent != null
? request.thresholdPercent
: balancerConfig.getThreshold();
validateThresholdPercent(thresholdPercent);
double thresholdRatio = thresholdPercent / 100.0;
Set<String> includeNodes = request.includeNodes != null
? request.includeNodes
: balancerConfig.getIncludeNodes();
Set<String> excludeNodes = request.excludeNodes != null
? request.excludeNodes
: balancerConfig.getExcludeNodes();

ContainerBalancerClusterSnapshot snapshot = ContainerBalancerClusterAnalyzer.analyze(
nodes, thresholdRatio, includeNodes, excludeNodes);

List<ContainerBalancerProfile> profiles = selectProfilesForRecommend(request);
if (isClusterAlreadyBalanced(snapshot)) {
return recommendationsForBalancedCluster(profiles, thresholdPercent);
}
validateSnapshotForEstimation(snapshot);
List<ContainerBalancerRecommendation> recommendations = new ArrayList<>(profiles.size());
for (ContainerBalancerProfile profile : profiles) {
recommendations.add(recommendForProfile(
conf, request, profile, snapshot, balancerConfig, thresholdPercent));
}
return Collections.unmodifiableList(recommendations);
}

private static boolean isClusterAlreadyBalanced(ContainerBalancerClusterSnapshot snapshot) {
return snapshot.getBytesToMove() <= 0
|| (snapshot.getSourceCount() < 1 && snapshot.getTargetCount() < 1);
}

private static List<ContainerBalancerRecommendation> recommendationsForBalancedCluster(
List<ContainerBalancerProfile> profiles, double thresholdPercent) {
List<ContainerBalancerRecommendation> recommendations = new ArrayList<>(profiles.size());
for (ContainerBalancerProfile profile : profiles) {
recommendations.add(ContainerBalancerRecommendation.newBuilder()
.setProfile(profile)
.setThresholdPercent(thresholdPercent)
.setClusterBalanced(true)
.build());
}
return Collections.unmodifiableList(recommendations);
}

private static ContainerBalancerRecommendation recommendForProfile(
OzoneConfiguration conf,
AdvisorRequest request,
ContainerBalancerProfile profile,
ContainerBalancerClusterSnapshot snapshot,
ContainerBalancerConfiguration balancerConfig,
double thresholdPercent) {

ContainerBalancerRecommendation.Builder builder = ContainerBalancerRecommendation.newBuilder()
.setProfile(profile)
.setThresholdPercent(thresholdPercent);

try {
long moveReplicationTimeoutMillis = balancerConfig.getMoveReplicationTimeout().toMillis();
long moveTimeoutMillis = balancerConfig.getMoveTimeout().toMillis();
long balancingIntervalMillis = balancerConfig.getBalancingInterval().toMillis();
validateMoveTimeouts(conf, moveReplicationTimeoutMillis, moveTimeoutMillis);
validateBalancingIntervalMillis(balancingIntervalMillis);

ContainerBalancerEstimation estimation = estimateForProfile(
conf, request, profile, snapshot, balancerConfig, thresholdPercent);
if (!estimation.succeeded()) {
throw new IllegalArgumentException(estimation.getFailureMessage());
}

long configCeiling = balancerConfig.getMaxSizeToMovePerIteration();
long recommendedMaxMove = computeRecommendedMaxSizeToMove(configCeiling, estimation);

validateResolvedMoveLimits(
conf,
estimation.getMaxSizeEnteringTarget(),
estimation.getMaxSizeLeavingSource(),
recommendedMaxMove);

int recommendedIterations = (int) Math.ceil(
estimation.getEstimatedIterations() * PLANNING_ITERATION_BUFFER);

Map<String, String> rationale = buildRationale(
request.thresholdPercent != null,
profile,
balancerConfig,
estimation,
moveTimeoutMillis,
moveReplicationTimeoutMillis,
balancingIntervalMillis);

return builder
.setMaxDatanodesPercentage(estimation.getMaxDatanodesPercentage())
.setMaxSizeToMovePerIteration(recommendedMaxMove)
.setMaxSizeEnteringTarget(estimation.getMaxSizeEnteringTarget())
.setMaxSizeLeavingSource(estimation.getMaxSizeLeavingSource())
.setMoveTimeoutMillis(moveTimeoutMillis)
.setMoveReplicationTimeoutMillis(moveReplicationTimeoutMillis)
.setBalancingIntervalMillis(balancingIntervalMillis)
.setRecommendedIterations(recommendedIterations)
.setRationale(rationale)
.setEstimation(estimation)
.build();
} catch (IllegalArgumentException e) {
return builder.setFailureMessage(e.getMessage()).build();
}
}

/**
* Recommended max-size-to-move for one profile:
* min(config ceiling, max(per-iteration estimate, profile entering, profile leaving)).
*/
static long computeRecommendedMaxSizeToMove(
long configCeiling, ContainerBalancerEstimation estimation) {
long floor = maxPositive(
estimation.getPerIterationBytes(),
estimation.getMaxSizeEnteringTarget(),
estimation.getMaxSizeLeavingSource());
return minPositive(configCeiling, floor);
}

private static Map<String, String> buildRationale(
boolean userProvidedThreshold,
ContainerBalancerProfile profile,
ContainerBalancerConfiguration balancerConfig,
ContainerBalancerEstimation estimation,
long moveTimeoutMillis,
long moveReplicationTimeoutMillis,
long balancingIntervalMillis) {
Map<String, String> rationale = new HashMap<>();
rationale.put("threshold", userProvidedThreshold
? "user override"
: "default from configuration");
int profileDatanodesPercentage = profile.getDatanodesMaxPercentage();
int usedDatanodesPercentage = estimation.getMaxDatanodesPercentage();
if (usedDatanodesPercentage == profileDatanodesPercentage) {
rationale.put("maxDatanodesPercentage", String.format(
"profile default (%d%%)", usedDatanodesPercentage));
} else {
rationale.put("maxDatanodesPercentage", String.format(
"profile %d%%, raised to %d%% (≥2 nodes)", profileDatanodesPercentage, usedDatanodesPercentage));
}
rationale.put("maxSizeToMovePerIteration", String.format(
"min(%s ceiling, max(per-iteration %s, entering %s, leaving %s))",
byteDesc(balancerConfig.getMaxSizeToMovePerIteration()),
byteDesc(estimation.getPerIterationBytes()),
byteDesc(estimation.getMaxSizeEnteringTarget()),
byteDesc(estimation.getMaxSizeLeavingSource())));
rationale.put("maxSizeEnteringTarget", String.format(
"profile default (%d GB)", estimation.getMaxSizeEnteringTarget() / OzoneConsts.GB));
rationale.put("maxSizeLeavingSource", String.format(
"profile default (%d GB)", estimation.getMaxSizeLeavingSource() / OzoneConsts.GB));
rationale.put("moveTimeout", String.format(
"configuration default (%d min)", moveTimeoutMillis / 60000));
rationale.put("moveReplicationTimeout", String.format(
"configuration default (%d min)", moveReplicationTimeoutMillis / 60000));
rationale.put("balancingInterval", String.format(
"configuration default (%d min)", balancingIntervalMillis / 60000));
long planningIterations = (long) Math.ceil(
estimation.getEstimatedIterations() * PLANNING_ITERATION_BUFFER);
rationale.put("iterations", String.format(
"includes +30%% buffer (planning estimate: %d)", planningIterations));
return rationale;
}

private static ContainerBalancerEstimation estimateForProfile(OzoneConfiguration conf, AdvisorRequest request,
ContainerBalancerProfile profile, ContainerBalancerClusterSnapshot snapshot,
ContainerBalancerConfiguration balancerConfig, double thresholdPercent) {
Expand Down Expand Up @@ -167,6 +354,7 @@ private static ContainerBalancerEstimation estimateForProfile(OzoneConfiguration
.build();
} catch (IllegalArgumentException e) {
return builder
.setThresholdPercent(thresholdPercent)
.setMaxDatanodesPercentage(maxDatanodesPercentage)
.setFailureMessage(e.getMessage())
.build();
Expand Down Expand Up @@ -245,7 +433,17 @@ private static long minPositive(long... values) {
return result;
}

private static void validateSnapshotForEstimation(ContainerBalancerClusterSnapshot snapshot) {
private static long maxPositive(long... values) {
long result = 0;
for (long value : values) {
if (value > result) {
result = value;
}
}
return result;
}

private static void validateSnapshotForEstimation(ContainerBalancerClusterSnapshot snapshot ) {
if (snapshot.getSourceCount() < 1) {
throw new IllegalArgumentException("No over-utilized datanodes (sources) found.");
}
Expand Down Expand Up @@ -351,6 +549,16 @@ private static List<ContainerBalancerProfile> selectProfiles(AdvisorRequest requ
return Collections.singletonList(ContainerBalancerProfile.MEDIUM);
}

private static List<ContainerBalancerProfile> selectProfilesForRecommend(AdvisorRequest request) {
if (request.profile != null) {
return Collections.singletonList(request.profile);
}
return Arrays.asList(
ContainerBalancerProfile.SLOW,
ContainerBalancerProfile.MEDIUM,
ContainerBalancerProfile.FAST);
}

/**
* Input for {@link ContainerBalancerAdvisor}: cluster usage data and optional overrides.
* Unset fields fall back to {@link ContainerBalancerConfiguration} or profile presets.
Expand Down
Loading