Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .mvn/extensions.xml
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
<extension>
<groupId>com.gradle</groupId>
<artifactId>develocity-maven-extension</artifactId>
<version>2.5.0</version>
<version>2.6.0</version>
</extension>
<extension>
<groupId>com.gradle</groupId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -44,6 +45,8 @@ public class DiskBalancerInfo {
private double idealUsage;
// Report-only: per-volume info. NOT persisted.
private List<VolumeReportProto> volumeInfo;
// Report-only: per-storage-type balancing info. NOT persisted.
private List<StorageTypeDiskBalancerInfoProto> storageTypeInfo;

public DiskBalancerInfo(DiskBalancerRunningStatus operationalState, double threshold,
long bandwidthInMB, int parallelThread, boolean stopAfterDiskEven) {
Expand Down Expand Up @@ -250,6 +253,14 @@ public void setVolumeInfo(List<VolumeReportProto> volumeInfo) {
this.volumeInfo = volumeInfo;
}

public List<StorageTypeDiskBalancerInfoProto> getStorageTypeInfo() {
return storageTypeInfo != null ? storageTypeInfo : Collections.emptyList();
}

public void setStorageTypeInfo(List<StorageTypeDiskBalancerInfoProto> storageTypeInfo) {
this.storageTypeInfo = storageTypeInfo;
}

@Override
public boolean equals(Object o) {
if (this == o) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ public DatanodeDiskBalancerInfoProto getDiskBalancerInfo(GetDiskBalancerInfoRequ
.setRunningStatus(info.getOperationalState())
.setIdealUsage(info.getIdealUsage())
.addAllVolumeInfo(info.getVolumeInfo())
.addAllStorageTypeInfo(info.getStorageTypeInfo())
.build();
}

Expand Down Expand Up @@ -167,4 +168,3 @@ public void close() {
}
}


Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -691,7 +699,6 @@ private void deleteContainer(Container container) {
}
}

@VisibleForTesting
public void cleanupPendingDeletionContainers() {
// delete all pending deletion containers before stop the service
boolean ret;
Expand All @@ -717,26 +724,59 @@ private boolean tryCleanupOnePendingDeletionContainer() {

public DiskBalancerInfo getDiskBalancerInfo() {
final List<VolumeFixedUsage> 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<StorageTypeDiskBalancerInfoProto> 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),
metrics.getSuccessCount(),
metrics.getFailureCount(), bytesToMove, metrics.getSuccessBytes(), volumeDataDensity);
info.setIdealUsage(getIdealUsage(volumeUsages));
info.setVolumeInfo(buildVolumeReportProto(volumeUsages));
info.setStorageTypeInfo(storageTypeInfo);
return info;
}

static List<StorageTypeDiskBalancerInfoProto> buildStorageTypeInfo(
List<VolumeFixedUsage> volumeUsages, double thresholdPercentage,
boolean calculateBytes) {
final Map<StorageType, List<VolumeFixedUsage>> usagesByStorageType =
getUsableVolumesByStorageType(volumeUsages);
final List<StorageTypeDiskBalancerInfoProto> result = new ArrayList<>();
for (StorageType storageType : StorageType.values()) {
final List<VolumeFixedUsage> 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,
Expand All @@ -755,6 +795,7 @@ public static List<VolumeReportProto> buildVolumeReportProto(List<VolumeFixedUsa
HddsVolume volume = v.getVolume();
VolumeReportProto.Builder builder = VolumeReportProto.newBuilder()
.setStorageId(volume.getStorageID())
.setStorageType(StorageTypeUtils.getStorageTypeProto(volume.getStorageType()))
.setTotalCapacity(v.getUsage().getCapacity())
.setOzoneAvailable(v.getUsage().getAvailable())
.setUsedSpace(v.getUsage().getUsedSpace())
Expand All @@ -770,13 +811,19 @@ public static List<VolumeReportProto> buildVolumeReportProto(List<VolumeFixedUsa
}

public long calculateBytesToMove(List<VolumeFixedUsage> 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<VolumeFixedUsage> 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;

Expand Down Expand Up @@ -806,7 +853,6 @@ public DiskBalancerServiceMetrics getMetrics() {
return metrics;
}

@VisibleForTesting
public void setBalancedBytesInLastWindow(long bytes) {
this.balancedBytesInLastWindow.set(bytes);
}
Expand All @@ -815,17 +861,14 @@ public ContainerChoosingPolicy getVolumeContainerChoosingPolicy() {
return volumeContainerChoosingPolicy;
}

@VisibleForTesting
public void setVolumeContainerChoosingPolicy(ContainerChoosingPolicy volumeContainerChoosingPolicy) {
this.volumeContainerChoosingPolicy = volumeContainerChoosingPolicy;
}

@VisibleForTesting
public Set<ContainerID> getInProgressContainers() {
return inProgressContainers;
}

@VisibleForTesting
public Map<HddsVolume, Long> getDeltaSizes() {
return deltaSizes;
}
Expand Down Expand Up @@ -878,7 +921,6 @@ public void shutdown() {
}
}

@VisibleForTesting
public static void setInjector(FaultInjector instance) {
injector = instance;
}
Expand All @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -59,7 +61,90 @@ public static List<VolumeFixedUsage> 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<StorageType, List<VolumeFixedUsage>> getUsableVolumesByStorageType(
List<VolumeFixedUsage> 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<VolumeFixedUsage> 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<VolumeFixedUsage> 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<VolumeFixedUsage> 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<VolumeFixedUsage> 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.
*
Expand Down
Loading
Loading