From b4f223c21d32863de8b5e222a9571b29aac7187d Mon Sep 17 00:00:00 2001 From: sravani-revuri Date: Sun, 27 Sep 2026 21:41:45 +0530 Subject: [PATCH 1/5] HDDS-16181. Add container balancer recommend CLI command to suggest config values for slow/medium/fast profiles --- .../balancer/ContainerBalancerAdvisor.java | 213 +++++++++++++++- .../ContainerBalancerRecommendation.java | 209 ++++++++++++++++ .../TestContainerBalancerAdvisor.java | 129 ++++++++++ .../scm/cli/ContainerBalancerCommands.java | 25 +- .../ContainerBalancerRecommendSubcommand.java | 228 ++++++++++++++++++ .../TestContainerBalancerSubCommand.java | 65 +++++ 6 files changed, 861 insertions(+), 8 deletions(-) create mode 100644 hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java create mode 100644 hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java index b7ba0b025cf2..cc9a8081ec34 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java @@ -21,7 +21,10 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Locale; +import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.concurrent.TimeUnit; @@ -29,6 +32,8 @@ 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; +import static org.apache.hadoop.util.StringUtils.byteDesc; /** * Orchestrates cluster analysis, estimation and recommendation for container balancer. @@ -36,6 +41,7 @@ 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"; @@ -72,7 +78,7 @@ public static List estimate(OzoneConfiguration conf ContainerBalancerClusterSnapshot snapshot = ContainerBalancerClusterAnalyzer.analyze(nodes, thresholdRatio, includeNodes, excludeNodes); - validateSnapshotForEstimation(snapshot); + validateSnapshotForEstimation(snapshot, conf); List profiles = selectProfiles(request); List estimations = new ArrayList<>(profiles.size()); @@ -82,6 +88,178 @@ public static List estimate(OzoneConfiguration conf return Collections.unmodifiableList(estimations); } + /** @see #estimate(OzoneConfiguration, AdvisorRequest) */ + public static List estimateDryRun(OzoneConfiguration conf, AdvisorRequest request) { + return estimate(conf, request); + } + + /** + * Recommends balancer configuration for one or more balancer profiles. + * + *

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 recommend( + OzoneConfiguration conf, AdvisorRequest request) { + Objects.requireNonNull(conf, "conf"); + Objects.requireNonNull(request, "request"); + List 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 includeNodes = request.includeNodes != null + ? request.includeNodes + : balancerConfig.getIncludeNodes(); + Set excludeNodes = request.excludeNodes != null + ? request.excludeNodes + : balancerConfig.getExcludeNodes(); + + ContainerBalancerClusterSnapshot snapshot = ContainerBalancerClusterAnalyzer.analyze( + nodes, thresholdRatio, includeNodes, excludeNodes); + validateSnapshotForEstimation(snapshot, conf); + + List profiles = selectProfilesForRecommend(request); + List recommendations = new ArrayList<>(profiles.size()); + for (ContainerBalancerProfile profile : profiles) { + recommendations.add(recommendForProfile( + conf, request, profile, balancerConfig, thresholdPercent, includeNodes, excludeNodes)); + } + return Collections.unmodifiableList(recommendations); + } + + private static ContainerBalancerRecommendation recommendForProfile( + OzoneConfiguration conf, + AdvisorRequest request, + ContainerBalancerProfile profile, + ContainerBalancerConfiguration balancerConfig, + double thresholdPercent, + Set includeNodes, + Set excludeNodes) { + + 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); + + AdvisorRequest estimateRequest = new AdvisorRequest() + .setNodes(request.nodes) + .setThresholdPercent(thresholdPercent) + .setIncludeNodes(includeNodes) + .setExcludeNodes(excludeNodes) + .setProfile(profile); + + ContainerBalancerEstimation estimation = estimate(conf, estimateRequest).get(0); + 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 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 buildRationale( + boolean userProvidedThreshold, + ContainerBalancerProfile profile, + ContainerBalancerConfiguration balancerConfig, + ContainerBalancerEstimation estimation, + long moveTimeoutMillis, + long moveReplicationTimeoutMillis, + long balancingIntervalMillis) { + Map 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(Locale.ENGLISH, + "profile default (%d%%)", usedDatanodesPercentage)); + } else { + rationale.put("maxDatanodesPercentage", String.format(Locale.ENGLISH, + "profile %d%%, raised to %d%% (≥2 nodes)", profileDatanodesPercentage, usedDatanodesPercentage)); + } + rationale.put("maxSizeToMovePerIteration", String.format(Locale.ENGLISH, + "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(Locale.ENGLISH, + "profile default (%d GB)", estimation.getMaxSizeEnteringTarget() / OzoneConsts.GB)); + rationale.put("maxSizeLeavingSource", String.format(Locale.ENGLISH, + "profile default (%d GB)", estimation.getMaxSizeLeavingSource() / OzoneConsts.GB)); + rationale.put("moveTimeout", String.format(Locale.ENGLISH, + "configuration default (%d min)", moveTimeoutMillis / 60000)); + rationale.put("moveReplicationTimeout", String.format(Locale.ENGLISH, + "configuration default (%d min)", moveReplicationTimeoutMillis / 60000)); + rationale.put("balancingInterval", String.format(Locale.ENGLISH, + "configuration default (%d min)", balancingIntervalMillis / 60000)); + long planningIterations = (long) Math.ceil( + estimation.getEstimatedIterations() * PLANNING_ITERATION_BUFFER); + rationale.put("iterations", String.format(Locale.ENGLISH, + "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) { @@ -167,6 +345,7 @@ private static ContainerBalancerEstimation estimateForProfile(OzoneConfiguration .build(); } catch (IllegalArgumentException e) { return builder + .setThresholdPercent(thresholdPercent) .setMaxDatanodesPercentage(maxDatanodesPercentage) .setFailureMessage(e.getMessage()) .build(); @@ -245,7 +424,22 @@ 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, + OzoneConfiguration conf) { + long containerSizeBytes = (long) conf.getStorageSize( + ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE, + ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT, + StorageUnit.BYTES); if (snapshot.getSourceCount() < 1) { throw new IllegalArgumentException("No over-utilized datanodes (sources) found."); } @@ -255,6 +449,11 @@ private static void validateSnapshotForEstimation(ContainerBalancerClusterSnapsh if (snapshot.getBytesToMove() <= 0) { throw new IllegalArgumentException("No bytes to move."); } + if (snapshot.getBytesToMove() < containerSizeBytes) { + throw new IllegalArgumentException( + "Bytes to move (" + snapshot.getBytesToMove() + + ") is less than container size (" + containerSizeBytes + ")."); + } if (snapshot.getTotalEligibleDatanodes() < 2) { throw new IllegalArgumentException(String.format( "Container Balancer found %d eligible datanode(s) but requires at least 2.", @@ -351,6 +550,16 @@ private static List selectProfiles(AdvisorRequest requ return Collections.singletonList(ContainerBalancerProfile.MEDIUM); } + private static List 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. diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java new file mode 100644 index 000000000000..eed12de14396 --- /dev/null +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java @@ -0,0 +1,209 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.hdds.scm.container.balancer; + +import java.util.Collections; +import java.util.Map; +import java.util.Objects; + +/** + * Recommended balancer configuration for a single profile. + */ +public final class ContainerBalancerRecommendation { + + private final ContainerBalancerProfile profile; + private final String failureMessage; + private final double thresholdPercent; + private final int maxDatanodesPercentage; + private final long maxSizeToMovePerIteration; + private final long maxSizeEnteringTarget; + private final long maxSizeLeavingSource; + private final long moveTimeoutMillis; + private final long moveReplicationTimeoutMillis; + private final long balancingIntervalMillis; + private final int recommendedIterations; + private final Map rationale; + private final ContainerBalancerEstimation estimation; + + private ContainerBalancerRecommendation(Builder b) { + this.profile = Objects.requireNonNull(b.profile, "profile == null"); + this.failureMessage = b.failureMessage; + this.thresholdPercent = b.thresholdPercent; + this.maxDatanodesPercentage = b.maxDatanodesPercentage; + this.maxSizeToMovePerIteration = b.maxSizeToMovePerIteration; + this.maxSizeEnteringTarget = b.maxSizeEnteringTarget; + this.maxSizeLeavingSource = b.maxSizeLeavingSource; + this.moveTimeoutMillis = b.moveTimeoutMillis; + this.moveReplicationTimeoutMillis = b.moveReplicationTimeoutMillis; + this.balancingIntervalMillis = b.balancingIntervalMillis; + this.recommendedIterations = b.recommendedIterations; + this.rationale = b.rationale == null + ? Collections.emptyMap() + : Collections.unmodifiableMap(b.rationale); + this.estimation = b.estimation; + } + + public static Builder newBuilder() { + return new Builder(); + } + + public ContainerBalancerProfile getProfile() { + return profile; + } + + public boolean succeeded() { + return failureMessage == null; + } + + public String getFailureMessage() { + return failureMessage; + } + + public double getThresholdPercent() { + return thresholdPercent; + } + + public int getMaxDatanodesPercentage() { + return maxDatanodesPercentage; + } + + public long getMaxSizeToMovePerIteration() { + return maxSizeToMovePerIteration; + } + + public long getMaxSizeEnteringTarget() { + return maxSizeEnteringTarget; + } + + public long getMaxSizeLeavingSource() { + return maxSizeLeavingSource; + } + + public long getMoveTimeoutMillis() { + return moveTimeoutMillis; + } + + public long getMoveReplicationTimeoutMillis() { + return moveReplicationTimeoutMillis; + } + + public long getBalancingIntervalMillis() { + return balancingIntervalMillis; + } + + public int getRecommendedIterations() { + return recommendedIterations; + } + + public Map getRationale() { + return rationale; + } + + public ContainerBalancerEstimation getEstimation() { + return estimation; + } + + /** Builder for {@link ContainerBalancerRecommendation}. */ + public static final class Builder { + private ContainerBalancerProfile profile; + private String failureMessage; + private double thresholdPercent; + private int maxDatanodesPercentage; + private long maxSizeToMovePerIteration; + private long maxSizeEnteringTarget; + private long maxSizeLeavingSource; + private long moveTimeoutMillis; + private long moveReplicationTimeoutMillis; + private long balancingIntervalMillis; + private int recommendedIterations; + private Map rationale; + private ContainerBalancerEstimation estimation; + + private Builder() { + } + + public Builder setProfile(ContainerBalancerProfile profileValue) { + this.profile = profileValue; + return this; + } + + public Builder setFailureMessage(String message) { + this.failureMessage = message; + return this; + } + + public Builder setThresholdPercent(double threshold) { + this.thresholdPercent = threshold; + return this; + } + + public Builder setMaxDatanodesPercentage(int percentage) { + this.maxDatanodesPercentage = percentage; + return this; + } + + public Builder setMaxSizeToMovePerIteration(long bytes) { + this.maxSizeToMovePerIteration = bytes; + return this; + } + + public Builder setMaxSizeEnteringTarget(long bytes) { + this.maxSizeEnteringTarget = bytes; + return this; + } + + public Builder setMaxSizeLeavingSource(long bytes) { + this.maxSizeLeavingSource = bytes; + return this; + } + + public Builder setMoveTimeoutMillis(long millis) { + this.moveTimeoutMillis = millis; + return this; + } + + public Builder setMoveReplicationTimeoutMillis(long millis) { + this.moveReplicationTimeoutMillis = millis; + return this; + } + + public Builder setBalancingIntervalMillis(long millis) { + this.balancingIntervalMillis = millis; + return this; + } + + public Builder setRecommendedIterations(int iterations) { + this.recommendedIterations = iterations; + return this; + } + + public Builder setRationale(Map rationaleMap) { + this.rationale = rationaleMap; + return this; + } + + public Builder setEstimation(ContainerBalancerEstimation estimationValue) { + this.estimation = estimationValue; + return this; + } + + public ContainerBalancerRecommendation build() { + return new ContainerBalancerRecommendation(this); + } + } +} diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java index 411183f45048..dcf55f7a9c29 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java @@ -20,7 +20,10 @@ import static org.apache.hadoop.ozone.ClientVersion.DEFAULT_VERSION; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; import java.util.ArrayList; import java.util.List; @@ -34,6 +37,19 @@ /** Tests for {@link ContainerBalancerAdvisor} estimation. */ public final class TestContainerBalancerAdvisor { + @Test + void testComputePerIterationBytesLimitedByEnteringTarget() { + int[] involved = {7, 7}; + long expected = 26L * OzoneConsts.GB * 7; + + assertEquals(expected, ContainerBalancerAdvisor.computePerIterationBytes( + expected * 10, + 500L * OzoneConsts.GB, + 26 * OzoneConsts.GB, + 26 * OzoneConsts.GB, + involved)); + } + @Test void testComputePerIterationBytesNeverExceedsBytesToMove() { int[] involved = {3, 3}; @@ -240,6 +256,119 @@ void testEstimateFailsWhenClusterBalanced() { new ContainerBalancerAdvisor.AdvisorRequest().setNodes(balanced))); } + @Test + void testRecommendReturnsThreeProfiles() { + OzoneConfiguration conf = new OzoneConfiguration(); + List results = ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest().setNodes(buildCluster(70, 14, 14))); + + assertEquals(3, results.size()); + assertEquals(ContainerBalancerProfile.SLOW, results.get(0).getProfile()); + assertEquals(ContainerBalancerProfile.MEDIUM, results.get(1).getProfile()); + assertEquals(ContainerBalancerProfile.FAST, results.get(2).getProfile()); + for (ContainerBalancerRecommendation result : results) { + assertTrue(result.succeeded()); + assertNotNull(result.getEstimation()); + assertTrue(result.getEstimation().succeeded()); + assertTrue(result.getRecommendedIterations() >= result.getEstimation().getEstimatedIterations()); + assertFalse(result.getRationale().isEmpty()); + } + } + + @Test + void testComputeRecommendedMaxSizeToMoveFloorsAtProfileLimits() { + ContainerBalancerEstimation estimation = ContainerBalancerEstimation.newBuilder() + .setProfile(ContainerBalancerProfile.SLOW) + .setPerIterationBytes(8L * OzoneConsts.GB) + .setMaxSizeEnteringTarget(10L * OzoneConsts.GB) + .setMaxSizeLeavingSource(10L * OzoneConsts.GB) + .build(); + + assertEquals(10L * OzoneConsts.GB, + ContainerBalancerAdvisor.computeRecommendedMaxSizeToMove(500L * OzoneConsts.GB, estimation)); + } + + @Test + void testComputeRecommendedMaxSizeToMoveUsesPerIterationWhenLarger() { + ContainerBalancerEstimation estimation = ContainerBalancerEstimation.newBuilder() + .setProfile(ContainerBalancerProfile.MEDIUM) + .setPerIterationBytes(182L * OzoneConsts.GB) + .setMaxSizeEnteringTarget(26L * OzoneConsts.GB) + .setMaxSizeLeavingSource(26L * OzoneConsts.GB) + .build(); + + assertEquals(182L * OzoneConsts.GB, + ContainerBalancerAdvisor.computeRecommendedMaxSizeToMove(500L * OzoneConsts.GB, estimation)); + } + + @Test + void testRecommendReturnsSingleProfileWhenProfileSet() { + OzoneConfiguration conf = new OzoneConfiguration(); + List results = ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest() + .setNodes(buildCluster(70, 14, 14)) + .setProfile(ContainerBalancerProfile.MEDIUM)); + + assertEquals(1, results.size()); + assertEquals(ContainerBalancerProfile.MEDIUM, results.get(0).getProfile()); + assertTrue(results.get(0).succeeded()); + } + + @Test + void testRecommendUsesProfileDefaults() { + OzoneConfiguration conf = new OzoneConfiguration(); + ContainerBalancerRecommendation slow = ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest().setNodes(buildCluster(70, 14, 14))) + .get(0); + + assertEquals(10, slow.getMaxDatanodesPercentage()); + assertEquals(10L * OzoneConsts.GB, slow.getMaxSizeEnteringTarget()); + assertEquals(10L * OzoneConsts.GB, slow.getMaxSizeLeavingSource()); + assertEquals(30L * OzoneConsts.GB, slow.getMaxSizeToMovePerIteration()); + assertTrue(slow.getMaxSizeToMovePerIteration() >= slow.getMaxSizeEnteringTarget()); + assertTrue(slow.getMaxSizeToMovePerIteration() >= slow.getMaxSizeLeavingSource()); + } + + @Test + void testRecommendRespectsThresholdOverride() { + OzoneConfiguration conf = new OzoneConfiguration(); + List nodes = buildCluster(70, 14, 14); + + long defaultBytesToMove = ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest().setNodes(nodes)) + .get(0) + .getEstimation() + .getBytesToMove(); + + long tighterBytesToMove = ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest() + .setNodes(nodes) + .setThresholdPercent(5.0)) + .get(0) + .getEstimation() + .getBytesToMove(); + + assertTrue(tighterBytesToMove > defaultBytesToMove); + } + + @Test + void testRecommendFailsWhenClusterBalanced() { + OzoneConfiguration conf = new OzoneConfiguration(); + List balanced = new ArrayList<>(); + balanced.add(proto("dn-1", OzoneConsts.TB, (long) (0.70 * OzoneConsts.TB))); + balanced.add(proto("dn-2", OzoneConsts.TB, (long) (0.70 * OzoneConsts.TB))); + + assertThrows(IllegalArgumentException.class, () -> + ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest().setNodes(balanced))); + } + @Test void testEstimateFailsWhenEnteringTargetTooSmall() { OzoneConfiguration conf = new OzoneConfiguration(); diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java index 1c2cd364a51c..56df98dc55b0 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java @@ -78,16 +78,28 @@ * estimate for the SLOW profile only * ozone admin containerbalancer estimate --profile fast -t 5 * estimate FAST profile with a 5% threshold + * To recommend: + * ozone admin containerbalancer recommend + * [ --profile {@literal } ] + * [ -t/--threshold {@literal } ] + * [ --include-datanodes {@literal } ] + * [ --exclude-datanodes {@literal } ] + * Examples: + * ozone admin containerbalancer recommend + * recommend balancer configuration for SLOW, MEDIUM, and FAST profiles + * ozone admin containerbalancer recommend --profile slow + * recommend for the SLOW profile only + * ozone admin containerbalancer recommend -t 5 + * recommend with a 5% threshold * To stop: * ozone admin containerbalancer stop * * *

DESCRIPTION - *

The estimate subcommand fetches datanode usage from SCM and estimates from - * local configurations, profile presets and cluster analysis made. It does not start the balancer. Start does not yet - * support {@code --profile}, compare estimate profiles to the config you plan - * to pass on start, or wait until profile support is added to start. - * estimate subcommand produces upper-bound estimates. + *

Estimate and recommend fetch datanode usage from SCM and use local configuration and cluster + * analysis. They do not start the balancer. Start does not yet support {@code --profile}; compare + * estimate/recommend profiles to the config you plan to pass on start. Estimate produces upper-bound + * duration estimates. *

The threshold parameter is a fraction in the range of (1%, 100%) with a * default value of 10%. The threshold sets a target for whether the cluster * is balanced. A cluster is balanced if for each datanode, the utilization @@ -112,7 +124,8 @@ ContainerBalancerStartSubcommand.class, ContainerBalancerStopSubcommand.class, ContainerBalancerStatusSubcommand.class, - ContainerBalancerEstimateSubcommand.class + ContainerBalancerEstimateSubcommand.class, + ContainerBalancerRecommendSubcommand.class }) @MetaInfServices(AdminSubcommand.class) public class ContainerBalancerCommands implements AdminSubcommand { diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java new file mode 100644 index 000000000000..e49ef92adaa8 --- /dev/null +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java @@ -0,0 +1,228 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.hdds.scm.cli; + +import static org.apache.hadoop.util.StringUtils.byteDesc; + +import java.io.IOException; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; +import org.apache.commons.lang3.StringUtils; +import org.apache.hadoop.hdds.cli.HddsVersionProvider; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto; +import org.apache.hadoop.hdds.scm.client.ScmClient; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerAdvisor; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerEstimation; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerProfile; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerRecommendation; +import picocli.CommandLine.Command; +import picocli.CommandLine.Option; + +/** + * Recommends container balancer configuration for SLOW, MEDIUM, and FAST profiles + * without starting the balancer. + */ +@Command( + name = "recommend", + description = "Recommend container balancer configuration based on current cluster imbalance. " + + "When --profile is omitted, SLOW, MEDIUM, and FAST are recommended. Does not start the balancer.", + mixinStandardHelpOptions = true, + versionProvider = HddsVersionProvider.class) +public class ContainerBalancerRecommendSubcommand extends ScmSubcommand { + + private static final double PLANNING_ITERATION_BUFFER = 1.3d; + private static final int PARAM_COLUMN_WIDTH = 42; + private static final int VALUE_COLUMN_WIDTH = 14; + + @Option(names = {"-t", "--threshold"}, + description = "Percentage deviation from average utilization of " + + "the cluster after which a datanode will be rebalanced. The value " + + "should be in the range [0.0, 100.0), with a default of 10 " + + "(specify '10' for 10%%).") + private Optional threshold; + + @Option(names = {"--include-datanodes"}, + description = "A list of Datanode hostnames or ip addresses separated by commas. Only the " + + "Datanodes specified in this list are balanced.") + private Optional includeNodes; + + @Option(names = {"--exclude-datanodes"}, + description = "A list of Datanode hostnames or ip addresses separated by commas. The " + + "Datanodes specified in this list are excluded from balancing.") + private Optional excludeNodes; + + @Option(names = {"--profile"}, + description = "Throttling profile: SLOW, MEDIUM, or FAST. When set, only this profile is " + + "recommended. When omitted, SLOW, MEDIUM, and FAST are recommended.") + private Optional profileName; + + @Override + public void execute(ScmClient scmClient) throws IOException { + List nodes = scmClient.getDatanodeUsageInfo(true, Integer.MAX_VALUE); + if (nodes == null || nodes.isEmpty()) { + throw new IOException("No datanode usage information available from SCM."); + } + + OzoneConfiguration conf = getOzoneConf(); + ContainerBalancerAdvisor.AdvisorRequest request = buildRequest(nodes); + List recommendations; + try { + recommendations = ContainerBalancerAdvisor.recommend(conf, request); + } catch (IllegalArgumentException e) { + throw new IOException(e.getMessage(), e); + } + + boolean anySucceeded = false; + for (ContainerBalancerRecommendation recommendation : recommendations) { + printRecommendation(recommendation); + if (recommendation.succeeded()) { + anySucceeded = true; + } + } + if (!anySucceeded) { + throw new IOException(recommendations.get(0).getFailureMessage()); + } + } + + private ContainerBalancerAdvisor.AdvisorRequest buildRequest(List nodes) + throws IOException { + ContainerBalancerAdvisor.AdvisorRequest request = + new ContainerBalancerAdvisor.AdvisorRequest().setNodes(nodes); + threshold.ifPresent(request::setThresholdPercent); + includeNodes.ifPresent(value -> request.setIncludeNodes(parseNodeSet(value))); + excludeNodes.ifPresent(value -> request.setExcludeNodes(parseNodeSet(value))); + if (profileName.isPresent()) { + request.setProfile(parseProfile(profileName.get())); + } + return request; + } + + private static ContainerBalancerProfile parseProfile(String name) throws IOException { + try { + return ContainerBalancerProfile.valueOf(name.trim().toUpperCase(Locale.ENGLISH)); + } catch (IllegalArgumentException e) { + throw new IOException("Invalid profile: " + name + ". Expected SLOW, MEDIUM, or FAST."); + } + } + + private void printRecommendation(ContainerBalancerRecommendation recommendation) { + out().printf("RECOMMENDED CONFIGURATION (profile: %s)%n", recommendation.getProfile().name()); + out().println(); + + if (!recommendation.succeeded()) { + out().printf("Recommendation failed: %s%n%n", recommendation.getFailureMessage()); + return; + } + + printRecommendedParameters(recommendation); + out().println(); + out().println(" Estimation:"); + printEstimation(recommendation.getEstimation()); + } + + private void printRecommendedParameters(ContainerBalancerRecommendation recommendation) { + Map rationale = recommendation.getRationale(); + long moveTimeoutMinutes = Math.round(recommendation.getMoveTimeoutMillis() / 60000d); + long moveReplicationTimeoutMinutes = + Math.round(recommendation.getMoveReplicationTimeoutMillis() / 60000d); + long balancingIntervalMinutes = Math.round(recommendation.getBalancingIntervalMillis() / 60000d); + + out().println(" Recommended parameters:"); + printParameterRow("Parameter", "Value", "Rationale"); + printParameterRow("--threshold", + String.format(Locale.ENGLISH, "%.1f%%", recommendation.getThresholdPercent()), + rationale.get("threshold")); + printParameterRow("--max-datanodes-percentage-to-involve", + String.format(Locale.ENGLISH, "%d%%", recommendation.getMaxDatanodesPercentage()), + rationale.get("maxDatanodesPercentage")); + printParameterRow("--max-size-to-move-per-iteration-in-gb", + byteDesc(recommendation.getMaxSizeToMovePerIteration()), + rationale.get("maxSizeToMovePerIteration")); + printParameterRow("--max-size-entering-target-in-gb", + byteDesc(recommendation.getMaxSizeEnteringTarget()) + " / node", + rationale.get("maxSizeEnteringTarget")); + printParameterRow("--max-size-leaving-source-in-gb", + byteDesc(recommendation.getMaxSizeLeavingSource()) + " / node", + rationale.get("maxSizeLeavingSource")); + printParameterRow("--move-timeout-minutes", + String.format(Locale.ENGLISH, "%d min", moveTimeoutMinutes), + rationale.get("moveTimeout")); + printParameterRow("--move-replication-timeout-minutes", + String.format(Locale.ENGLISH, "%d min", moveReplicationTimeoutMinutes), + rationale.get("moveReplicationTimeout")); + printParameterRow("--balancing-iteration-interval-minutes", + String.format(Locale.ENGLISH, "%d min", balancingIntervalMinutes), + rationale.get("balancingInterval")); + printParameterRow("--iterations", + String.valueOf(recommendation.getRecommendedIterations()), + rationale.get("iterations")); + } + + private void printParameterRow(String parameter, String value, String rationaleText) { + String rationale = rationaleText == null ? "" : rationaleText; + out().printf(Locale.ENGLISH, " %-" + PARAM_COLUMN_WIDTH + "s %-" + VALUE_COLUMN_WIDTH + "s %s%n", + parameter, value, rationale); + } + + private void printEstimation(ContainerBalancerEstimation estimation) { + long estimatedIterations = estimation.getEstimatedIterations(); + long planningIterations = (long) Math.ceil(estimatedIterations * PLANNING_ITERATION_BUFFER); + long cycleTimeMillis = estimation.getMoveTimeoutMillis() + estimation.getBalancingIntervalMillis(); + long baseDurationMillis = estimation.getEstimatedDurationMillis(); + long planningDurationMillis = planningIterations * cycleTimeMillis; + out().printf(" Bytes to move: %s%n", byteDesc(estimation.getBytesToMove())); + out().printf(" Per iteration (estimate): ~%s%n", byteDesc(estimation.getPerIterationBytes())); + out().printf(" Estimated iterations: %d (planning estimate: %d, includes +30%% buffer)%n", + estimatedIterations, planningIterations); + out().printf(" Estimated duration: upper bound %s (planning estimate: %s, includes +30%% buffer)%n", + formatEstimatedDuration(baseDurationMillis), + formatEstimatedDuration(planningDurationMillis)); + out().println(" (assumes full move timeout + interval each cycle)"); + out().println(); + } + + private static String formatEstimatedDuration(long durationMillis) { + double days = durationMillis / 86400000d; + if (days >= 1) { + return String.format(Locale.ENGLISH, "~%.1f days", days); + } + double hours = durationMillis / 3600000d; + if (hours >= 1) { + return String.format(Locale.ENGLISH, "~%.1f hours", hours); + } + long minutes = durationMillis / 60000; + return String.format(Locale.ENGLISH, "~%d min", minutes); + } + + private static Set parseNodeSet(String nodes) { + if (StringUtils.isBlank(nodes)) { + return Collections.emptySet(); + } + return Arrays.stream(nodes.split(",")) + .map(String::trim) + .filter(s -> !s.isEmpty()) + .collect(Collectors.toSet()); + } +} diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java index 92aae28b2f36..73d488e6d1db 100644 --- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java +++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java @@ -41,6 +41,7 @@ import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoResponseProto; import org.apache.hadoop.hdds.scm.cli.ContainerBalancerEstimateSubcommand; +import org.apache.hadoop.hdds.scm.cli.ContainerBalancerRecommendSubcommand; import org.apache.hadoop.hdds.scm.cli.ContainerBalancerStartSubcommand; import org.apache.hadoop.hdds.scm.cli.ContainerBalancerStatusSubcommand; import org.apache.hadoop.hdds.scm.cli.ContainerBalancerStopSubcommand; @@ -171,6 +172,7 @@ class TestContainerBalancerSubCommand { private ContainerBalancerStartSubcommand startCmd; private ContainerBalancerStatusSubcommand statusCmd; private ContainerBalancerEstimateSubcommand estimateCmd; + private ContainerBalancerRecommendSubcommand recommendCmd; private GenericTestUtils.PrintStreamCapturer out; private GenericTestUtils.PrintStreamCapturer err; private AtomicBoolean verbose; @@ -382,6 +384,7 @@ protected boolean isVerbose() { }; parseSubcommand(startCmd); estimateCmd = new ContainerBalancerEstimateSubcommand(); + recommendCmd = new ContainerBalancerRecommendSubcommand(); out = GenericTestUtils.captureOut(); err = GenericTestUtils.captureErr(); } @@ -1032,6 +1035,68 @@ void testContainerBalancerEstimateSubcommandPartialFailureWhenMaxMoveOverrideCon .doesNotContain("Per iteration (estimate):"); } + @Test + void testContainerBalancerRecommendSubcommandDefaultShowsAllProfiles() throws IOException { + ScmClient scmClient = mock(ScmClient.class); + when(scmClient.getDatanodeUsageInfo(true, Integer.MAX_VALUE)) + .thenReturn(buildImbalancedCluster()); + + parseSubcommand(recommendCmd); + recommendCmd.execute(scmClient); + + String output = out.get(); + assertThat(output) + .contains("RECOMMENDED CONFIGURATION (profile: SLOW)") + .contains("RECOMMENDED CONFIGURATION (profile: MEDIUM)") + .contains("RECOMMENDED CONFIGURATION (profile: FAST)") + .contains("--threshold") + .contains("--iterations") + .contains(" Estimation:") + .contains("Bytes to move:") + .contains("planning estimate:") + .contains("assumes full move timeout + interval each cycle") + .doesNotContain("Recommendation failed:"); + } + + @Test + void testContainerBalancerRecommendSubcommandWithThresholdOverride() throws IOException { + ScmClient scmClient = mock(ScmClient.class); + when(scmClient.getDatanodeUsageInfo(true, Integer.MAX_VALUE)) + .thenReturn(buildImbalancedCluster()); + + parseSubcommand(recommendCmd, "-t", "5"); + recommendCmd.execute(scmClient); + + assertThat(out.get()).contains("RECOMMENDED CONFIGURATION (profile: SLOW)"); + } + + @Test + void testContainerBalancerRecommendSubcommandWithProfileShowsOneProfile() throws IOException { + ScmClient scmClient = mock(ScmClient.class); + when(scmClient.getDatanodeUsageInfo(true, Integer.MAX_VALUE)) + .thenReturn(buildImbalancedCluster()); + + parseSubcommand(recommendCmd, "--profile", "medium"); + recommendCmd.execute(scmClient); + + String output = out.get(); + assertThat(output) + .contains("RECOMMENDED CONFIGURATION (profile: MEDIUM)") + .doesNotContain("RECOMMENDED CONFIGURATION (profile: SLOW)") + .doesNotContain("RECOMMENDED CONFIGURATION (profile: FAST)"); + } + + @Test + void testContainerBalancerRecommendSubcommandInvalidThresholdFails() throws IOException { + ScmClient scmClient = mock(ScmClient.class); + when(scmClient.getDatanodeUsageInfo(true, Integer.MAX_VALUE)) + .thenReturn(buildImbalancedCluster()); + + parseSubcommand(recommendCmd, "-t", "-1"); + IOException ex = assertThrows(IOException.class, () -> recommendCmd.execute(scmClient)); + assertThat(ex.getMessage()).contains("Threshold should be specified in the range [0.0, 100.0)."); + } + /** * Imbalanced cluster for estimate CLI tests. * From b60228594b517f6f648eabb395c1318122981463 Mon Sep 17 00:00:00 2001 From: sravani-revuri Date: Sun, 27 Sep 2026 21:57:07 +0530 Subject: [PATCH 2/5] checkstyle --- .../hdds/scm/container/balancer/ContainerBalancerAdvisor.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java index cc9a8081ec34..73364317f644 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java @@ -17,6 +17,8 @@ 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; @@ -33,7 +35,6 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.ozone.OzoneConsts; -import static org.apache.hadoop.util.StringUtils.byteDesc; /** * Orchestrates cluster analysis, estimation and recommendation for container balancer. From 45f95076becf9ce3d1678262b03ce8ccb07bdccb Mon Sep 17 00:00:00 2001 From: sravani-revuri Date: Tue, 29 Sep 2026 11:19:25 +0530 Subject: [PATCH 3/5] added helper --- .../balancer/ContainerBalancerAdvisor.java | 5 -- .../scm/cli/ContainerBalancerCliHelper.java | 89 +++++++++++++++++++ .../scm/cli/ContainerBalancerCommands.java | 10 +-- .../cli/ContainerBalancerConfigOptions.java | 19 +--- .../ContainerBalancerEstimateSubcommand.java | 45 +--------- .../ContainerBalancerRecommendSubcommand.java | 63 +------------ .../cli/TestContainerBalancerSubCommand.java | 5 -- 7 files changed, 102 insertions(+), 134 deletions(-) create mode 100644 hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java index 73364317f644..f54038a73f2a 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java @@ -89,11 +89,6 @@ public static List estimate(OzoneConfiguration conf return Collections.unmodifiableList(estimations); } - /** @see #estimate(OzoneConfiguration, AdvisorRequest) */ - public static List estimateDryRun(OzoneConfiguration conf, AdvisorRequest request) { - return estimate(conf, request); - } - /** * Recommends balancer configuration for one or more balancer profiles. * diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java new file mode 100644 index 000000000000..868940d7d675 --- /dev/null +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java @@ -0,0 +1,89 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.hdds.scm.cli; + +import static org.apache.hadoop.util.StringUtils.byteDesc; + +import java.io.IOException; +import java.io.PrintWriter; +import java.util.Arrays; +import java.util.Collections; +import java.util.Locale; +import java.util.Set; +import java.util.stream.Collectors; +import org.apache.commons.lang3.StringUtils; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerEstimation; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerProfile; + +/** Shared parsing and CLI output for container balancer subcommands. */ +public final class ContainerBalancerCliHelper { + + /** Multiplier for recommended/planning iteration counts and CLI planning display (+30%). */ + public static final double PLANNING_ITERATION_BUFFER = 1.3d; + + private ContainerBalancerCliHelper() { + } + + public static ContainerBalancerProfile parseProfile(String name) throws IOException { + try { + return ContainerBalancerProfile.valueOf(name.trim().toUpperCase(Locale.ENGLISH)); + } catch (IllegalArgumentException e) { + throw new IOException("Invalid profile: " + name + ". Expected SLOW, MEDIUM, or FAST."); + } + } + + public static Set parseNodeSet(String nodes) { + if (StringUtils.isBlank(nodes)) { + return Collections.emptySet(); + } + return Arrays.stream(nodes.split(",")) + .map(String::trim) + .filter(s -> !s.isEmpty()) + .collect(Collectors.toSet()); + } + + public static void printEstimation(PrintWriter out, ContainerBalancerEstimation estimation) { + long estimatedIterations = estimation.getEstimatedIterations(); + long planningIterations = (long) Math.ceil(estimatedIterations * PLANNING_ITERATION_BUFFER); + long cycleTimeMillis = estimation.getMoveTimeoutMillis() + estimation.getBalancingIntervalMillis(); + long baseDurationMillis = estimation.getEstimatedDurationMillis(); + long planningDurationMillis = planningIterations * cycleTimeMillis; + out.printf(" Bytes to move: %s%n", byteDesc(estimation.getBytesToMove())); + out.printf(" Per iteration (estimate): ~%s%n", byteDesc(estimation.getPerIterationBytes())); + out.printf(" Estimated iterations: %d (planning estimate: %d, includes +30%% buffer)%n", + estimatedIterations, planningIterations); + out.printf(" Estimated duration: upper bound %s (planning estimate: %s, includes +30%% buffer)%n", + formatEstimatedDuration(baseDurationMillis), + formatEstimatedDuration(planningDurationMillis)); + out.println(" (assumes full move timeout + interval each cycle)"); + out.println(); + } + + public static String formatEstimatedDuration(long durationMillis) { + double days = durationMillis / 86400000d; + if (days >= 1) { + return String.format(Locale.ENGLISH, "~%.1f days", days); + } + double hours = durationMillis / 3600000d; + if (hours >= 1) { + return String.format(Locale.ENGLISH, "~%.1f hours", hours); + } + long minutes = durationMillis / 60000; + return String.format(Locale.ENGLISH, "~%d min", minutes); + } +} diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java index 8ebdd5d4d9c4..2ca6eaf0c207 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCommands.java @@ -54,6 +54,11 @@ * involved in balancing * ozone admin containerbalancer start -s 10 * start balancer with maximum size of 10GB to move in one iteration + * To assess: + * ozone admin containerbalancer assessment + * [ -t/--threshold {@literal }] + * [ --include-datanodes {@literal }] + * [ --exclude-datanodes {@literal }] * To estimate: * ozone admin containerbalancer estimate * [ --profile {@literal } ] @@ -93,11 +98,6 @@ * recommend with a 5% threshold * To stop: * ozone admin containerbalancer stop - * To assess: - * ozone admin containerbalancer assessment - * [ -t/--threshold {@literal }] - * [ --include-datanodes {@literal }] - * [ --exclude-datanodes {@literal }] * * *

