diff --git a/.mvn/extensions.xml b/.mvn/extensions.xml index 36f013dd928d..781053b29fca 100644 --- a/.mvn/extensions.xml +++ b/.mvn/extensions.xml @@ -24,7 +24,7 @@ com.gradle develocity-maven-extension - 2.5.0 + 2.6.0 com.gradle diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerInfo.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerInfo.java index 46c21bc4d65c..a9ce007a5401 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerInfo.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerInfo.java @@ -21,11 +21,12 @@ import java.util.List; import java.util.Objects; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerRunningStatus; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.VolumeReportProto; /** * DiskBalancer's information to persist and for report. - * Report-only fields (idealUsage, volumeInfo) are NOT persisted to YAML. + * Report-only fields (idealUsage, volumeInfo, storageTypeInfo) are NOT persisted to YAML. */ public class DiskBalancerInfo { private DiskBalancerRunningStatus operationalState; @@ -44,6 +45,8 @@ public class DiskBalancerInfo { private double idealUsage; // Report-only: per-volume info. NOT persisted. private List volumeInfo; + // Report-only: per-storage-type balancing info. NOT persisted. + private List storageTypeInfo; public DiskBalancerInfo(DiskBalancerRunningStatus operationalState, double threshold, long bandwidthInMB, int parallelThread, boolean stopAfterDiskEven) { @@ -250,6 +253,14 @@ public void setVolumeInfo(List volumeInfo) { this.volumeInfo = volumeInfo; } + public List getStorageTypeInfo() { + return storageTypeInfo != null ? storageTypeInfo : Collections.emptyList(); + } + + public void setStorageTypeInfo(List storageTypeInfo) { + this.storageTypeInfo = storageTypeInfo; + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java index 770f7b67414b..801c0f952dea 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java @@ -85,6 +85,7 @@ public DatanodeDiskBalancerInfoProto getDiskBalancerInfo(GetDiskBalancerInfoRequ .setRunningStatus(info.getOperationalState()) .setIdealUsage(info.getIdealUsage()) .addAllVolumeInfo(info.getVolumeInfo()) + .addAllStorageTypeInfo(info.getStorageTypeInfo()) .build(); } @@ -167,4 +168,3 @@ public void close() { } } - diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerService.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerService.java index 7ab0243db064..1105947d5e82 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerService.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerService.java @@ -21,9 +21,10 @@ import static org.apache.hadoop.ozone.container.common.volume.StorageVolume.TMP_DIR_NAME; import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.calculateVolumeDataDensity; import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.getIdealUsage; +import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.getUsableVolumesByStorageType; import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.getVolumeUsages; +import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.isBalanceable; -import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Strings; import java.io.File; import java.io.IOException; @@ -47,11 +48,14 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Collectors; import org.apache.commons.io.FileUtils; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.fs.SpaceUsageSource; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerRunningStatus; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.VolumeReportProto; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; @@ -71,6 +75,7 @@ import org.apache.hadoop.ozone.container.common.utils.StorageVolumeUtil; import org.apache.hadoop.ozone.container.common.volume.HddsVolume; import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.ThresholdRange; import org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.VolumeFixedUsage; import org.apache.hadoop.ozone.container.diskbalancer.policy.ContainerCandidate; import org.apache.hadoop.ozone.container.diskbalancer.policy.ContainerChoosingPolicy; @@ -458,8 +463,11 @@ public BackgroundTaskQueue getTasks() { if (queue.isEmpty() && inProgressContainers.isEmpty()) { if (stopAfterDiskEven) { - LOG.info("Disk balancer is stopped due to disk even as" + - " the property StopAfterDiskEven is set to true."); + // Containers only move between volumes of the same StorageType, so no volume pair is also + // reported when every StorageType has fewer than two volumes, even if those volumes are + // unevenly used. + LOG.info("Disk balancer is stopped because no volume pair of the same StorageType needs" + + " balancing as the property StopAfterDiskEven is set to true."); this.operationalState = DiskBalancerRunningStatus.STOPPED; try { // Persist the updated shouldRun status into the YAML file @@ -691,7 +699,6 @@ private void deleteContainer(Container container) { } } - @VisibleForTesting public void cleanupPendingDeletionContainers() { // delete all pending deletion containers before stop the service boolean ret; @@ -717,16 +724,15 @@ private boolean tryCleanupOnePendingDeletionContainer() { public DiskBalancerInfo getDiskBalancerInfo() { final List volumeUsages = getVolumeUsages(volumeSet, deltaSizes); - - // Calculate volumeDataDensity - final double volumeDataDensity = calculateVolumeDataDensity(volumeUsages); - - long bytesToMove = 0; - if (this.operationalState == DiskBalancerRunningStatus.RUNNING) { - // this calculates live changes in bytesToMove - // calculate bytes to move if the balancer is in a running state, else 0. - bytesToMove = calculateBytesToMove(volumeUsages); - } + final List storageTypeInfo = + buildStorageTypeInfo(volumeUsages, threshold, + this.operationalState == DiskBalancerRunningStatus.RUNNING); + final double volumeDataDensity = storageTypeInfo.stream() + .mapToDouble(StorageTypeDiskBalancerInfoProto::getCurrentVolumeDensitySum) + .sum(); + final long bytesToMove = storageTypeInfo.stream() + .mapToLong(StorageTypeDiskBalancerInfoProto::getBytesToMove) + .sum(); DiskBalancerInfo info = new DiskBalancerInfo(operationalState, threshold, bandwidthInMB, parallelThread, stopAfterDiskEven, version, formatContainerStates(movableContainerStates), @@ -734,9 +740,43 @@ parallelThread, stopAfterDiskEven, version, formatContainerStates(movableContain metrics.getFailureCount(), bytesToMove, metrics.getSuccessBytes(), volumeDataDensity); info.setIdealUsage(getIdealUsage(volumeUsages)); info.setVolumeInfo(buildVolumeReportProto(volumeUsages)); + info.setStorageTypeInfo(storageTypeInfo); return info; } + static List buildStorageTypeInfo( + List volumeUsages, double thresholdPercentage, + boolean calculateBytes) { + final Map> usagesByStorageType = + getUsableVolumesByStorageType(volumeUsages); + final List result = new ArrayList<>(); + for (StorageType storageType : StorageType.values()) { + final List sameTypeUsages = usagesByStorageType.get(storageType); + if (sameTypeUsages == null) { + continue; + } + + final boolean balanceable = isBalanceable(sameTypeUsages); + final StorageTypeDiskBalancerInfoProto.Builder builder = + StorageTypeDiskBalancerInfoProto.newBuilder() + .setStorageType(StorageTypeUtils.getStorageTypeProto(storageType)) + .setUsableVolumeCount(sameTypeUsages.size()) + .setBalanceable(balanceable); + if (balanceable) { + // One snapshot per storage type, shared by the reported ideal usage and the byte estimate. + final ThresholdRange thresholdRange = + ThresholdRange.of(sameTypeUsages, thresholdPercentage); + builder.setIdealUsage(thresholdRange.getIdealUsage()) + .setCurrentVolumeDensitySum(calculateVolumeDataDensity(sameTypeUsages)); + if (calculateBytes) { + builder.setBytesToMove(calculateBytesToMove(sameTypeUsages, thresholdRange)); + } + } + result.add(builder.build()); + } + return result; + } + /** * Build a list of VolumeReportProto from a list of VolumeFixedUsage. * VolumeReportProto consists of information like StorageID, @@ -755,6 +795,7 @@ public static List buildVolumeReportProto(List buildVolumeReportProto(List inputVolumeSet) { - // If there are no available volumes or only one volume, return 0 bytes to move - if (inputVolumeSet.isEmpty() || inputVolumeSet.size() < 2) { + if (!isBalanceable(inputVolumeSet)) { return 0; } + return calculateBytesToMove(inputVolumeSet, ThresholdRange.of(inputVolumeSet, threshold)); + } - // Calculate actual threshold - final double actualThreshold = getIdealUsage(inputVolumeSet) + threshold / 100.0; + /** + * @param thresholdRange range already computed for {@code inputVolumeSet} + */ + private static long calculateBytesToMove(List inputVolumeSet, + ThresholdRange thresholdRange) { + // Volumes above the upper threshold are the sources that need to give up data. + final double actualThreshold = thresholdRange.getUpperThreshold(); long totalBytesToMove = 0; @@ -806,7 +853,6 @@ public DiskBalancerServiceMetrics getMetrics() { return metrics; } - @VisibleForTesting public void setBalancedBytesInLastWindow(long bytes) { this.balancedBytesInLastWindow.set(bytes); } @@ -815,17 +861,14 @@ public ContainerChoosingPolicy getVolumeContainerChoosingPolicy() { return volumeContainerChoosingPolicy; } - @VisibleForTesting public void setVolumeContainerChoosingPolicy(ContainerChoosingPolicy volumeContainerChoosingPolicy) { this.volumeContainerChoosingPolicy = volumeContainerChoosingPolicy; } - @VisibleForTesting public Set getInProgressContainers() { return inProgressContainers; } - @VisibleForTesting public Map getDeltaSizes() { return deltaSizes; } @@ -878,7 +921,6 @@ public void shutdown() { } } - @VisibleForTesting public static void setInjector(FaultInjector instance) { injector = instance; } @@ -894,17 +936,14 @@ private static void pauseInjector() { } } - @VisibleForTesting public void setReplicaDeletionDelay(long durationMills) { this.replicaDeletionDelay = durationMills; } - @VisibleForTesting public int getPendingDeletionDeadlineCount() { return pendingDeletionContainers.size(); } - @VisibleForTesting public int getPendingDeletionQueueSize() { return pendingDeletionContainers.values().stream() .mapToInt(Queue::size) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerVolumeCalculation.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerVolumeCalculation.java index 071e43c51877..7899e40f3f3f 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerVolumeCalculation.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerVolumeCalculation.java @@ -20,9 +20,11 @@ import static org.apache.ratis.util.Preconditions.assertInstanceOf; import static org.apache.ratis.util.Preconditions.assertTrue; +import java.util.EnumMap; import java.util.List; import java.util.Map; import java.util.stream.Collectors; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.fs.SpaceUsageSource; import org.apache.hadoop.ozone.container.common.volume.HddsVolume; import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; @@ -59,7 +61,90 @@ public static List getVolumeUsages(MutableVolumeSet volumeSet, .map(v -> newVolumeFixedUsage(v, deltas)) .collect(Collectors.toList()); } - + + /** + * Group positive-capacity volumes by storage type. These are the volumes eligible for disk + * balancing calculations. + */ + public static Map> getUsableVolumesByStorageType( + List volumes) { + return volumes.stream() + .filter(v -> v.getUsage().getCapacity() > 0) + .collect(Collectors.groupingBy(v -> v.getVolume().getStorageType(), + () -> new EnumMap<>(StorageType.class), Collectors.toList())); + } + + /** + * Whether a set of volumes can be balanced against each other. A single volume has no second + * volume to move containers to, so there is nothing to balance. + */ + public static boolean isBalanceable(List volumes) { + return volumes != null && volumes.size() >= 2; + } + + /** + * The ideal usage of a set of volumes and the acceptable band around it, computed once from a + * single volume snapshot. Volumes outside the band are candidates for balancing. + */ + public static final class ThresholdRange { + + private final double idealUsage; + private final double lowerThreshold; + private final double upperThreshold; + + private ThresholdRange(double idealUsage, double thresholdPercentage) { + final double threshold = thresholdPercentage / 100.0; + this.idealUsage = idealUsage; + this.lowerThreshold = idealUsage - threshold; + this.upperThreshold = idealUsage + threshold; + } + + /** + * @param volumes volumes that share a storage type, so the ideal usage is a target that + * container moves among them can actually reach + * @param thresholdPercentage acceptable deviation from ideal usage, in percent + */ + public static ThresholdRange of(List volumes, double thresholdPercentage) { + return new ThresholdRange( + DiskBalancerVolumeCalculation.getIdealUsage(volumes), thresholdPercentage); + } + + public double getIdealUsage() { + return idealUsage; + } + + public double getLowerThreshold() { + return lowerThreshold; + } + + public double getUpperThreshold() { + return upperThreshold; + } + + /** + * How far the most deviant volume lies outside the band. Zero or negative means every volume + * is within the band and nothing needs to move. + * + * @param sortedVolumes volumes sorted ascending by utilization + */ + public double getViolation(List sortedVolumes) { + final double lowestUsage = sortedVolumes.get(0).getUtilization(); + final double highestUsage = sortedVolumes.get(sortedVolumes.size() - 1).getUtilization(); + return Math.max(highestUsage - upperThreshold, lowerThreshold - lowestUsage); + } + + /** + * Whether both ends of the utilization range sit inside the band. + * + * @param sortedVolumes volumes sorted ascending by utilization + */ + public boolean isWithinRange(List sortedVolumes) { + final double lowestUsage = sortedVolumes.get(0).getUtilization(); + final double highestUsage = sortedVolumes.get(sortedVolumes.size() - 1).getUtilization(); + return highestUsage < upperThreshold && lowestUsage > lowerThreshold; + } + } + /** * Get ideal usage from an immutable list of volumes. * diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/policy/DefaultContainerChoosingPolicy.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/policy/DefaultContainerChoosingPolicy.java index baee04bb9768..aeda1a2c0809 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/policy/DefaultContainerChoosingPolicy.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/policy/DefaultContainerChoosingPolicy.java @@ -19,7 +19,8 @@ import static java.util.concurrent.TimeUnit.HOURS; import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.computeUtilization; -import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.getIdealUsage; +import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.getUsableVolumesByStorageType; +import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.isBalanceable; import static org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.newVolumeFixedUsage; import com.google.common.cache.Cache; @@ -32,6 +33,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.fs.SpaceUsageSource; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State; import org.apache.hadoop.hdds.scm.container.ContainerID; @@ -40,16 +42,23 @@ import org.apache.hadoop.ozone.container.common.volume.HddsVolume; import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.ThresholdRange; import org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerVolumeCalculation.VolumeFixedUsage; import org.apache.hadoop.ozone.container.ozoneimpl.OzoneContainer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** - * First chooses a source volume and destination volume pair based on ideal utilization and threshold, - * then chooses a container from the source volume that can be moved to the destination without - * exceeding the upper threshold. Space is reserved on the destination only when a container is - * chosen, using the actual container size. + * Volumes are grouped by {@link StorageType} and each group is balanced on its own, so a container + * never moves between volumes of different storage types. Within a group, this first chooses a + * source volume and destination volume pair based on ideal utilization and threshold, then chooses + * a container from the source volume that can be moved to the destination without exceeding the + * upper threshold. Space is reserved on the destination only when a container is chosen, using the + * actual container size. + * + * A storage type with fewer than two usable volumes is skipped rather than paired with a volume of + * another type, because moving a container across types would break the storage policy its data was + * placed for. * * Which container states may move is determined by {@code DiskBalancerConfiguration#getMovableContainerStates()}. */ @@ -91,62 +100,140 @@ public ContainerCandidate chooseVolumesAndContainer(OzoneContainer ozoneContaine return null; } - // Calculate ideal usage and threshold range (once) - final double idealUsage = getIdealUsage(volumeUsages); - final double actualThreshold = thresholdPercentage / 100.0; - final double lowerThreshold = idealUsage - actualThreshold; - final double upperThreshold = idealUsage + actualThreshold; - - if (LOG.isDebugEnabled()) { - logVolumeBalancingState(volumeUsages, idealUsage, thresholdPercentage, - lowerThreshold, upperThreshold, deltaMap); + // Balance each storage type independently so a container never crosses types. Try the type + // with the largest threshold violation first. StorageType order is only a deterministic + // tie-breaker. Each type's threshold range is computed once from the snapshot above and + // reused for both the ordering and the balancing attempt. + final List storageTypeGroups = + getUsableVolumesByStorageType(volumeUsages).entrySet().stream() + .filter(entry -> isGroupBalanceable(entry.getKey(), entry.getValue())) + .map(entry -> new StorageTypeGroup(entry.getKey(), entry.getValue(), + ThresholdRange.of(entry.getValue(), thresholdPercentage))) + .sorted(Comparator.comparingDouble(StorageTypeGroup::getViolation).reversed() + .thenComparingInt(group -> group.getStorageType().ordinal())) + .collect(Collectors.toList()); + for (StorageTypeGroup group : storageTypeGroups) { + final ContainerCandidate candidate = chooseWithinStorageType(ozoneContainer, group, + deltaMap, inProgressContainerIDs, thresholdPercentage, movableContainerStates); + if (candidate != null) { + return candidate; + } } + LOG.debug("Failed to find appropriate destination volume and container in any storage type."); + return null; + } finally { + lock.unlock(); + } + } + + private static boolean isGroupBalanceable(StorageType storageType, + List volumeUsages) { + if (isBalanceable(volumeUsages)) { + return true; + } + LOG.debug("Skipping storage type {} for disk balancing: only {} usable volume(s)", + storageType, volumeUsages.size()); + return false; + } - // Get highest and lowest utilization volumes - final VolumeFixedUsage highestUsage = volumeUsages.get(volumeUsages.size() - 1); - final VolumeFixedUsage lowestUsage = volumeUsages.get(0); + /** + * Volumes of a single storage type with the threshold range they are measured against. The + * range is computed once per balancing cycle so the ordering and the balancing attempt agree. + */ + private static final class StorageTypeGroup { - // Only return null if highest is below upper threshold AND lowest is above lower threshold - if (highestUsage.getUtilization() < upperThreshold && - lowestUsage.getUtilization() > lowerThreshold) { - return null; - } + private final StorageType storageType; + private final List volumeUsages; + private final ThresholdRange thresholdRange; + private final double violation; - // Determine source volume: highest utilization volume - final VolumeFixedUsage srcUsage = highestUsage; - final HddsVolume src = srcUsage.getVolume(); - - // Find destination volume and container: try each dest with lower utilization than source - for (int i = 0; i < volumeUsages.size() - 1; i++) { - final VolumeFixedUsage dstUsage = volumeUsages.get(i); - final HddsVolume dst = dstUsage.getVolume(); - - // Check if destination has lower utilization than source and some usable space - if (dstUsage.getUtilization() < srcUsage.getUtilization() && - dstUsage.computeUsableSpace() > 0) { - ContainerData containerData = chooseContainer(ozoneContainer, - src, dst, dstUsage, inProgressContainerIDs, upperThreshold, movableContainerStates); - if (containerData != null) { - long containerSize = containerData.getBytesUsed(); - dst.incCommittedBytes(containerSize); - LOG.debug("Chosen volume pair for disk balancing: source={} (utilization={}), " - + "destination={} (utilization={})", - src.getStorageDir().getPath(), srcUsage.getUtilization(), - dst.getStorageDir().getPath(), dstUsage.getUtilization()); - return new ContainerCandidate(containerData, src, dst); - } - LOG.debug("No container to move for destination {}, trying next volume.", - dst.getStorageDir().getPath()); - } else { - LOG.debug("Destination volume {} does not have enough space, trying next volume.", - dst.getStorageDir().getPath()); + private StorageTypeGroup(StorageType storageType, List volumeUsages, + ThresholdRange thresholdRange) { + this.storageType = storageType; + this.volumeUsages = volumeUsages; + this.thresholdRange = thresholdRange; + this.violation = thresholdRange.getViolation(volumeUsages); + } + + private StorageType getStorageType() { + return storageType; + } + + private List getVolumeUsages() { + return volumeUsages; + } + + private ThresholdRange getThresholdRange() { + return thresholdRange; + } + + private double getViolation() { + return violation; + } + } + + /** + * Chooses a source volume, destination volume and container among volumes that all share one + * storage type. + * + * @param group usable volumes of a single storage type, sorted ascending by utilization, with + * the threshold range already computed for them + * @return a candidate, or null if this storage type has nothing to move + */ + private ContainerCandidate chooseWithinStorageType(OzoneContainer ozoneContainer, + StorageTypeGroup group, Map deltaMap, + Set inProgressContainerIDs, double thresholdPercentage, + Set movableContainerStates) { + final StorageType storageType = group.getStorageType(); + final List volumeUsages = group.getVolumeUsages(); + final ThresholdRange thresholdRange = group.getThresholdRange(); + final double upperThreshold = thresholdRange.getUpperThreshold(); + + if (LOG.isDebugEnabled()) { + logVolumeBalancingState(volumeUsages, thresholdRange.getIdealUsage(), thresholdPercentage, + thresholdRange.getLowerThreshold(), upperThreshold, deltaMap); + } + + // Nothing to move while every volume sits inside the threshold range. + if (thresholdRange.isWithinRange(volumeUsages)) { + return null; + } + + final VolumeFixedUsage highestUsage = volumeUsages.get(volumeUsages.size() - 1); + + // Determine source volume: highest utilization volume + final VolumeFixedUsage srcUsage = highestUsage; + final HddsVolume src = srcUsage.getVolume(); + + // Find destination volume and container: try each dest with lower utilization than source + for (int i = 0; i < volumeUsages.size() - 1; i++) { + final VolumeFixedUsage dstUsage = volumeUsages.get(i); + final HddsVolume dst = dstUsage.getVolume(); + + // Check if destination has lower utilization than source and some usable space + if (dstUsage.getUtilization() < srcUsage.getUtilization() && + dstUsage.computeUsableSpace() > 0) { + ContainerData containerData = chooseContainer(ozoneContainer, + src, dst, dstUsage, inProgressContainerIDs, upperThreshold, movableContainerStates); + if (containerData != null) { + long containerSize = containerData.getBytesUsed(); + dst.incCommittedBytes(containerSize); + LOG.debug("Chosen volume pair for disk balancing on storage type {}: source={} " + + "(utilization={}), destination={} (utilization={})", + storageType, src.getStorageDir().getPath(), srcUsage.getUtilization(), + dst.getStorageDir().getPath(), dstUsage.getUtilization()); + return new ContainerCandidate(containerData, src, dst); } + LOG.debug("No container to move for destination {}, trying next volume.", + dst.getStorageDir().getPath()); + } else { + LOG.debug("Destination volume {} does not have enough space, trying next volume.", + dst.getStorageDir().getPath()); } - LOG.debug("Failed to find appropriate destination volume and container."); - return null; - } finally { - lock.unlock(); } + LOG.debug("Failed to find appropriate destination volume and container on storage type {}.", + storageType); + return null; } private static boolean hasPositiveCapacity(VolumeFixedUsage volumeUsage) { diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDefaultContainerChoosingPolicy.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDefaultContainerChoosingPolicy.java index 2dbeab1c5bf7..caf9e7d1e2e0 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDefaultContainerChoosingPolicy.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDefaultContainerChoosingPolicy.java @@ -41,6 +41,7 @@ import java.util.UUID; import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Stream; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.fs.MockSpaceUsageCheckFactory; @@ -257,6 +258,11 @@ public String toString() { private HddsVolume createVolume(String name, double utilization, long capacity) throws IOException { + return createVolume(name, utilization, capacity, StorageType.DEFAULT); + } + + private HddsVolume createVolume(String name, double utilization, long capacity, + StorageType storageType) throws IOException { long usedSpace = (long) (capacity * utilization); Path volumePath = baseDir.resolve(name); @@ -272,6 +278,7 @@ private HddsVolume createVolume(String name, double utilization, long capacity) return new HddsVolume.Builder(volumePath.toString()) .conf(volumeConf) .usageCheckFactory(factory) + .storageType(storageType) .build(); } @@ -452,6 +459,80 @@ public void testChooseVolumesSkipsZeroCapacityVolume() throws IOException { assertNull(result); } + /** + * When multiple storage types need balancing, the type with the largest threshold violation is + * selected even when another eligible type appears earlier in {@link StorageType#values()}. + * Source and destination remain within the selected storage type. + */ + @Test + public void testChoosesMostImbalancedStorageTypeWithoutCrossingTypes() throws IOException { + HddsVolume ssdHigh = createVolume("ssd-high", 0.75, VOLUME_CAPACITY, StorageType.SSD); + HddsVolume ssdLow = createVolume("ssd-low", 0.50, VOLUME_CAPACITY, StorageType.SSD); + HddsVolume diskHigh = createVolume("disk-high", 0.60, VOLUME_CAPACITY, StorageType.DISK); + HddsVolume diskMid = createVolume("disk-mid", 0.30, VOLUME_CAPACITY, StorageType.DISK); + HddsVolume diskLow = createVolume("disk-low", 0.20, VOLUME_CAPACITY, StorageType.DISK); + volumeSet = createVolumeSetForUsages( + Arrays.asList(ssdHigh, ssdLow, diskHigh, diskMid, diskLow)); + + containerSet = newContainerSet(); + createContainer(1L, DEFAULT_CONTAINER_SIZE, ssdHigh, containerSet); + createContainer(2L, DEFAULT_CONTAINER_SIZE, diskHigh, containerSet); + mockContainerSet(containerSet); + + ContainerCandidate result = policy.chooseVolumesAndContainer(ozoneContainer, + volumeSet, deltaMap, inProgressContainerIDs, THRESHOLD, DEFAULT_MOVABLE_STATES); + + assertNotNull(result); + assertEquals(diskHigh, result.getSourceVolume()); + assertEquals(diskLow, result.getDestVolume()); + assertEquals(StorageType.DISK, result.getSourceVolume().getStorageType()); + assertEquals(StorageType.DISK, result.getDestVolume().getStorageType()); + } + + /** + * A StorageType with a single volume has no same-type peer, so it is skipped instead of being + * paired with a volume of another type. Here the lone SSD volume is left alone and the DISK pair + * is balanced. + */ + @Test + public void testChooseVolumesSkipsStorageTypeWithSingleVolume() throws IOException { + HddsVolume lonelySsd = createVolume("ssd-only", 0.20, VOLUME_CAPACITY, StorageType.SSD); + HddsVolume diskHigh = createVolume("disk-high", 0.75, VOLUME_CAPACITY, StorageType.DISK); + HddsVolume diskLow = createVolume("disk-low", 0.50, VOLUME_CAPACITY, StorageType.DISK); + volumeSet = createVolumeSetForUsages(Arrays.asList(lonelySsd, diskHigh, diskLow)); + + containerSet = newContainerSet(); + createContainer(1L, DEFAULT_CONTAINER_SIZE, diskHigh, containerSet); + mockContainerSet(containerSet); + + ContainerCandidate result = policy.chooseVolumesAndContainer(ozoneContainer, + volumeSet, deltaMap, inProgressContainerIDs, THRESHOLD, DEFAULT_MOVABLE_STATES); + + assertNotNull(result); + assertEquals(diskHigh, result.getSourceVolume()); + assertEquals(diskLow, result.getDestVolume()); + } + + /** + * When every StorageType has a single volume there is no valid pair, even though the volumes are + * unevenly used. Nothing moves across types. + */ + @Test + public void testChooseVolumesReturnsNullWhenEveryStorageTypeHasOneVolume() throws IOException { + HddsVolume ssd = createVolume("ssd-only", 0.95, VOLUME_CAPACITY, StorageType.SSD); + HddsVolume disk = createVolume("disk-only", 0.10, VOLUME_CAPACITY, StorageType.DISK); + volumeSet = createVolumeSetForUsages(Arrays.asList(ssd, disk)); + + containerSet = newContainerSet(); + createContainer(1L, DEFAULT_CONTAINER_SIZE, ssd, containerSet); + mockContainerSet(containerSet); + + ContainerCandidate result = policy.chooseVolumesAndContainer(ozoneContainer, + volumeSet, deltaMap, inProgressContainerIDs, THRESHOLD, DEFAULT_MOVABLE_STATES); + + assertNull(result); + } + /** * Generic test method that can be reused for different scenarios. * diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerProtocolServer.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerProtocolServer.java index 03494eb78f9d..547e2a8ba8c4 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerProtocolServer.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerProtocolServer.java @@ -35,6 +35,8 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeDiskBalancerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerConfigurationProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerRunningStatus; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.VolumeReportProto; import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; import org.apache.hadoop.ozone.container.diskbalancer.DiskBalancerProtocolServer.PrivilegedOperation; @@ -101,6 +103,14 @@ void setup() throws IOException { TEST_VOLUME_DENSITY ); diskBalancerInfo.setIdealUsage(TEST_IDEAL_USAGE); + diskBalancerInfo.setStorageTypeInfo(Arrays.asList( + StorageTypeDiskBalancerInfoProto.newBuilder() + .setStorageType(StorageTypeProto.DISK) + .setCurrentVolumeDensitySum(TEST_VOLUME_DENSITY) + .setIdealUsage(TEST_IDEAL_USAGE) + .setUsableVolumeCount(TEST_VOLUME_INFO_COUNT) + .setBalanceable(true) + .build())); diskBalancerInfo.setVolumeInfo(Arrays.asList( VolumeReportProto.newBuilder() .setStorageId(TEST_STORAGE_ID_1) @@ -154,6 +164,8 @@ void testGetDiskBalancerInfoReport() throws IOException { assertEquals(TEST_VOLUME_DENSITY, report.getCurrentVolumeDensitySum()); assertEquals(TEST_IDEAL_USAGE, report.getIdealUsage()); assertEquals(TEST_VOLUME_INFO_COUNT, report.getVolumeInfoCount()); + assertEquals(1, report.getStorageTypeInfoCount()); + assertEquals(StorageTypeProto.DISK, report.getStorageTypeInfo(0).getStorageType()); assertEquals(TEST_STORAGE_ID_1, volReport0.getStorageId()); assertEquals(TEST_STORAGE_PATH_1, volReport0.getStoragePath()); assertEquals(TEST_UTILIZATION_1, volReport0.getUtilization()); @@ -325,4 +337,3 @@ void testUpdateRequiresAdmin() { exception.getMessage()); } } - diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerVolumeCalculation.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerVolumeCalculation.java index 4289af7afb7f..ab6b1f3ef916 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerVolumeCalculation.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/diskbalancer/TestDiskBalancerVolumeCalculation.java @@ -17,20 +17,28 @@ package org.apache.hadoop.ozone.container.diskbalancer; +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.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; import java.io.IOException; import java.nio.file.Path; import java.time.Duration; import java.util.Arrays; import java.util.Collections; +import java.util.List; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.fs.MockSpaceUsageCheckFactory; import org.apache.hadoop.hdds.fs.MockSpaceUsageSource; import org.apache.hadoop.hdds.fs.SpaceUsageCheckFactory; import org.apache.hadoop.hdds.fs.SpaceUsagePersistence; import org.apache.hadoop.hdds.fs.SpaceUsageSource; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.VolumeReportProto; import org.apache.hadoop.ozone.container.common.volume.HddsVolume; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; @@ -49,6 +57,66 @@ void getIdealUsageReturnsZeroForEmptyVolumeList() { Collections.emptyList())); } + @Test + void isBalanceableRequiresAtLeastTwoVolumes() throws IOException { + assertFalse(DiskBalancerVolumeCalculation.isBalanceable(null)); + assertFalse(DiskBalancerVolumeCalculation.isBalanceable(Collections.emptyList())); + assertFalse(DiskBalancerVolumeCalculation.isBalanceable( + Collections.singletonList(usage("solo", 100, 50)))); + assertTrue(DiskBalancerVolumeCalculation.isBalanceable( + Arrays.asList(usage("a", 100, 50), usage("b", 100, 50)))); + } + + /** + * The range is the ideal usage plus or minus the threshold, so a 10% threshold around a 50% + * ideal accepts everything between 40% and 60%. + */ + @Test + void thresholdRangeBracketsIdealUsage() throws IOException { + // 20 of 100 used on one volume, 80 of 100 on the other -> ideal 50%. + List volumes = + Arrays.asList(usage("low", 100, 80), usage("high", 100, 20)); + + DiskBalancerVolumeCalculation.ThresholdRange range = + DiskBalancerVolumeCalculation.ThresholdRange.of(volumes, 10.0); + + // Tolerance covers the small amount of space the volume builder reserves. + assertEquals(0.5, range.getIdealUsage(), 0.01); + assertEquals(0.4, range.getLowerThreshold(), 0.01); + assertEquals(0.6, range.getUpperThreshold(), 0.01); + } + + /** + * Violation measures how far the worst volume sits outside the band, and is not positive while + * every volume is inside it. + */ + @Test + void thresholdRangeReportsViolationOutsideBand() throws IOException { + // 20% and 80% against a 50% ideal with a 10% threshold -> 20 points past the upper bound. + List spread = + Arrays.asList(usage("low", 100, 80), usage("high", 100, 20)); + DiskBalancerVolumeCalculation.ThresholdRange wideRange = + DiskBalancerVolumeCalculation.ThresholdRange.of(spread, 10.0); + + assertEquals(0.2, wideRange.getViolation(spread), 0.01); + assertFalse(wideRange.isWithinRange(spread)); + + // Both volumes at 50% -> nothing outside the band. + List even = + Arrays.asList(usage("a", 100, 50), usage("b", 100, 50)); + DiskBalancerVolumeCalculation.ThresholdRange evenRange = + DiskBalancerVolumeCalculation.ThresholdRange.of(even, 10.0); + + assertThat(evenRange.getViolation(even)).isLessThanOrEqualTo(0.0); + assertTrue(evenRange.isWithinRange(even)); + } + + private DiskBalancerVolumeCalculation.VolumeFixedUsage usage(String name, long capacity, + long available) throws IOException { + return DiskBalancerVolumeCalculation.newVolumeFixedUsage( + createVolume(name, capacity, available), null); + } + @Test void getIdealUsageReturnsZeroForZeroTotalCapacity() throws IOException { HddsVolume zeroCapacityVolume = createVolume("zero-capacity", 0, 0); @@ -93,14 +161,53 @@ void getUtilizationReturnsZeroForZeroCapacityVolume() } @Test - void buildVolumeReportProtoReportsZeroUtilizationForZeroCapacityVolume() + void buildVolumeReportProtoIncludesStorageTypeForZeroCapacityVolume() throws IOException { - HddsVolume volume = createVolume("zero-capacity-report", 0, 0); + HddsVolume volume = createVolume("zero-capacity-report", 0, 0, StorageType.SSD); - assertEquals(0.0, DiskBalancerService.buildVolumeReportProto( + VolumeReportProto report = DiskBalancerService.buildVolumeReportProto( Collections.singletonList( - DiskBalancerVolumeCalculation.newVolumeFixedUsage(volume, null))) - .get(0).getUtilization()); + DiskBalancerVolumeCalculation.newVolumeFixedUsage(volume, null))).get(0); + + assertEquals(StorageTypeProto.SSD, report.getStorageType()); + assertEquals(0.0, report.getUtilization()); + } + + @Test + void buildStorageTypeInfoCalculatesEachTypeIndependently() throws IOException { + List volumeUsages = Arrays.asList( + DiskBalancerVolumeCalculation.newVolumeFixedUsage( + createVolume("ssd-1", 100, 20, StorageType.SSD), null), + DiskBalancerVolumeCalculation.newVolumeFixedUsage( + createVolume("ssd-2", 100, 20, StorageType.SSD), null), + DiskBalancerVolumeCalculation.newVolumeFixedUsage( + createVolume("disk-1", 100, 80, StorageType.DISK), null), + DiskBalancerVolumeCalculation.newVolumeFixedUsage( + createVolume("disk-2", 100, 80, StorageType.DISK), null)); + + List result = + DiskBalancerService.buildStorageTypeInfo(volumeUsages, 10.0, true); + + assertEquals(2, result.size()); + assertStorageTypeInfo(result.get(0), StorageTypeProto.SSD, 0.8); + assertStorageTypeInfo(result.get(1), StorageTypeProto.DISK, 0.2); + } + + @Test + void buildStorageTypeInfoMarksSingleVolumeTypeNotBalanceable() throws IOException { + List volumeUsages = Collections.singletonList( + DiskBalancerVolumeCalculation.newVolumeFixedUsage( + createVolume("ssd", 100, 5, StorageType.SSD), null)); + + StorageTypeDiskBalancerInfoProto result = + DiskBalancerService.buildStorageTypeInfo(volumeUsages, 10.0, true).get(0); + + assertEquals(StorageTypeProto.SSD, result.getStorageType()); + assertEquals(1, result.getUsableVolumeCount()); + assertFalse(result.getBalanceable()); + assertFalse(result.hasIdealUsage()); + assertEquals(0, result.getBytesToMove()); + assertEquals(0.0, result.getCurrentVolumeDensitySum()); } @Test @@ -155,6 +262,11 @@ void getIdealUsageRejectsEffectiveUsedGreaterThanCapacity() private HddsVolume createVolume(String name, long capacity, long available) throws IOException { + return createVolume(name, capacity, available, StorageType.DEFAULT); + } + + private HddsVolume createVolume(String name, long capacity, long available, + StorageType storageType) throws IOException { OzoneConfiguration conf = new OzoneConfiguration(); SpaceUsageSource source = MockSpaceUsageSource.fixed(capacity, available); SpaceUsageCheckFactory factory = MockSpaceUsageCheckFactory.of( @@ -163,6 +275,17 @@ private HddsVolume createVolume(String name, long capacity, long available) return new HddsVolume.Builder(tempDir.resolve(name).toString()) .conf(conf) .usageCheckFactory(factory) + .storageType(storageType) .build(); } + + private static void assertStorageTypeInfo(StorageTypeDiskBalancerInfoProto info, + StorageTypeProto expectedType, double expectedIdealUsage) { + assertEquals(expectedType, info.getStorageType()); + assertEquals(2, info.getUsableVolumeCount()); + assertTrue(info.getBalanceable()); + assertEquals(expectedIdealUsage, info.getIdealUsage(), 0.01); + assertEquals(0, info.getBytesToMove()); + assertEquals(0.0, info.getCurrentVolumeDensitySum()); + } } diff --git a/hadoop-hdds/interface-client/src/main/proto/hdds.proto b/hadoop-hdds/interface-client/src/main/proto/hdds.proto index b396af1966ff..6fc515ec8cbe 100644 --- a/hadoop-hdds/interface-client/src/main/proto/hdds.proto +++ b/hadoop-hdds/interface-client/src/main/proto/hdds.proto @@ -608,6 +608,16 @@ message VolumeReportProto { optional uint64 effectiveUsedSpace = 6; optional double utilization = 7; optional uint64 ozoneAvailable = 8; + optional StorageTypeProto storageType = 9; +} + +message StorageTypeDiskBalancerInfoProto { + optional StorageTypeProto storageType = 1; + optional double currentVolumeDensitySum = 2; + optional uint64 bytesToMove = 3; + optional double idealUsage = 4; + optional uint32 usableVolumeCount = 5; + optional bool balanceable = 6; } message DatanodeDiskBalancerInfoProto { @@ -621,4 +631,5 @@ message DatanodeDiskBalancerInfoProto { optional uint64 bytesMoved = 8; optional double idealUsage = 9; repeated VolumeReportProto volumeInfo = 10; + repeated StorageTypeDiskBalancerInfoProto storageTypeInfo = 11; } diff --git a/hadoop-hdds/interface-client/src/main/resources/proto.lock b/hadoop-hdds/interface-client/src/main/resources/proto.lock index 005b6ca73710..e7d6b6d86187 100644 --- a/hadoop-hdds/interface-client/src/main/resources/proto.lock +++ b/hadoop-hdds/interface-client/src/main/resources/proto.lock @@ -4625,6 +4625,53 @@ "name": "ozoneAvailable", "type": "uint64", "optional": true + }, + { + "id": 9, + "name": "storageType", + "type": "StorageTypeProto", + "optional": true + } + ] + }, + { + "name": "StorageTypeDiskBalancerInfoProto", + "fields": [ + { + "id": 1, + "name": "storageType", + "type": "StorageTypeProto", + "optional": true + }, + { + "id": 2, + "name": "currentVolumeDensitySum", + "type": "double", + "optional": true + }, + { + "id": 3, + "name": "bytesToMove", + "type": "uint64", + "optional": true + }, + { + "id": 4, + "name": "idealUsage", + "type": "double", + "optional": true + }, + { + "id": 5, + "name": "usableVolumeCount", + "type": "uint32", + "optional": true + }, + { + "id": 6, + "name": "balanceable", + "type": "bool", + "optional": true } ] }, @@ -4690,6 +4737,12 @@ "name": "volumeInfo", "type": "VolumeReportProto", "is_repeated": true + }, + { + "id": 11, + "name": "storageTypeInfo", + "type": "StorageTypeDiskBalancerInfoProto", + "is_repeated": true } ] } @@ -4718,4 +4771,4 @@ } } ] -} \ No newline at end of file +} diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java index 023d42758bba..d07efe520f8c 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java @@ -30,6 +30,8 @@ import org.apache.hadoop.hdds.cli.HddsVersionProvider; import org.apache.hadoop.hdds.protocol.DiskBalancerProtocol; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeDiskBalancerInfoProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.VolumeReportProto; import org.apache.hadoop.util.StringUtils; import picocli.CommandLine.Command; @@ -123,7 +125,10 @@ private String generateReport( .append(formatPercent(p.getCurrentVolumeDensitySum())) .append(System.lineSeparator()); - if (p.hasIdealUsage() && p.hasDiskBalancerConf() + if (p.getStorageTypeInfoCount() > 0 && p.hasDiskBalancerConf() + && p.getDiskBalancerConf().hasThreshold()) { + appendStorageTypeDetails(header, p); + } else if (p.hasIdealUsage() && p.hasDiskBalancerConf() && p.getDiskBalancerConf().hasThreshold()) { double idealUsage = p.getIdealUsage(); double threshold = p.getDiskBalancerConf().getThreshold(); @@ -141,8 +146,9 @@ private String generateReport( formatBuilder.append("%s%n"); contentList.add(header.toString()); - if (p.getVolumeInfoCount() > 0 && p.hasIdealUsage()) { - formatBuilder.append("%-45s %-40s %15s %15s %15s %30s %20s %15s %15s%n"); + if (p.getVolumeInfoCount() > 0 && (p.hasIdealUsage() || p.getStorageTypeInfoCount() > 0)) { + formatBuilder.append("%-12s %-45s %-40s %15s %15s %15s %30s %20s %15s %15s%n"); + contentList.add("StorageType"); contentList.add("StorageID"); contentList.add("StoragePath"); contentList.add("OzoneCapacity"); @@ -153,9 +159,11 @@ private String generateReport( contentList.add("Utilization"); contentList.add("VolumeDensity"); - double ideal = p.getIdealUsage(); + Map storageTypeInfo = + getStorageTypeInfo(p); for (VolumeReportProto v : p.getVolumeInfoList()) { - formatBuilder.append("%-45s %-40s %15s %15s %15s %30s %20s %15s %15s%n"); + formatBuilder.append("%-12s %-45s %-40s %15s %15s %15s %30s %20s %15s %15s%n"); + contentList.add(v.hasStorageType() ? v.getStorageType().name() : "-"); contentList.add(v.hasStorageId() ? v.getStorageId() : "-"); contentList.add(v.hasStoragePath() ? v.getStoragePath() : "-"); contentList.add(v.hasTotalCapacity() ? StringUtils.byteDesc(v.getTotalCapacity()) : "-"); @@ -164,7 +172,7 @@ private String generateReport( contentList.add(StringUtils.byteDesc(v.getCommittedBytes())); contentList.add(v.hasEffectiveUsedSpace() ? StringUtils.byteDesc(v.getEffectiveUsedSpace()) : "-"); contentList.add(formatPercent(v.getUtilization())); - contentList.add(formatPercent(Math.abs(v.getUtilization() - ideal))); + contentList.add(formatVolumeDensity(p, v, storageTypeInfo)); } formatBuilder.append("%n"); } @@ -175,7 +183,7 @@ private String generateReport( } formatBuilder.append("%nNote:%n") - .append(" - Aggregate VolumeDataDensity: Sum of per-volume density (deviation from ideal);") + .append(" - Aggregate VolumeDataDensity: Sum of per-volume density from each storage type's ideal;") .append(" higher means more imbalance.%n") .append(" - IdealUsage: Target utilization (0-100%%) when volumes are evenly balanced.%n") .append(" - ThresholdRange: Acceptable deviation (percent); volumes within") @@ -196,6 +204,53 @@ private String generateReport( return String.format(formatBuilder.toString(), contentList.toArray(new Object[0])); } + private static void appendStorageTypeDetails(StringBuilder header, + DatanodeDiskBalancerInfoProto report) { + double threshold = report.getDiskBalancerConf().getThreshold(); + header.append("Storage Type Details:").append(System.lineSeparator()); + for (StorageTypeDiskBalancerInfoProto info : report.getStorageTypeInfoList()) { + header.append(" ").append(info.getStorageType()).append(": "); + if (!info.getBalanceable() || !info.hasIdealUsage()) { + header.append("not balanceable (").append(info.getUsableVolumeCount()) + .append(" usable volume(s))").append(System.lineSeparator()); + continue; + } + double idealUsage = info.getIdealUsage(); + double lowerThreshold = Math.max(0.0, idealUsage - threshold / 100.0); + double upperThreshold = Math.min(1.0, idealUsage + threshold / 100.0); + header.append("IdealUsage: ").append(formatPercent(idealUsage)) + .append(" | ThresholdRange: (").append(formatPercent(lowerThreshold)) + .append(", ").append(formatPercent(upperThreshold)).append(')') + .append(" | VolumeDataDensity: ") + .append(formatPercent(info.getCurrentVolumeDensitySum())) + .append(" | EstBytesToMove: ").append(StringUtils.byteDesc(info.getBytesToMove())) + .append(System.lineSeparator()); + } + header.append(System.lineSeparator()).append("Volume Details:").append(System.lineSeparator()); + } + + private static Map getStorageTypeInfo( + DatanodeDiskBalancerInfoProto report) { + Map result = new LinkedHashMap<>(); + for (StorageTypeDiskBalancerInfoProto info : report.getStorageTypeInfoList()) { + result.put(info.getStorageType(), info); + } + return result; + } + + private static String formatVolumeDensity(DatanodeDiskBalancerInfoProto report, + VolumeReportProto volume, + Map storageTypeInfo) { + if (volume.hasStorageType()) { + StorageTypeDiskBalancerInfoProto info = storageTypeInfo.get(volume.getStorageType()); + if (info != null && info.hasIdealUsage()) { + return formatPercent(Math.abs(volume.getUtilization() - info.getIdealUsage())); + } + } + return report.hasIdealUsage() + ? formatPercent(Math.abs(volume.getUtilization() - report.getIdealUsage())) : "-"; + } + @Override protected String getActionName() { return "report"; @@ -218,7 +273,35 @@ private Map toJson(String hostName, DatanodeDiskBalancerInfoProt result.put("status", "success"); result.put("volumeDensity", formatPercent(report.getCurrentVolumeDensitySum())); - if (report.hasIdealUsage() && report.hasDiskBalancerConf() + Map storageTypeInfo = + getStorageTypeInfo(report); + // Report ideal usage per storage type when the datanode sends it. The node-level + // idealUsage averages across storage types, which is not a target any move can reach + // on a datanode with more than one type, so it is only reported as a fallback for + // datanodes that predate the per-storage-type fields. + if (!storageTypeInfo.isEmpty()) { + double threshold = report.getDiskBalancerConf().getThreshold(); + List> storageTypes = new ArrayList<>(); + for (StorageTypeDiskBalancerInfoProto info : report.getStorageTypeInfoList()) { + Map storageType = new LinkedHashMap<>(); + storageType.put("storageType", info.getStorageType().name()); + storageType.put("balanceable", info.getBalanceable()); + storageType.put("usableVolumeCount", info.getUsableVolumeCount()); + storageType.put("volumeDensity", formatPercent(info.getCurrentVolumeDensitySum())); + storageType.put("estBytesToMove", StringUtils.byteDesc(info.getBytesToMove())); + if (info.hasIdealUsage()) { + double idealUsage = info.getIdealUsage(); + double lowerThreshold = Math.max(0.0, idealUsage - threshold / 100.0); + double upperThreshold = Math.min(1.0, idealUsage + threshold / 100.0); + storageType.put("idealUsage", formatPercent(idealUsage)); + storageType.put("thresholdRange", String.format("(%s, %s)", + formatPercent(lowerThreshold), formatPercent(upperThreshold))); + } + storageTypes.add(storageType); + } + result.put("storageTypes", storageTypes); + result.put("threshold %", String.format(Locale.ROOT, PERCENT_FORMAT, threshold)); + } else if (report.hasIdealUsage() && report.hasDiskBalancerConf() && report.getDiskBalancerConf().hasThreshold()) { double idealUsage = report.getIdealUsage(); double threshold = report.getDiskBalancerConf().getThreshold(); @@ -231,10 +314,10 @@ private Map toJson(String hostName, DatanodeDiskBalancerInfoProt } if (report.getVolumeInfoCount() > 0) { - double ideal = report.hasIdealUsage() ? report.getIdealUsage() : 0.0; List> vols = new ArrayList<>(); for (VolumeReportProto v : report.getVolumeInfoList()) { Map vm = new LinkedHashMap<>(); + vm.put("storageType", v.hasStorageType() ? v.getStorageType().name() : "-"); vm.put("storageId", v.getStorageId()); vm.put("storagePath", v.hasStoragePath() ? v.getStoragePath() : "-"); vm.put("ozoneCapacity", v.hasTotalCapacity() ? StringUtils.byteDesc(v.getTotalCapacity()) : "-"); @@ -244,7 +327,7 @@ private Map toJson(String hostName, DatanodeDiskBalancerInfoProt vm.put("effectiveUsedSpace", v.hasEffectiveUsedSpace() ? StringUtils.byteDesc(v.getEffectiveUsedSpace()) : "-"); vm.put("utilization", formatPercent(v.getUtilization())); - vm.put("volumeDensity", formatPercent(Math.abs(v.getUtilization() - ideal))); + vm.put("volumeDensity", formatVolumeDensity(report, v, storageTypeInfo)); vols.add(vm); } diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java index 8e1dacd76117..1458c8f87a8b 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java @@ -17,6 +17,7 @@ package org.apache.hadoop.hdds.scm.cli.datanode; +import static java.util.stream.Collectors.joining; import static java.util.stream.Collectors.toList; import java.io.IOException; @@ -28,6 +29,7 @@ import org.apache.hadoop.hdds.protocol.DiskBalancerProtocol; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeDiskBalancerInfoProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; import picocli.CommandLine.Command; /** @@ -102,7 +104,7 @@ protected void displayResults(List successNodes, List failedNode private String generateStatus( List protos, List datanodeDisplayNames) { StringBuilder formatBuilder = new StringBuilder("Status result:%n" + - "%-60s %-10s %-15s %-15s %-10s %-18s %-30s %-12s %-12s %-15s %-18s %-20s%n"); + "%-60s %-10s %-15s %-15s %-10s %-18s %-30s %-12s %-12s %-15s %-18s %-40s %-20s%n"); List contentList = new ArrayList<>(); contentList.add("Datanode"); @@ -116,11 +118,13 @@ private String generateStatus( contentList.add("FailureMove"); contentList.add("BytesMoved(MB)"); contentList.add("EstBytesToMove(MB)"); + contentList.add("StorageTypeEstBytesToMove(MB)"); contentList.add("EstTimeLeft(min)"); for (int i = 0; i < protos.size(); i++) { HddsProtos.DatanodeDiskBalancerInfoProto proto = protos.get(i); - formatBuilder.append("%-60s %-10s %-15s %-15s %-10s %-18s %-30s %-12s %-12s %-15s %-18s %-20s%n"); + formatBuilder.append( + "%-60s %-10s %-15s %-15s %-10s %-18s %-30s %-12s %-12s %-15s %-18s %-40s %-20s%n"); long estimatedTimeLeft = calculateEstimatedTimeLeft(proto); long bytesMovedMB = (long) Math.ceil(proto.getBytesMoved() / (1024.0 * 1024.0)); long bytesToMoveMB = (long) Math.ceil(proto.getBytesToMove() / (1024.0 * 1024.0)); @@ -142,6 +146,7 @@ private String generateStatus( contentList.add(String.valueOf(proto.getFailureMoveCount())); contentList.add(String.valueOf(bytesMovedMB)); contentList.add(String.valueOf(bytesToMoveMB)); + contentList.add(formatStorageTypeBytesToMove(proto)); contentList.add(estimatedTimeLeft >= 0 ? String.valueOf(estimatedTimeLeft) : "N/A"); } @@ -185,11 +190,41 @@ private Map createStatusResult( result.put("failureMove", status.getFailureMoveCount()); result.put("bytesMovedMB", (long) Math.ceil(status.getBytesMoved() / (1024.0 * 1024.0))); result.put("estBytesToMoveMB", (long) Math.ceil(status.getBytesToMove() / (1024.0 * 1024.0))); + if (status.getStorageTypeInfoCount() > 0) { + result.put("storageTypes", createStorageTypeResults(status)); + } long estimatedTimeLeft = calculateEstimatedTimeLeft(status); result.put("estTimeLeftMin", estimatedTimeLeft >= 0 ? estimatedTimeLeft : null); return result; } + private static String formatStorageTypeBytesToMove(DatanodeDiskBalancerInfoProto status) { + if (status.getStorageTypeInfoCount() == 0) { + return "-"; + } + return status.getStorageTypeInfoList().stream() + .map(info -> info.getBalanceable() + ? String.format("%s=%d", info.getStorageType(), + (long) Math.ceil(info.getBytesToMove() / (1024.0 * 1024.0))) + : info.getStorageType() + "=N/A") + .collect(joining(", ")); + } + + private static List> createStorageTypeResults( + DatanodeDiskBalancerInfoProto status) { + List> storageTypes = new ArrayList<>(); + for (StorageTypeDiskBalancerInfoProto info : status.getStorageTypeInfoList()) { + Map storageType = new LinkedHashMap<>(); + storageType.put("storageType", info.getStorageType().name()); + storageType.put("balanceable", info.getBalanceable()); + storageType.put("usableVolumeCount", info.getUsableVolumeCount()); + storageType.put("estBytesToMoveMB", + (long) Math.ceil(info.getBytesToMove() / (1024.0 * 1024.0))); + storageTypes.add(storageType); + } + return storageTypes; + } + private long calculateEstimatedTimeLeft(DatanodeDiskBalancerInfoProto proto) { long bytesToMove = proto.getBytesToMove(); diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java index 1b5ef945ed11..60666684907f 100644 --- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java +++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java @@ -53,6 +53,8 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeDiskBalancerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerConfigurationProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerRunningStatus; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeDiskBalancerInfoProto; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.VolumeReportProto; import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; import org.junit.jupiter.api.AfterEach; @@ -455,6 +457,8 @@ public void testStatusDiskBalancerWithInServiceDatanodes() throws Exception { String output = outContent.toString(DEFAULT_ENCODING); assertTrue(output.contains("Status result")); + assertTrue(output.contains("StorageTypeEstBytesToMove(MB)")); + assertTrue(output.contains("DISK=")); assertTrue(output.contains("host-1")); assertTrue(output.contains("host-2")); assertTrue(output.contains("host-3")); @@ -483,6 +487,8 @@ public void testStatusDiskBalancerWithJson() throws Exception { assertTrue(output.contains("\"bandwidthInMB\"")); assertTrue(output.contains("\"threads\"")); assertTrue(output.contains("\"stopAfterDiskEven\"")); + assertTrue(output.contains("\"storageTypes\"")); + assertTrue(output.contains("\"storageType\"")); } } @@ -885,6 +891,33 @@ public void testReportThresholdRangeClamped(double idealUsage, } } + /** + * On a datanode with more than one storage type, the JSON report must carry the per-type ideal + * usages and not the node-level average, which is not a target any container move can reach. + */ + @Test + public void testReportJsonOmitsCrossStorageTypeIdealUsage() throws Exception { + DiskBalancerReportSubcommand cmd = new DiskBalancerReportSubcommand(); + + when(mockProtocol.getDiskBalancerInfo()) + .thenReturn(createMixedStorageTypeReportProto("host-1")); + + try (DiskBalancerMocks mocks = setupAllMocks()) { + CommandLine c = new CommandLine(cmd); + c.parseArgs("--json", "host-1"); + cmd.call(); + + String output = outContent.toString(DEFAULT_ENCODING); + // Per-type ideal usages are reported. + assertThat(output).contains("\"storageTypes\""); + assertThat(output).contains("20.00%"); + assertThat(output).contains("80.00%"); + // The node-level average across storage types is not. + assertThat(output).doesNotContain("50.00%"); + assertThat(output).doesNotContain("\"thresholdRange\" : \"(40.00%, 60.00%)\""); + } + } + @Test public void testReportDiskBalancerWithInServiceDatanodes() throws Exception { DiskBalancerReportSubcommand cmd = new DiskBalancerReportSubcommand(); @@ -905,6 +938,8 @@ public void testReportDiskBalancerWithInServiceDatanodes() throws Exception { String output = outContent.toString(DEFAULT_ENCODING); assertTrue(output.contains("Report result")); + assertTrue(output.contains("Storage Type Details:")); + assertTrue(output.contains("StorageType")); assertTrue(output.contains("host-1")); assertTrue(output.contains("host-2")); assertTrue(output.contains("host-3")); @@ -930,7 +965,9 @@ public void testReportDiskBalancerWithJson() throws Exception { assertTrue(output.contains("\"datanode\"")); assertTrue(output.contains("\"volumeDensity\"")); assertTrue(output.contains("\"idealUsage\"")); + assertTrue(output.contains("\"storageTypes\"")); assertTrue(output.contains("\"volumes\"")); + assertTrue(output.contains("\"storageType\"")); assertTrue(output.contains("\"storageId\"")); assertTrue(output.contains("\"storagePath\"")); assertTrue(output.contains("\"ozoneCapacity\"")); @@ -1061,6 +1098,12 @@ private DatanodeDiskBalancerInfoProto createStatusProto(String hostname, .setFailureMoveCount(failureMove) .setBytesMoved(bytesMoved) .setBytesToMove(bytesToMove) + .addStorageTypeInfo(StorageTypeDiskBalancerInfoProto.newBuilder() + .setStorageType(StorageTypeProto.DISK) + .setBalanceable(true) + .setUsableVolumeCount(2) + .setBytesToMove(bytesToMove) + .build()) .build(); } @@ -1120,6 +1163,7 @@ private DatanodeDiskBalancerInfoProto generateRandomReportProto(String hostname) String path1 = "/data/hdds-" + hostname + "-1"; String path2 = "/data/hdds-" + hostname + "-2"; VolumeReportProto vol1 = VolumeReportProto.newBuilder() + .setStorageType(StorageTypeProto.DISK) .setStorageId("DISK-" + hostname + "-vol1") .setStoragePath(path1) .setUtilization(util1) @@ -1130,6 +1174,7 @@ private DatanodeDiskBalancerInfoProto generateRandomReportProto(String hostname) .setEffectiveUsedSpace(effective1) .build(); VolumeReportProto vol2 = VolumeReportProto.newBuilder() + .setStorageType(StorageTypeProto.DISK) .setStorageId("DISK-" + hostname + "-vol2") .setStoragePath(path2) .setUtilization(util2) @@ -1145,6 +1190,13 @@ private DatanodeDiskBalancerInfoProto generateRandomReportProto(String hostname) .setCurrentVolumeDensitySum(volumeDensity) .setIdealUsage(idealUsage) .setDiskBalancerConf(configProto) + .addStorageTypeInfo(StorageTypeDiskBalancerInfoProto.newBuilder() + .setStorageType(StorageTypeProto.DISK) + .setCurrentVolumeDensitySum(volumeDensity) + .setIdealUsage(idealUsage) + .setUsableVolumeCount(2) + .setBalanceable(true) + .build()) .addVolumeInfo(vol1) .addVolumeInfo(vol2) .build(); @@ -1169,6 +1221,42 @@ private DatanodeDiskBalancerInfoProto createReportProto(String hostname, double .build(); } + /** + * A datanode with SSD volumes at 20% and DISK volumes at 80%. Each storage type is balanced + * within itself, but the node-level idealUsage averages to 50%, which no move can reach. + */ + private DatanodeDiskBalancerInfoProto createMixedStorageTypeReportProto(String hostname) { + DatanodeDetailsProto nodeProto = DatanodeDetailsProto.newBuilder() + .setHostName(hostname) + .setIpAddress("127.0.0.1") + .addPorts(HddsProtos.Port.newBuilder() + .setName("CLIENT_RPC") + .setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT) + .build()) + .build(); + + return DatanodeDiskBalancerInfoProto.newBuilder() + .setNode(nodeProto) + .setCurrentVolumeDensitySum(0.0) + .setIdealUsage(0.5) + .setDiskBalancerConf(createConfigProto(10.0, 100L, 5, true)) + .addStorageTypeInfo(StorageTypeDiskBalancerInfoProto.newBuilder() + .setStorageType(StorageTypeProto.SSD) + .setUsableVolumeCount(2) + .setBalanceable(true) + .setIdealUsage(0.2) + .setCurrentVolumeDensitySum(0.0) + .build()) + .addStorageTypeInfo(StorageTypeDiskBalancerInfoProto.newBuilder() + .setStorageType(StorageTypeProto.DISK) + .setUsableVolumeCount(2) + .setBalanceable(true) + .setIdealUsage(0.8) + .setCurrentVolumeDensitySum(0.0) + .build()) + .build(); + } + private DiskBalancerConfigurationProto createConfigProto(double threshold, long bandwidthInMB, int parallelThread, boolean stopAfterDiskEven) { return DiskBalancerConfigurationProto.newBuilder()