DESCRIPTION diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerConfigOptions.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerConfigOptions.java index e82a954f2fd3..eb0a7bd0bb51 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerConfigOptions.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerConfigOptions.java @@ -18,12 +18,7 @@ package org.apache.hadoop.hdds.scm.cli; import java.time.Duration; -import java.util.Arrays; -import java.util.Collections; import java.util.Optional; -import java.util.Set; -import java.util.stream.Collectors; -import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerAdvisor; import org.apache.hadoop.ozone.OzoneConsts; import picocli.CommandLine.Option; @@ -159,17 +154,7 @@ public void applyToEstimateRequest(ContainerBalancerAdvisor.AdvisorRequest reque request.setMoveTimeoutMillis(Duration.ofMinutes(minutes).toMillis())); moveReplicationTimeout.ifPresent(minutes -> request.setMoveReplicationTimeoutMillis(Duration.ofMinutes(minutes).toMillis())); - includeNodes.ifPresent(value -> request.setIncludeNodes(parseNodeSet(value))); - excludeNodes.ifPresent(value -> request.setExcludeNodes(parseNodeSet(value))); - } - - private static Set parseNodeSet(String nodes) { - if (StringUtils.isBlank(nodes)) { - return Collections.emptySet(); - } - return Arrays.stream(nodes.split(",")) - .map(String::trim) - .filter(s -> !s.isEmpty()) - .collect(Collectors.toSet()); + includeNodes.ifPresent(value -> request.setIncludeNodes(ContainerBalancerCliHelper.parseNodeSet(value))); + excludeNodes.ifPresent(value -> request.setExcludeNodes(ContainerBalancerCliHelper.parseNodeSet(value))); } } diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java index 3e0b984b0508..a0f71791a2a7 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java @@ -29,7 +29,6 @@ import org.apache.hadoop.hdds.scm.client.ScmClient; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerAdvisor; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerEstimation; -import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerProfile; import picocli.CommandLine; import picocli.CommandLine.Command; import picocli.CommandLine.Option; @@ -47,8 +46,6 @@ versionProvider = HddsVersionProvider.class) public class ContainerBalancerEstimateSubcommand extends ScmSubcommand { - private static final double PLANNING_ITERATION_BUFFER = 1.3d; - @CommandLine.Mixin private ContainerBalancerConfigOptions configOptions; @@ -83,7 +80,7 @@ public void execute(ScmClient scmClient) throws IOException { printBasedOn(result); if (result.succeeded()) { anySucceeded = true; - printEstimation(result); + ContainerBalancerCliHelper.printEstimation(out(), result); } else if (multipleProfiles) { out().printf(" Estimation failed: %s%n%n", result.getFailureMessage()); } @@ -104,20 +101,12 @@ private ContainerBalancerAdvisor.AdvisorRequest buildRequest(List= 1) { - return String.format(Locale.ENGLISH, "~%.1f days", days); - } - double hours = durationMillis / 3600000d; - if (hours >= 1) { - return String.format(Locale.ENGLISH, "~%.1f hours", hours); - } - long minutes = durationMillis / 60000; - return String.format(Locale.ENGLISH, "~%d min", minutes); - } - static class ProfileSelection { @Option(names = {"--profile"}, description = "Throttling profile: SLOW, MEDIUM, or FAST profiles. When set, only this profile is estimated. " diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java index e49ef92adaa8..1601fa49943e 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java @@ -20,22 +20,15 @@ import static org.apache.hadoop.util.StringUtils.byteDesc; import java.io.IOException; -import java.util.Arrays; -import java.util.Collections; import java.util.List; import java.util.Locale; import java.util.Map; import java.util.Optional; -import java.util.Set; -import java.util.stream.Collectors; -import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hdds.cli.HddsVersionProvider; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto; import org.apache.hadoop.hdds.scm.client.ScmClient; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerAdvisor; -import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerEstimation; -import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerProfile; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerRecommendation; import picocli.CommandLine.Command; import picocli.CommandLine.Option; @@ -52,7 +45,6 @@ versionProvider = HddsVersionProvider.class) public class ContainerBalancerRecommendSubcommand extends ScmSubcommand { - private static final double PLANNING_ITERATION_BUFFER = 1.3d; private static final int PARAM_COLUMN_WIDTH = 42; private static final int VALUE_COLUMN_WIDTH = 14; @@ -111,22 +103,14 @@ private ContainerBalancerAdvisor.AdvisorRequest buildRequest(List request.setIncludeNodes(parseNodeSet(value))); - excludeNodes.ifPresent(value -> request.setExcludeNodes(parseNodeSet(value))); + includeNodes.ifPresent(value -> request.setIncludeNodes(ContainerBalancerCliHelper.parseNodeSet(value))); + excludeNodes.ifPresent(value -> request.setExcludeNodes(ContainerBalancerCliHelper.parseNodeSet(value))); if (profileName.isPresent()) { - request.setProfile(parseProfile(profileName.get())); + request.setProfile(ContainerBalancerCliHelper.parseProfile(profileName.get())); } return request; } - private static ContainerBalancerProfile parseProfile(String name) throws IOException { - try { - return ContainerBalancerProfile.valueOf(name.trim().toUpperCase(Locale.ENGLISH)); - } catch (IllegalArgumentException e) { - throw new IOException("Invalid profile: " + name + ". Expected SLOW, MEDIUM, or FAST."); - } - } - private void printRecommendation(ContainerBalancerRecommendation recommendation) { out().printf("RECOMMENDED CONFIGURATION (profile: %s)%n", recommendation.getProfile().name()); out().println(); @@ -139,7 +123,7 @@ private void printRecommendation(ContainerBalancerRecommendation recommendation) printRecommendedParameters(recommendation); out().println(); out().println(" Estimation:"); - printEstimation(recommendation.getEstimation()); + ContainerBalancerCliHelper.printEstimation(out(), recommendation.getEstimation()); } private void printRecommendedParameters(ContainerBalancerRecommendation recommendation) { @@ -186,43 +170,4 @@ private void printParameterRow(String parameter, String value, String rationaleT parameter, value, rationale); } - private void printEstimation(ContainerBalancerEstimation estimation) { - long estimatedIterations = estimation.getEstimatedIterations(); - long planningIterations = (long) Math.ceil(estimatedIterations * PLANNING_ITERATION_BUFFER); - long cycleTimeMillis = estimation.getMoveTimeoutMillis() + estimation.getBalancingIntervalMillis(); - long baseDurationMillis = estimation.getEstimatedDurationMillis(); - long planningDurationMillis = planningIterations * cycleTimeMillis; - out().printf(" Bytes to move: %s%n", byteDesc(estimation.getBytesToMove())); - out().printf(" Per iteration (estimate): ~%s%n", byteDesc(estimation.getPerIterationBytes())); - out().printf(" Estimated iterations: %d (planning estimate: %d, includes +30%% buffer)%n", - estimatedIterations, planningIterations); - out().printf(" Estimated duration: upper bound %s (planning estimate: %s, includes +30%% buffer)%n", - formatEstimatedDuration(baseDurationMillis), - formatEstimatedDuration(planningDurationMillis)); - out().println(" (assumes full move timeout + interval each cycle)"); - out().println(); - } - - private static String formatEstimatedDuration(long durationMillis) { - double days = durationMillis / 86400000d; - if (days >= 1) { - return String.format(Locale.ENGLISH, "~%.1f days", days); - } - double hours = durationMillis / 3600000d; - if (hours >= 1) { - return String.format(Locale.ENGLISH, "~%.1f hours", hours); - } - long minutes = durationMillis / 60000; - return String.format(Locale.ENGLISH, "~%d min", minutes); - } - - private static Set parseNodeSet(String nodes) { - if (StringUtils.isBlank(nodes)) { - return Collections.emptySet(); - } - return Arrays.stream(nodes.split(",")) - .map(String::trim) - .filter(s -> !s.isEmpty()) - .collect(Collectors.toSet()); - } } diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java index a4d410810eee..342a3d0aceef 100644 --- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java +++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java @@ -40,11 +40,6 @@ import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos; import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoResponseProto; -import org.apache.hadoop.hdds.scm.cli.ContainerBalancerEstimateSubcommand; -import org.apache.hadoop.hdds.scm.cli.ContainerBalancerRecommendSubcommand; -import org.apache.hadoop.hdds.scm.cli.ContainerBalancerStartSubcommand; -import org.apache.hadoop.hdds.scm.cli.ContainerBalancerStatusSubcommand; -import org.apache.hadoop.hdds.scm.cli.ContainerBalancerStopSubcommand; import org.apache.hadoop.hdds.scm.client.ScmClient; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerConfiguration; import org.apache.hadoop.hdds.utils.IOUtils; From 2601f356542b79932085213999013566cebff7ec Mon Sep 17 00:00:00 2001 From: sravani-revuri Date: Wed, 30 Sep 2026 12:11:30 +0530 Subject: [PATCH 4/5] no balancing mssg + start cmmd + remove exception --- .../balancer/ContainerBalancerAdvisor.java | 56 +++++++++++-------- .../ContainerBalancerRecommendation.java | 12 ++++ .../TestContainerBalancerAdvisor.java | 15 +++-- .../scm/cli/ContainerBalancerCliHelper.java | 54 ++++++++++++++++-- .../ContainerBalancerEstimateSubcommand.java | 5 +- .../ContainerBalancerRecommendSubcommand.java | 25 ++++++--- .../cli/TestContainerBalancerSubCommand.java | 30 +++++++++- 7 files changed, 152 insertions(+), 45 deletions(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java index f54038a73f2a..fb8267893e48 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java @@ -25,7 +25,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; -import java.util.Locale; import java.util.Map; import java.util.Objects; import java.util.Set; @@ -79,7 +78,7 @@ public static List estimate(OzoneConfiguration conf ContainerBalancerClusterSnapshot snapshot = ContainerBalancerClusterAnalyzer.analyze(nodes, thresholdRatio, includeNodes, excludeNodes); - validateSnapshotForEstimation(snapshot, conf); + validateSnapshotForEstimation(snapshot); List profiles = selectProfiles(request); List estimations = new ArrayList<>(profiles.size()); @@ -118,9 +117,12 @@ public static List recommend( ContainerBalancerClusterSnapshot snapshot = ContainerBalancerClusterAnalyzer.analyze( nodes, thresholdRatio, includeNodes, excludeNodes); - validateSnapshotForEstimation(snapshot, conf); List profiles = selectProfilesForRecommend(request); + if (isClusterAlreadyBalanced(snapshot)) { + return recommendationsForBalancedCluster(profiles, thresholdPercent); + } + validateSnapshotForEstimation(snapshot); List recommendations = new ArrayList<>(profiles.size()); for (ContainerBalancerProfile profile : profiles) { recommendations.add(recommendForProfile( @@ -129,6 +131,24 @@ public static List recommend( return Collections.unmodifiableList(recommendations); } + private static boolean isClusterAlreadyBalanced(ContainerBalancerClusterSnapshot snapshot) { + return snapshot.getBytesToMove() <= 0 + || (snapshot.getSourceCount() < 1 && snapshot.getTargetCount() < 1); + } + + private static List recommendationsForBalancedCluster( + List profiles, double thresholdPercent) { + List 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, @@ -227,31 +247,31 @@ private static Map buildRationale( int profileDatanodesPercentage = profile.getDatanodesMaxPercentage(); int usedDatanodesPercentage = estimation.getMaxDatanodesPercentage(); if (usedDatanodesPercentage == profileDatanodesPercentage) { - rationale.put("maxDatanodesPercentage", String.format(Locale.ENGLISH, + rationale.put("maxDatanodesPercentage", String.format( "profile default (%d%%)", usedDatanodesPercentage)); } else { - rationale.put("maxDatanodesPercentage", String.format(Locale.ENGLISH, + rationale.put("maxDatanodesPercentage", String.format( "profile %d%%, raised to %d%% (≥2 nodes)", profileDatanodesPercentage, usedDatanodesPercentage)); } - rationale.put("maxSizeToMovePerIteration", String.format(Locale.ENGLISH, + 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(Locale.ENGLISH, + rationale.put("maxSizeEnteringTarget", String.format( "profile default (%d GB)", estimation.getMaxSizeEnteringTarget() / OzoneConsts.GB)); - rationale.put("maxSizeLeavingSource", String.format(Locale.ENGLISH, + rationale.put("maxSizeLeavingSource", String.format( "profile default (%d GB)", estimation.getMaxSizeLeavingSource() / OzoneConsts.GB)); - rationale.put("moveTimeout", String.format(Locale.ENGLISH, + rationale.put("moveTimeout", String.format( "configuration default (%d min)", moveTimeoutMillis / 60000)); - rationale.put("moveReplicationTimeout", String.format(Locale.ENGLISH, + rationale.put("moveReplicationTimeout", String.format( "configuration default (%d min)", moveReplicationTimeoutMillis / 60000)); - rationale.put("balancingInterval", String.format(Locale.ENGLISH, + 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(Locale.ENGLISH, + rationale.put("iterations", String.format( "includes +30%% buffer (planning estimate: %d)", planningIterations)); return rationale; } @@ -430,12 +450,7 @@ private static long maxPositive(long... values) { return result; } - private static void validateSnapshotForEstimation(ContainerBalancerClusterSnapshot snapshot, - OzoneConfiguration conf) { - long containerSizeBytes = (long) conf.getStorageSize( - ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE, - ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT, - StorageUnit.BYTES); + private static void validateSnapshotForEstimation(ContainerBalancerClusterSnapshot snapshot ) { if (snapshot.getSourceCount() < 1) { throw new IllegalArgumentException("No over-utilized datanodes (sources) found."); } @@ -445,11 +460,6 @@ private static void validateSnapshotForEstimation(ContainerBalancerClusterSnapsh if (snapshot.getBytesToMove() <= 0) { throw new IllegalArgumentException("No bytes to move."); } - if (snapshot.getBytesToMove() < containerSizeBytes) { - throw new IllegalArgumentException( - "Bytes to move (" + snapshot.getBytesToMove() - + ") is less than container size (" + containerSizeBytes + ")."); - } if (snapshot.getTotalEligibleDatanodes() < 2) { throw new IllegalArgumentException(String.format( "Container Balancer found %d eligible datanode(s) but requires at least 2.", diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java index eed12de14396..bcd3197a57e0 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerRecommendation.java @@ -39,6 +39,7 @@ public final class ContainerBalancerRecommendation { private final int recommendedIterations; private final Map rationale; private final ContainerBalancerEstimation estimation; + private final boolean clusterBalanced; private ContainerBalancerRecommendation(Builder b) { this.profile = Objects.requireNonNull(b.profile, "profile == null"); @@ -56,6 +57,7 @@ private ContainerBalancerRecommendation(Builder b) { ? Collections.emptyMap() : Collections.unmodifiableMap(b.rationale); this.estimation = b.estimation; + this.clusterBalanced = b.clusterBalanced; } public static Builder newBuilder() { @@ -70,6 +72,10 @@ public boolean succeeded() { return failureMessage == null; } + public boolean isClusterBalanced() { + return clusterBalanced; + } + public String getFailureMessage() { return failureMessage; } @@ -133,6 +139,7 @@ public static final class Builder { private int recommendedIterations; private Map rationale; private ContainerBalancerEstimation estimation; + private boolean clusterBalanced; private Builder() { } @@ -202,6 +209,11 @@ public Builder setEstimation(ContainerBalancerEstimation estimationValue) { return this; } + public Builder setClusterBalanced(boolean balanced) { + this.clusterBalanced = balanced; + return this; + } + public ContainerBalancerRecommendation build() { return new ContainerBalancerRecommendation(this); } diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java index dcf55f7a9c29..46cb5933f955 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerAdvisor.java @@ -357,16 +357,21 @@ void testRecommendRespectsThresholdOverride() { } @Test - void testRecommendFailsWhenClusterBalanced() { + void testRecommendSucceedsWhenClusterBalanced() { OzoneConfiguration conf = new OzoneConfiguration(); List balanced = new ArrayList<>(); balanced.add(proto("dn-1", OzoneConsts.TB, (long) (0.70 * OzoneConsts.TB))); balanced.add(proto("dn-2", OzoneConsts.TB, (long) (0.70 * OzoneConsts.TB))); - assertThrows(IllegalArgumentException.class, () -> - ContainerBalancerAdvisor.recommend( - conf, - new ContainerBalancerAdvisor.AdvisorRequest().setNodes(balanced))); + List results = ContainerBalancerAdvisor.recommend( + conf, + new ContainerBalancerAdvisor.AdvisorRequest().setNodes(balanced)); + + assertEquals(3, results.size()); + for (ContainerBalancerRecommendation result : results) { + assertThat(result.succeeded()).isTrue(); + assertThat(result.isClusterBalanced()).isTrue(); + } } @Test diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java index 868940d7d675..3cb8961959a9 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerCliHelper.java @@ -23,12 +23,14 @@ import java.io.PrintWriter; import java.util.Arrays; import java.util.Collections; -import java.util.Locale; +import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerEstimation; import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerProfile; +import org.apache.hadoop.hdds.scm.container.balancer.ContainerBalancerRecommendation; +import org.apache.hadoop.ozone.OzoneConsts; /** Shared parsing and CLI output for container balancer subcommands. */ public final class ContainerBalancerCliHelper { @@ -41,7 +43,7 @@ private ContainerBalancerCliHelper() { public static ContainerBalancerProfile parseProfile(String name) throws IOException { try { - return ContainerBalancerProfile.valueOf(name.trim().toUpperCase(Locale.ENGLISH)); + return ContainerBalancerProfile.valueOf(name.trim().toUpperCase()); } catch (IllegalArgumentException e) { throw new IOException("Invalid profile: " + name + ". Expected SLOW, MEDIUM, or FAST."); } @@ -77,13 +79,55 @@ public static void printEstimation(PrintWriter out, ContainerBalancerEstimation public static String formatEstimatedDuration(long durationMillis) { double days = durationMillis / 86400000d; if (days >= 1) { - return String.format(Locale.ENGLISH, "~%.1f days", days); + return String.format("~%.1f days", days); } double hours = durationMillis / 3600000d; if (hours >= 1) { - return String.format(Locale.ENGLISH, "~%.1f hours", hours); + return String.format("~%.1f hours", hours); } long minutes = durationMillis / 60000; - return String.format(Locale.ENGLISH, "~%d min", minutes); + return String.format("~%d min", minutes); + } + + public static String formatStartCommand( + ContainerBalancerRecommendation recommendation, + Optional includeNodes, + Optional excludeNodes) { + StringBuilder sb = new StringBuilder("ozone admin containerbalancer start"); + sb.append(" -t ").append(formatThreshold(recommendation.getThresholdPercent())); + sb.append(" -i ").append(recommendation.getRecommendedIterations()); + sb.append(" -d ").append(recommendation.getMaxDatanodesPercentage()); + sb.append(" -s ").append(bytesToGb(recommendation.getMaxSizeToMovePerIteration())); + sb.append(" -e ").append(bytesToGb(recommendation.getMaxSizeEnteringTarget())); + sb.append(" -l ").append(bytesToGb(recommendation.getMaxSizeLeavingSource())); + sb.append(" --balancing-iteration-interval-minutes ") + .append(recommendation.getBalancingIntervalMillis() / 60000); + sb.append(" --move-timeout-minutes ") + .append(recommendation.getMoveTimeoutMillis() / 60000); + sb.append(" --move-replication-timeout-minutes ") + .append(recommendation.getMoveReplicationTimeoutMillis() / 60000); + appendNodeFilters(sb, includeNodes, excludeNodes); + return sb.toString(); + } + + private static void appendNodeFilters( + StringBuilder sb, + Optional includeNodes, + Optional excludeNodes) { + includeNodes.filter(s -> !StringUtils.isBlank(s)) + .ifPresent(s -> sb.append(" --include-datanodes ").append(s.trim())); + excludeNodes.filter(s -> !StringUtils.isBlank(s)) + .ifPresent(s -> sb.append(" --exclude-datanodes ").append(s.trim())); + } + + private static long bytesToGb(long bytes) { + return bytes / OzoneConsts.GB; + } + + private static String formatThreshold(double thresholdPercent) { + if (thresholdPercent == Math.rint(thresholdPercent)) { + return String.valueOf((long) thresholdPercent); + } + return String.valueOf(thresholdPercent); } } diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java index a0f71791a2a7..8278080ce2be 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerEstimateSubcommand.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.util.List; -import java.util.Locale; import java.util.Optional; import org.apache.hadoop.hdds.cli.HddsVersionProvider; import org.apache.hadoop.hdds.conf.OzoneConfiguration; @@ -111,8 +110,8 @@ private void printBasedOn(ContainerBalancerEstimation estimation) { long moveTimeoutMinutes = Math.round(estimation.getMoveTimeoutMillis() / 60000d); long balancingIntervalMinutes = Math.round(estimation.getBalancingIntervalMillis() / 60000d); out().println(" Based on:"); - out().printf(Locale.ENGLISH, " Threshold: %.1f%%%n", estimation.getThresholdPercent()); - out().printf(Locale.ENGLISH, " Datanode involvement: %d%%%n", + out().printf(" Threshold: %.1f%%%n", estimation.getThresholdPercent()); + out().printf(" Datanode involvement: %d%%%n", estimation.getMaxDatanodesPercentage()); out().printf(" Max entering target: %s / node%n", byteDesc(estimation.getMaxSizeEnteringTarget())); out().printf(" Max leaving source: %s / node%n", byteDesc(estimation.getMaxSizeLeavingSource())); diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java index 1601fa49943e..46ce3da80eef 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerRecommendSubcommand.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.util.List; -import java.util.Locale; import java.util.Map; import java.util.Optional; import org.apache.hadoop.hdds.cli.HddsVersionProvider; @@ -86,6 +85,13 @@ public void execute(ScmClient scmClient) throws IOException { throw new IOException(e.getMessage(), e); } + if (!recommendations.isEmpty() + && recommendations.stream().allMatch(ContainerBalancerRecommendation::isClusterBalanced)) { + out().println("Cluster is already balanced within the configured threshold. " + + "No balancing recommended."); + return; + } + boolean anySucceeded = false; for (ContainerBalancerRecommendation recommendation : recommendations) { printRecommendation(recommendation); @@ -124,6 +130,11 @@ private void printRecommendation(ContainerBalancerRecommendation recommendation) out().println(); out().println(" Estimation:"); ContainerBalancerCliHelper.printEstimation(out(), recommendation.getEstimation()); + + out().println(" Suggested commands:"); + out().println(" " + ContainerBalancerCliHelper.formatStartCommand( + recommendation, includeNodes, excludeNodes)); + out().println(); } private void printRecommendedParameters(ContainerBalancerRecommendation recommendation) { @@ -136,10 +147,10 @@ private void printRecommendedParameters(ContainerBalancerRecommendation recommen out().println(" Recommended parameters:"); printParameterRow("Parameter", "Value", "Rationale"); printParameterRow("--threshold", - String.format(Locale.ENGLISH, "%.1f%%", recommendation.getThresholdPercent()), + String.format("%.1f%%", recommendation.getThresholdPercent()), rationale.get("threshold")); printParameterRow("--max-datanodes-percentage-to-involve", - String.format(Locale.ENGLISH, "%d%%", recommendation.getMaxDatanodesPercentage()), + String.format("%d%%", recommendation.getMaxDatanodesPercentage()), rationale.get("maxDatanodesPercentage")); printParameterRow("--max-size-to-move-per-iteration-in-gb", byteDesc(recommendation.getMaxSizeToMovePerIteration()), @@ -151,13 +162,13 @@ private void printRecommendedParameters(ContainerBalancerRecommendation recommen byteDesc(recommendation.getMaxSizeLeavingSource()) + " / node", rationale.get("maxSizeLeavingSource")); printParameterRow("--move-timeout-minutes", - String.format(Locale.ENGLISH, "%d min", moveTimeoutMinutes), + String.format("%d min", moveTimeoutMinutes), rationale.get("moveTimeout")); printParameterRow("--move-replication-timeout-minutes", - String.format(Locale.ENGLISH, "%d min", moveReplicationTimeoutMinutes), + String.format("%d min", moveReplicationTimeoutMinutes), rationale.get("moveReplicationTimeout")); printParameterRow("--balancing-iteration-interval-minutes", - String.format(Locale.ENGLISH, "%d min", balancingIntervalMinutes), + String.format("%d min", balancingIntervalMinutes), rationale.get("balancingInterval")); printParameterRow("--iterations", String.valueOf(recommendation.getRecommendedIterations()), @@ -166,7 +177,7 @@ private void printRecommendedParameters(ContainerBalancerRecommendation recommen private void printParameterRow(String parameter, String value, String rationaleText) { String rationale = rationaleText == null ? "" : rationaleText; - out().printf(Locale.ENGLISH, " %-" + PARAM_COLUMN_WIDTH + "s %-" + VALUE_COLUMN_WIDTH + "s %s%n", + out().printf(" %-" + PARAM_COLUMN_WIDTH + "s %-" + VALUE_COLUMN_WIDTH + "s %s%n", parameter, value, rationale); } diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java index 342a3d0aceef..8a7a24d6df26 100644 --- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java +++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/TestContainerBalancerSubCommand.java @@ -1050,7 +1050,10 @@ void testContainerBalancerRecommendSubcommandDefaultShowsAllProfiles() throws IO .contains("Bytes to move:") .contains("planning estimate:") .contains("assumes full move timeout + interval each cycle") + .contains("Suggested commands:") + .contains("ozone admin containerbalancer start") .doesNotContain("Recommendation failed:"); + assertThat(output.split("Suggested commands:")).hasSize(4); } @Test @@ -1062,7 +1065,9 @@ void testContainerBalancerRecommendSubcommandWithThresholdOverride() throws IOEx parseSubcommand(recommendCmd, "-t", "5"); recommendCmd.execute(scmClient); - assertThat(out.get()).contains("RECOMMENDED CONFIGURATION (profile: SLOW)"); + assertThat(out.get()) + .contains("RECOMMENDED CONFIGURATION (profile: SLOW)") + .contains("ozone admin containerbalancer start -t 5"); } @Test @@ -1078,7 +1083,10 @@ void testContainerBalancerRecommendSubcommandWithProfileShowsOneProfile() throws assertThat(output) .contains("RECOMMENDED CONFIGURATION (profile: MEDIUM)") .doesNotContain("RECOMMENDED CONFIGURATION (profile: SLOW)") - .doesNotContain("RECOMMENDED CONFIGURATION (profile: FAST)"); + .doesNotContain("RECOMMENDED CONFIGURATION (profile: FAST)") + .contains("Suggested commands:") + .contains("ozone admin containerbalancer start"); + assertThat(output.split("Suggested commands:")).hasSize(2); } @Test @@ -1092,6 +1100,24 @@ void testContainerBalancerRecommendSubcommandInvalidThresholdFails() throws IOEx assertThat(ex.getMessage()).contains("Threshold should be specified in the range [0.0, 100.0)."); } + @Test + void testContainerBalancerRecommendSubcommandWhenClusterBalanced() throws IOException { + ScmClient scmClient = mock(ScmClient.class); + List balanced = new ArrayList<>(); + long capacity = OzoneConsts.TB; + balanced.add(datanodeUsageProto("dn-1", capacity, (long) (capacity * 0.70))); + balanced.add(datanodeUsageProto("dn-2", capacity, (long) (capacity * 0.70))); + when(scmClient.getDatanodeUsageInfo(true, Integer.MAX_VALUE)).thenReturn(balanced); + + parseSubcommand(recommendCmd); + recommendCmd.execute(scmClient); + + assertThat(out.get()) + .contains("Cluster is already balanced within the configured threshold.") + .doesNotContain("Suggested commands:") + .doesNotContain("ozone admin containerbalancer start"); + } + /** * Imbalanced cluster for estimate CLI tests. * From 272af8e21c6053f7863973453be0ef1e2534431a Mon Sep 17 00:00:00 2001 From: sravani-revuri Date: Wed, 7 Oct 2026 18:44:56 +0530 Subject: [PATCH 5/5] removed advisor request , reusing estimate thing --- .../balancer/ContainerBalancerAdvisor.java | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java index fb8267893e48..3a8a19062137 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerAdvisor.java @@ -126,7 +126,7 @@ public static List recommend( List recommendations = new ArrayList<>(profiles.size()); for (ContainerBalancerProfile profile : profiles) { recommendations.add(recommendForProfile( - conf, request, profile, balancerConfig, thresholdPercent, includeNodes, excludeNodes)); + conf, request, profile, snapshot, balancerConfig, thresholdPercent)); } return Collections.unmodifiableList(recommendations); } @@ -153,10 +153,9 @@ private static ContainerBalancerRecommendation recommendForProfile( OzoneConfiguration conf, AdvisorRequest request, ContainerBalancerProfile profile, + ContainerBalancerClusterSnapshot snapshot, ContainerBalancerConfiguration balancerConfig, - double thresholdPercent, - Set includeNodes, - Set excludeNodes) { + double thresholdPercent) { ContainerBalancerRecommendation.Builder builder = ContainerBalancerRecommendation.newBuilder() .setProfile(profile) @@ -169,14 +168,8 @@ private static ContainerBalancerRecommendation recommendForProfile( validateMoveTimeouts(conf, moveReplicationTimeoutMillis, moveTimeoutMillis); validateBalancingIntervalMillis(balancingIntervalMillis); - AdvisorRequest estimateRequest = new AdvisorRequest() - .setNodes(request.nodes) - .setThresholdPercent(thresholdPercent) - .setIncludeNodes(includeNodes) - .setExcludeNodes(excludeNodes) - .setProfile(profile); - - ContainerBalancerEstimation estimation = estimate(conf, estimateRequest).get(0); + ContainerBalancerEstimation estimation = estimateForProfile( + conf, request, profile, snapshot, balancerConfig, thresholdPercent); if (!estimation.succeeded()) { throw new IllegalArgumentException(estimation.getFailureMessage()); }