From 03ce5393f091c482dbb2f7a858aa68e98e15e64a Mon Sep 17 00:00:00 2001 From: Devesh Singh Date: Sat, 3 Oct 2026 19:01:36 +0530 Subject: [PATCH 1/2] HDDS-16653. SCM ContainerBalancer supports balancing replicas within the same StorageType. Co-authored-by: https://github.com/xichen01 --- .../container/keyvalue/KeyValueHandler.java | 4 +- .../replication/ContainerImporter.java | 29 +++- .../replication/ContainerUploader.java | 9 +- .../replication/GrpcContainerUploader.java | 8 +- .../container/replication/PushReplicator.java | 3 +- .../replication/ReplicationTask.java | 12 ++ .../SendContainerOutputStream.java | 28 +++- .../SendContainerRequestHandler.java | 10 +- .../commands/ReplicateContainerCommand.java | 29 +++- .../replication/TestContainerImporter.java | 22 +++ .../TestGrpcContainerUploader.java | 2 +- .../TestGrpcReplicationService.java | 4 +- .../replication/TestPushReplicator.java | 2 +- .../TestReplicationSupervisor.java | 2 +- .../TestSendContainerOutputStream.java | 50 +++++++ .../TestSendContainerRequestHandler.java | 2 +- .../TestReplicateContainerCommand.java | 98 ++++++++++++++ .../main/proto/DatanodeClientProtocol.proto | 3 + .../ScmServerDatanodeHeartbeatProtocol.proto | 4 + .../hdds/scm/SCMCommonPlacementPolicy.java | 38 +++++- .../hdds/scm/container/ContainerReplica.java | 13 ++ .../balancer/ContainerBalancerMetrics.java | 13 ++ .../ContainerBalancerSelectionCriteria.java | 41 ++++++ .../balancer/ContainerBalancerTask.java | 128 ++++++++++++++++-- .../scm/container/balancer/MoveManager.java | 18 ++- .../SCMContainerPlacementCapacity.java | 21 ++- .../SCMContainerPlacementRandom.java | 4 +- .../replication/ECMisReplicationHandler.java | 3 +- .../ECUnderReplicationHandler.java | 3 +- ...asiClosedStuckUnderReplicationHandler.java | 9 +- .../RatisMisReplicationHandler.java | 12 +- .../RatisUnderReplicationHandler.java | 37 ++++- .../replication/ReplicationManager.java | 39 ++++++ .../hdds/scm/node/DatanodeUsageInfo.java | 52 +++++++ .../scm/pipeline/PipelinePlacementPolicy.java | 2 +- ...estContainerBalancerSelectionCriteria.java | 61 +++++++++ .../container/balancer/TestMoveManager.java | 46 ++++++- .../TestSCMContainerPlacementCapacity.java | 98 ++++++++++++++ .../replication/ReplicationTestUtil.java | 4 +- .../TestECMisReplicationHandler.java | 2 +- .../TestECUnderReplicationHandler.java | 2 +- .../TestRatisMisReplicationHandler.java | 2 +- .../TestRatisUnderReplicationHandler.java | 2 +- .../hdds/scm/node/TestDatanodeUsageInfo.java | 57 ++++++++ 44 files changed, 967 insertions(+), 61 deletions(-) create mode 100644 hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java index 788718535381..db58e62540ce 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java @@ -551,7 +551,9 @@ private void populateContainerPathFields(KeyValueContainer container, HddsVolume hddsVolume) throws IOException { volumeSet.readLock(); try { - // TODO StoragePolicy Check whether need to adapt storageType + // The caller has already chosen hddsVolume for the container's storage + // type and recorded that type on the container data, so there is nothing + // storage-type specific to adapt here. String idDir = VersionedDatanodeFeatures.ScmHA.chooseContainerPathID( hddsVolume, clusterId); container.populatePathFields(idDir, hddsVolume); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java index 682e119c5537..de1144f3083a 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java @@ -18,6 +18,7 @@ package org.apache.hadoop.ozone.container.replication; import jakarta.annotation.Nonnull; +import jakarta.annotation.Nullable; import java.io.IOException; import java.io.InputStream; import java.nio.file.Files; @@ -26,6 +27,7 @@ import java.util.Collections; import java.util.HashSet; import java.util.Set; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; @@ -116,6 +118,12 @@ public void importContainer(long containerID, Path tarFilePath, } ContainerUtils.verifyContainerFileChecksum(containerData, conf); containerData.setVolume(targetVolume); + if (targetVolume != null) { + // The descriptor carries the source volume's storage type. Record the + // type of the volume actually chosen here, so the replica reports where + // it really lives rather than where its source lived. + containerData.setStorageType(targetVolume.getStorageType()); + } // lastDataScanTime should be cleared for an imported container containerData.setDataScanTimestamp(null); @@ -147,13 +155,26 @@ private static void deleteFileQuietely(Path tarFilePath) { } HddsVolume chooseNextVolume(long spaceToReserve) throws IOException { + return chooseNextVolume(spaceToReserve, null); + } + + /** + * Chooses a volume for an incoming container. + * + * @param spaceToReserve space the container needs in both tmp and dest dirs + * @param storageType storage type the container should be placed on, or + * null to allow any volume. Null is used for containers + * replicated from a datanode without storage type + * support. + */ + HddsVolume chooseNextVolume(long spaceToReserve, + @Nullable StorageType storageType) throws IOException { // Choose volume that can hold both container in tmp and dest directory - LOG.debug("Choosing volume to reserve space : {}", spaceToReserve); - // TODO: Use the target container storage type once replication/import - // requests carry it. Null preserves the existing any-volume behavior. + LOG.debug("Choosing volume to reserve space : {}, storageType: {}", + spaceToReserve, storageType); return volumeChoosingPolicy.chooseVolume( StorageVolumeUtil.getHddsVolumesList(volumeSet.getVolumesList()), - spaceToReserve, null); + spaceToReserve, storageType); } public static Path getUntarDirectory(HddsVolume hddsVolume) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerUploader.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerUploader.java index 55874511ba7c..3a929b3f2e70 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerUploader.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerUploader.java @@ -17,16 +17,23 @@ package org.apache.hadoop.ozone.container.replication; +import jakarta.annotation.Nullable; import java.io.IOException; import java.io.OutputStream; import java.util.concurrent.CompletableFuture; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.protocol.DatanodeDetails; /** * Client-side interface for sending a container to a target datanode. */ public interface ContainerUploader { + /** + * @param targetVolumeStorageType StorageType the target should place the + * container on, or null to let the target choose any volume + */ OutputStream startUpload(long containerId, DatanodeDetails target, - CompletableFuture callback, CopyContainerCompression compression) + CompletableFuture callback, CopyContainerCompression compression, + @Nullable StorageType targetVolumeStorageType) throws IOException; } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java index 3e726671cca5..376d8c535af1 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java @@ -18,10 +18,12 @@ package org.apache.hadoop.ozone.container.replication; import com.google.common.annotations.VisibleForTesting; +import jakarta.annotation.Nullable; import java.io.IOException; import java.io.OutputStream; import java.util.concurrent.CompletableFuture; import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeDetails.Port; @@ -59,7 +61,8 @@ public GrpcContainerUploader( @Override public OutputStream startUpload(long containerId, DatanodeDetails target, - CompletableFuture callback, CopyContainerCompression compression) throws IOException { + CompletableFuture callback, CopyContainerCompression compression, + @Nullable StorageType targetVolumeStorageType) throws IOException { // Get container size from local datanode instead of using passed replicateSize Long containerSize = null; @@ -82,7 +85,8 @@ public OutputStream startUpload(long containerId, DatanodeDetails target, (CallStreamObserver) client.upload( responseObserver), responseObserver); return new SendContainerOutputStream(requestStream, containerId, - GrpcReplicationService.BUFFER_SIZE, compression, containerSize) { + GrpcReplicationService.BUFFER_SIZE, compression, containerSize, + targetVolumeStorageType) { @Override public void close() throws IOException { try { diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/PushReplicator.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/PushReplicator.java index 759aff722baf..004106369c58 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/PushReplicator.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/PushReplicator.java @@ -61,7 +61,8 @@ public void replicate(ReplicationTask task) { CountingOutputStream output = null; try { output = new CountingOutputStream( - uploader.startUpload(containerID, target, fut, compression)); + uploader.startUpload(containerID, target, fut, compression, + task.getTargetVolumeStorageType())); source.copyData(containerID, output, compression); fut.get(); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java index ce8f535e0c4a..9af99c09438a 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java @@ -17,7 +17,9 @@ package org.apache.hadoop.ozone.container.replication; +import jakarta.annotation.Nullable; import java.util.Objects; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; @@ -79,6 +81,16 @@ public long getContainerId() { return cmd.getContainerID(); } + /** + * StorageType the replica should land on at the target, as chosen by SCM. Null + * when the command carries no storage type, which lets the target pick any + * volume. + */ + @Nullable + public StorageType getTargetVolumeStorageType() { + return cmd.getTargetVolumeStorageType(); + } + @Override protected Object getCommandForDebug() { return debugString; diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java index 3bb7e463d9d3..d33cdfcd71cd 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java @@ -17,6 +17,9 @@ package org.apache.hadoop.ozone.container.replication; +import jakarta.annotation.Nullable; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerRequest; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.grpc.stub.CallStreamObserver; @@ -28,14 +31,27 @@ class SendContainerOutputStream extends GrpcOutputStream { private final CopyContainerCompression compression; private final Long size; + /** + * StorageType the target should place the container on. Sent alongside size on + * the first request. Null lets the target choose any volume. + */ + private final StorageType targetVolumeStorageType; SendContainerOutputStream( CallStreamObserver streamObserver, long containerId, int bufferSize, CopyContainerCompression compression, Long size) { + this(streamObserver, containerId, bufferSize, compression, size, null); + } + + SendContainerOutputStream( + CallStreamObserver streamObserver, + long containerId, int bufferSize, CopyContainerCompression compression, + Long size, @Nullable StorageType targetVolumeStorageType) { super(streamObserver, containerId, bufferSize); this.compression = compression; this.size = size; + this.targetVolumeStorageType = targetVolumeStorageType; } @Override @@ -46,9 +62,15 @@ protected void sendPart(boolean eof, int length, ByteString data) { .setOffset(getWrittenBytes()) .setCompression(compression.toProto()); - // Include container size in the first request - if (getWrittenBytes() == 0 && size != null) { - requestBuilder.setSize(size); + // Include container size and target storage type in the first request + if (getWrittenBytes() == 0) { + if (size != null) { + requestBuilder.setSize(size); + } + if (targetVolumeStorageType != null) { + requestBuilder.setStorageTypeID( + StorageTypeUtils.getID(targetVolumeStorageType)); + } } getStreamObserver().onNext(requestBuilder.build()); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerRequestHandler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerRequestHandler.java index 0824341127c3..cde899222ab8 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerRequestHandler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerRequestHandler.java @@ -23,6 +23,8 @@ import java.io.OutputStream; import java.nio.file.Files; import java.nio.file.Path; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerRequest; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerResponse; @@ -90,7 +92,13 @@ public void onNext(SendContainerRequest req) { spaceToReserve = importer.getSpaceToReserve( req.hasSize() ? req.getSize() : null); - volume = importer.chooseNextVolume(spaceToReserve); + // Keep the replica on the same storage type as its source. Absent for + // senders without storage type support, which allows any volume. + StorageType storageType = req.hasStorageTypeID() + ? StorageTypeUtils.getStorageTypeFromID(req.getStorageTypeID()) + : null; + + volume = importer.chooseNextVolume(spaceToReserve, storageType); Path dir = ContainerImporter.getUntarDirectory(volume); Files.createDirectories(dir); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java index c35439096532..8399a48440e7 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java @@ -17,7 +17,10 @@ package org.apache.hadoop.ozone.protocol.commands; +import jakarta.annotation.Nullable; import java.util.Objects; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto.Builder; @@ -36,6 +39,12 @@ public final class ReplicateContainerCommand private int replicaIndex = 0; private ReplicationCommandPriority priority = ReplicationCommandPriority.NORMAL; + /** + * StorageType the replica should be written to on the target, so a copy stays + * on the same tier as its source. Null for containers created before storage + * policy support, which leaves the target free to choose any volume. + */ + private StorageType targetVolumeStorageType; public static ReplicateContainerCommand toTarget(long containerID, DatanodeDetails target) { @@ -63,6 +72,10 @@ public void setPriority(ReplicationCommandPriority priority) { this.priority = priority; } + public void setTargetVolumeStorageType(@Nullable StorageType storageType) { + this.targetVolumeStorageType = storageType; + } + @Override public Type getType() { return SCMCommandProto.Type.replicateContainerCommand; @@ -81,6 +94,10 @@ public ReplicateContainerCommandProto getProto() { .setReplicaIndex(replicaIndex) .setTarget(targetDatanode.getProtoBufMessage()) .setPriority(priority); + if (targetVolumeStorageType != null) { + builder.setVolumeStorageType( + StorageTypeUtils.getStorageTypeProto(targetVolumeStorageType)); + } return builder.build(); } @@ -100,6 +117,10 @@ public static ReplicateContainerCommand getFromProtobuf( if (protoMessage.hasPriority()) { cmd.setPriority(protoMessage.getPriority()); } + if (protoMessage.hasVolumeStorageType()) { + cmd.setTargetVolumeStorageType( + StorageTypeUtils.getFromProtobuf(protoMessage.getVolumeStorageType())); + } return cmd; } @@ -119,6 +140,11 @@ public ReplicationCommandPriority getPriority() { return priority; } + @Nullable + public StorageType getTargetVolumeStorageType() { + return targetVolumeStorageType; + } + @Override public String toString() { return getType() @@ -129,6 +155,7 @@ public String toString() { + ", containerId=" + getContainerID() + ", replicaIndex=" + getReplicaIndex() + ", targetNode=" + targetDatanode - + ", priority=" + priority; + + ", priority=" + priority + + ", targetVolumeStorageType=" + targetVolumeStorageType; } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java index 553fb2678826..c299163b1bb9 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java @@ -24,6 +24,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.isNull; import static org.mockito.Mockito.anyLong; import static org.mockito.Mockito.atLeastOnce; @@ -50,6 +51,7 @@ import org.apache.commons.compress.archivers.tar.TarArchiveEntry; import org.apache.commons.compress.archivers.tar.TarArchiveOutputStream; import org.apache.commons.io.IOUtils; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.scm.ScmConfigKeys; @@ -249,6 +251,26 @@ public void testChooseNextVolumeStorageType() throws Exception { assertEquals(expectedVolume, importer.chooseNextVolume(spaceToReserve)); } + /** + * A requested storage type must reach the volume choosing policy, so an + * imported replica lands on the same tier as its source. + */ + @Test + public void testChooseNextVolumeHonoursRequestedStorageType() throws Exception { + VolumeChoosingPolicy policy = mock(VolumeChoosingPolicy.class); + HddsVolume ssdVolume = mock(HddsVolume.class); + long spaceToReserve = 100L; + when(policy.chooseVolume(anyList(), anyLong(), eq(StorageType.SSD))) + .thenReturn(ssdVolume); + ContainerImporter importer = new ContainerImporter(conf, containerSet, + controllerMock, volumeSet, policy); + + assertEquals(ssdVolume, + importer.chooseNextVolume(spaceToReserve, StorageType.SSD)); + verify(policy).chooseVolume(anyList(), eq(spaceToReserve), + eq(StorageType.SSD)); + } + private File containerTarFile(long id, ContainerData data) throws IOException { File yamlFile = new File(tempDir, "container.yaml"); ContainerDataYaml.createContainerFile(data, yamlFile); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcContainerUploader.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcContainerUploader.java index 4e53206e3778..7361d66c288a 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcContainerUploader.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcContainerUploader.java @@ -130,7 +130,7 @@ protected GrpcReplicationClient createReplicationClient( private static OutputStream startUpload(GrpcContainerUploader subject, CompletableFuture callback) throws IOException { DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails(); - return subject.startUpload(1, target, callback, NO_COMPRESSION); + return subject.startUpload(1, target, callback, NO_COMPRESSION, null); } /** diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java index 4d90a8328be8..6316236d3380 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java @@ -146,7 +146,7 @@ public void init() throws Exception { }).when(importer).importContainer(anyLong(), any(), any(), any()); doReturn(true).when(importer).isAllowedContainerImport(eq( CONTAINER_ID)); - when(importer.chooseNextVolume(anyLong())).thenReturn(new HddsVolume.Builder( + when(importer.chooseNextVolume(anyLong(), any())).thenReturn(new HddsVolume.Builder( Files.createDirectory(tempDir.resolve("ImporterDir")).toString()).conf( conf).build()); @@ -196,7 +196,7 @@ public void copyData(long containerId, OutputStream destination, OutputStream uploadStream = mock(OutputStream.class); ContainerUploader uploader = mock(ContainerUploader.class); - when(uploader.startUpload(anyLong(), any(), any(), any())) + when(uploader.startUpload(anyLong(), any(), any(), any(), any())) .thenReturn(uploadStream); PushReplicator pushReplicator = new PushReplicator(conf, source, uploader); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java index a4463410cea5..2758e59074d6 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java @@ -140,7 +140,7 @@ private ContainerReplicator createSubject( when( uploader.startUpload(eq(containerID), eq(target), - futureArgument.capture(), compressionArgument.capture() + futureArgument.capture(), compressionArgument.capture(), any() )) .thenReturn(outputStream); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java index 21efce17dee9..af5389226f7d 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java @@ -329,7 +329,7 @@ public void testPushReplicatorTargetReturnsError(ContainerLayoutVersion layout) ContainerUploader uploader = mock(ContainerUploader.class); // Have the uploader's startUpload immediately complete the future exceptionally, // simulating the target returning an error before any data is transferred. - when(uploader.startUpload(anyLong(), any(), any(), any())) + when(uploader.startUpload(anyLong(), any(), any(), any(), any())) .thenAnswer(invocation -> { CompletableFuture fut = invocation.getArgument(2); fut.completeExceptionally( diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java index 716bf4d3aebc..5e2fed3b1280 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java @@ -19,13 +19,18 @@ import static org.apache.hadoop.ozone.container.replication.CopyContainerCompression.NO_COMPRESSION; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.mockito.Mockito.verify; import java.io.OutputStream; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerRequest; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; +import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.EnumSource; +import org.mockito.ArgumentCaptor; /** * Test for {@link SendContainerOutputStream}. @@ -64,6 +69,51 @@ void usesCompression(CopyContainerCompression compression) throws Exception { verify(getObserver()).onCompleted(); } + /** + * The target storage type travels on the first request, next to size, so the + * importing datanode can place the container on a matching volume. + */ + @ParameterizedTest + @EnumSource(StorageType.class) + void sendsTargetVolumeStorageType(StorageType storageType) throws Exception { + OutputStream subject = new SendContainerOutputStream( + getObserver(), getContainerId(), getBufferSize(), NO_COMPRESSION, null, + storageType); + + byte[] bytes = getRandomBytes(16); + subject.write(bytes, 0, bytes.length); + subject.close(); + + SendContainerRequest req = SendContainerRequest.newBuilder() + .setContainerID(getContainerId()) + .setOffset(0) + .setData(ByteString.copyFrom(bytes)) + .setCompression(NO_COMPRESSION.toProto()) + .setStorageTypeID(StorageTypeUtils.getID(storageType)) + .build(); + + verify(getObserver()).onNext(req); + verify(getObserver()).onCompleted(); + } + + /** + * Without a target storage type the request must not claim one, so the target + * keeps its any-volume behaviour. + */ + @Test + void omitsStorageTypeWhenNotRequested() throws Exception { + OutputStream subject = createSubject(); + + byte[] bytes = getRandomBytes(16); + subject.write(bytes, 0, bytes.length); + subject.close(); + + ArgumentCaptor captor = + ArgumentCaptor.forClass(SendContainerRequest.class); + verify(getObserver()).onNext(captor.capture()); + assertFalse(captor.getValue().hasStorageTypeID()); + } + @Override protected ByteString verifyPart(SendContainerRequest response, int expectedOffset, int size) { diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerRequestHandler.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerRequestHandler.java index 6a7bfa1063eb..c0f90a63aa61 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerRequestHandler.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerRequestHandler.java @@ -113,7 +113,7 @@ void testNoSpaceOnTargetVolume() throws Exception { // import the incoming container by having the volume chooser throw. DiskOutOfSpaceException noSpace = new DiskOutOfSpaceException("No volumes have enough space for a new container"); - doThrow(noSpace).when(importer).chooseNextVolume(anyLong()); + doThrow(noSpace).when(importer).chooseNextVolume(anyLong(), any()); doAnswer(invocation -> { Object arg = invocation.getArgument(0); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java new file mode 100644 index 000000000000..f99314a96958 --- /dev/null +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java @@ -0,0 +1,98 @@ +/* + * 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.ozone.protocol.commands; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; + +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; + +/** + * Tests that {@link ReplicateContainerCommand} carries the target volume + * storage type across protobuf serialization, so a replica can be placed on the + * same tier as its source. + */ +public class TestReplicateContainerCommand { + + @ParameterizedTest + @EnumSource(StorageType.class) + public void targetVolumeStorageTypeSurvivesRoundTrip(StorageType storageType) { + ReplicateContainerCommand command = ReplicateContainerCommand.toTarget( + 1L, MockDatanodeDetails.randomDatanodeDetails()); + command.setTargetVolumeStorageType(storageType); + + ReplicateContainerCommandProto proto = command.getProto(); + assertEquals(storageType, + ReplicateContainerCommand.getFromProtobuf(proto) + .getTargetVolumeStorageType()); + } + + /** + * A command without a storage type must not claim one on the wire, so the + * importing datanode keeps its any-volume behaviour. + */ + @Test + public void unsetStorageTypeIsAbsentFromProto() { + ReplicateContainerCommand command = ReplicateContainerCommand.toTarget( + 1L, MockDatanodeDetails.randomDatanodeDetails()); + + assertNull(command.getTargetVolumeStorageType()); + assertFalse(command.getProto().hasVolumeStorageType()); + assertNull(ReplicateContainerCommand + .getFromProtobuf(command.getProto()).getTargetVolumeStorageType()); + } + + /** + * Guards the upgrade path: a command from an SCM without storage type support + * deserializes with no storage type rather than defaulting to one. + */ + @Test + public void protoWithoutStorageTypeDeserializesAsNull() { + ReplicateContainerCommandProto proto = ReplicateContainerCommandProto + .newBuilder() + .setCmdId(1L) + .setContainerID(2L) + .setTarget(MockDatanodeDetails.randomDatanodeDetails() + .getProtoBufMessage()) + .build(); + + assertNull(ReplicateContainerCommand.getFromProtobuf(proto) + .getTargetVolumeStorageType()); + } + + /** + * Setting null explicitly clears the type rather than being ignored, so a + * caller can opt out after the fact. + */ + @Test + public void settingNullClearsStorageType() { + ReplicateContainerCommand command = ReplicateContainerCommand.toTarget( + 1L, MockDatanodeDetails.randomDatanodeDetails()); + command.setTargetVolumeStorageType(StorageType.SSD); + command.setTargetVolumeStorageType(null); + + assertNull(command.getTargetVolumeStorageType()); + assertFalse(command.getProto().hasVolumeStorageType()); + } +} diff --git a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto index 9989dddc2bfd..96c8a6e9e094 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -571,6 +571,9 @@ message SendContainerRequest { optional int64 checksum = 4; optional CopyContainerCompressProto compression = 5; optional int64 size = 6; + // StorageType the importing datanode should place the container on. Sent on + // the first request only, alongside size. Absent from older senders. + optional int32 storageTypeID = 7; } message SendContainerResponse { diff --git a/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto b/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto index 2b0f285ca9a3..bc7bef6cc42f 100644 --- a/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto +++ b/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto @@ -433,6 +433,10 @@ message ReplicateContainerCommandProto { optional int32 replicaIndex = 4; optional DatanodeDetailsProto target = 5; optional ReplicationCommandPriority priority = 6 [default = NORMAL]; + // StorageType the replica should be written to on the target datanode, so the + // copy stays on the same tier as the source. Absent for commands from an SCM + // without storage type support, which leaves the target free to pick a volume. + optional StorageTypeProto volumeStorageType = 7; } /** diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java index e39ee266f977..cdd8aaef27ed 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java @@ -21,6 +21,7 @@ import com.google.common.base.Preconditions; import com.google.common.collect.Maps; import com.google.common.collect.Sets; +import jakarta.annotation.Nullable; import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; @@ -387,10 +388,27 @@ public static boolean hasEnoughSpace(DatanodeDetails datanodeDetails, public List getResultSet( int nodesRequired, List healthyNodes) throws SCMException { + return getResultSet(nodesRequired, healthyNodes, null); + } + + /** + * Picks the required number of nodes, comparing candidates on the given + * storage type where the policy ranks nodes by usage. + * + * @param storageType storage type the container will be placed on, or null to + * compare nodes on their overall usage. The healthy node + * list is already filtered to this storage type by + * {@link #chooseDatanodesInternal}; this only affects how + * the remaining candidates are ranked against each other. + */ + public List getResultSet( + int nodesRequired, List healthyNodes, + @Nullable StorageType storageType) + throws SCMException { List results = new ArrayList<>(); for (int x = 0; x < nodesRequired; x++) { // invoke the choose function defined in the derived classes. - DatanodeDetails nodeId = chooseNode(healthyNodes); + DatanodeDetails nodeId = chooseNode(healthyNodes, storageType); if (nodeId != null) { removePeers(nodeId, healthyNodes); results.add(nodeId); @@ -418,6 +436,24 @@ public List getResultSet( public abstract DatanodeDetails chooseNode( List healthyNodes); + /** + * Choose a datanode according to the policy, ranking candidates by their usage + * of the given storage type. + * + * Policies that rank nodes by usage should override this to compare on + * {@code storageType}; the default ignores it, which is correct for policies + * that choose on topology or at random. + * + * @param healthyNodes - Set of healthy nodes we can choose from. + * @param storageType - storage type to compare node usage on, or null to + * compare overall usage. + * @return DatanodeDetails + */ + public DatanodeDetails chooseNode(List healthyNodes, + @Nullable StorageType storageType) { + return chooseNode(healthyNodes); + } + /** * Default implementation to return the number of racks containers should span * to meet the placement policy. For simple policies that are not rack aware diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplica.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplica.java index bb6b3a6ec47f..74ad58597005 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplica.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplica.java @@ -150,6 +150,19 @@ public StorageType getVolumeStorageType() { return volumeStorageType; } + /** + * The storage type a copy of this replica should be placed on to stay on the + * same tier. Prefers the volume the replica actually lives on, since that is + * ground truth, and falls back to the type the container was created for. + * + * @return the storage type to target, or null when neither was reported, which + * leaves the target free to choose any volume + */ + @Nullable + public StorageType getTargetStorageTypeForCopy() { + return volumeStorageType != null ? volumeStorageType : storageType; + } + @Override public int hashCode() { return new HashCodeBuilder(61, 71) diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerMetrics.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerMetrics.java index b4acf2a2fe75..c373b797e9ae 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerMetrics.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerMetrics.java @@ -90,6 +90,11 @@ public final class ContainerBalancerMetrics { " all iterations of Container Balancer.") private MutableCounterLong numContainerMovesScheduled; + @Metric(about = "Number of times Container Balancer stopped scheduling " + + "moves in an iteration because it reached the maximum size to move " + + "per iteration, while sources and targets were still available.") + private MutableCounterLong numIterationLimitsSkipped; + /** * Create and register metrics named {@link ContainerBalancerMetrics#NAME} * for {@link ContainerBalancer}. @@ -363,4 +368,12 @@ public void resetNumContainerMovesFailedInLatestIteration() { numContainerMovesFailedInLatestIteration.incr( -getNumContainerMovesFailedInLatestIteration()); } + + public long getNumIterationLimitsSkipped() { + return numIterationLimitsSkipped.value(); + } + + public void incrementNumIterationLimitsSkipped() { + numIterationLimitsSkipped.incr(); + } } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerSelectionCriteria.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerSelectionCriteria.java index c6828e635bd1..c046c601f5d0 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerSelectionCriteria.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerSelectionCriteria.java @@ -17,6 +17,7 @@ package org.apache.hadoop.hdds.scm.container.balancer; +import jakarta.annotation.Nullable; import java.util.Collections; import java.util.Comparator; import java.util.HashMap; @@ -25,6 +26,7 @@ import java.util.NavigableSet; import java.util.Set; import java.util.TreeSet; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto; @@ -185,6 +187,21 @@ private Comparator orderContainersByUsedBytes() { */ public boolean shouldBeExcluded(ContainerID containerID, DatanodeDetails node, long sizeMovedAlready) { + return shouldBeExcluded(containerID, node, sizeMovedAlready, null); + } + + /** + * As {@link #shouldBeExcluded(ContainerID, DatanodeDetails, long)}, but also + * excludes containers whose replica on {@code node} is not on + * {@code storageType}. The balancer uses this to move data within one storage + * tier at a time, so a move never relocates data to a different tier. + * + * @param storageType the storage type being balanced, or null to consider + * containers on any storage type + */ + public boolean shouldBeExcluded(ContainerID containerID, + DatanodeDetails node, long sizeMovedAlready, + @Nullable StorageType storageType) { ContainerInfo container; //If includeContainers is specified, exclude containers not in the include list if (!includeContainers.isEmpty() && !includeContainers.contains(containerID)) { @@ -219,6 +236,13 @@ public boolean shouldBeExcluded(ContainerID containerID, return true; } + if (storageType != null + && !isReplicaOnStorageType(replicas, node, storageType)) { + // The replica on this source is not on the tier being balanced. Moving it + // would relocate data across tiers, so leave it for that tier's pass. + return true; + } + if (balancerConfiguration.getIncludeNonStandardContainers()) { return !isContainerClosedRelaxed(container, node, replicas) || !isContainerHealthyForMoveRelaxed(container, replicas) || @@ -406,6 +430,23 @@ Set getExcludeNotFoundContainers() { return excludeContainersNotFound; } + /** + * Whether the container's replica on {@code node} sits on {@code storageType}. + * A replica that does not report a storage type is treated as a match, so + * balancing still works against datanodes that predate storage type reporting. + */ + private static boolean isReplicaOnStorageType(Set replicas, + DatanodeDetails node, StorageType storageType) { + for (ContainerReplica replica : replicas) { + if (replica.getDatanodeDetails().equals(node)) { + StorageType replicaType = replica.getTargetStorageTypeForCopy(); + return replicaType == null || replicaType == storageType; + } + } + // No replica on this node; nothing to move from here. + return false; + } + private NavigableSet getCandidateContainers(DatanodeDetails node) { NavigableSet newSet = new TreeSet<>(orderContainersByUsedBytes().reversed()); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java index 7d879d34018d..782160f243d8 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java @@ -23,12 +23,14 @@ import static org.apache.hadoop.util.StringUtils.byteDesc; import com.google.common.annotations.VisibleForTesting; +import jakarta.annotation.Nullable; import java.io.IOException; import java.time.Duration; import java.time.OffsetDateTime; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.EnumSet; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; @@ -46,6 +48,7 @@ import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.stream.Collectors; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -96,6 +99,9 @@ public class ContainerBalancerTask implements Runnable { private ContainerBalancerConfiguration config; private ContainerBalancerMetrics metrics; private final ContainerMoveFailureTracker moveFailureTracker; + // Utilization bounds for the whole cluster. Storage types are balanced in + // separate passes, but all passes share these bounds, since a datanode's + // utilization is measured across all of its volumes. private double upperLimit; private double lowerLimit; private ContainerBalancerSelectionCriteria selectionCriteria; @@ -659,15 +665,84 @@ private boolean isValidSCMState() { } private IterationResult doIteration() { + moveSelectionToFutureMap = new ConcurrentHashMap<>(); + iterationResult = IterationResult.ITERATION_COMPLETED; + boolean isMoveGeneratedInThisIteration = false; + + // Balance one storage type at a time so a container never moves to a + // datanode volume of a different tier. Iterating StorageType.values() keeps + // the order stable across iterations. + for (StorageType storageType : storageTypesToBalance()) { + if (balanceStorageType(storageType)) { + isMoveGeneratedInThisIteration = true; + } + // Stop early when the balancer was stopped, or when the iteration-wide + // size limit is reached, since both apply across all storage types. + if (iterationResult == IterationResult.ITERATION_INTERRUPTED + || reachedMaxSizeToMovePerIteration()) { + break; + } + } + + checkIterationResults(isMoveGeneratedInThisIteration); + return iterationResult; + } + + /** + * Storage types worth balancing, that is those present on at least one of the + * unbalanced datanodes. Returns a single null element when no datanode reports + * a storage type, which runs one tier-agnostic pass and preserves the behaviour + * from before storage type support. + */ + private List storageTypesToBalance() { + Set present = EnumSet.noneOf(StorageType.class); + for (DatanodeUsageInfo node : overUtilizedNodes) { + present.addAll(node.getStorageTypes()); + } + for (DatanodeUsageInfo node : underUtilizedNodes) { + present.addAll(node.getStorageTypes()); + } + if (present.isEmpty()) { + return Collections.singletonList(null); + } + List ordered = new ArrayList<>(); + for (StorageType type : StorageType.values()) { + if (present.contains(type)) { + ordered.add(type); + } + } + return ordered; + } + + /** + * Runs one source-to-target matching pass restricted to a single storage type. + * + * @param storageType the tier to balance, or null to consider all tiers + * @return true if at least one move was generated + */ + private boolean balanceStorageType(@Nullable StorageType storageType) { + // Only consider nodes that actually have a volume of this storage type, + // otherwise a source with no capacity on this tier would be picked and then + // find no container to move. + List potentialTargets = + nodesWithStorageType(getPotentialTargets(), storageType); + List potentialSources = + nodesWithStorageType(getPotentialSources(), storageType); + if (potentialSources.isEmpty() || potentialTargets.isEmpty()) { + if (LOG.isDebugEnabled()) { + LOG.debug("Skipping storage type {}: {} sources and {} targets have a " + + "volume of this type.", storageType, potentialSources.size(), + potentialTargets.size()); + } + return false; + } + // note that potential and selected targets are updated in the following // loop - List potentialTargets = getPotentialTargets(); findTargetStrategy.reInitialize(potentialTargets, config, upperLimit); - findSourceStrategy.reInitialize(getPotentialSources(), config, lowerLimit); + findSourceStrategy.reInitialize(potentialSources, config, lowerLimit); - moveSelectionToFutureMap = new ConcurrentHashMap<>(); boolean isMoveGeneratedInThisIteration = false; - iterationResult = IterationResult.ITERATION_COMPLETED; boolean canAdaptWhenNearingLimits = true; boolean canAdaptOnReachingLimits = true; @@ -680,6 +755,9 @@ private IterationResult doIteration() { // break out if we've reached max size to move limit if (reachedMaxSizeToMovePerIteration()) { + // Record that this pass was cut short by the size limit rather than by + // running out of sources or targets. + metrics.incrementNumIterationLimitsSkipped(); break; } @@ -706,7 +784,8 @@ involve, take some action in adaptWhenNearingIterationLimits() break; } - ContainerMoveSelection moveSelection = matchSourceWithTarget(source); + ContainerMoveSelection moveSelection = + matchSourceWithTarget(source, storageType); if (moveSelection != null) { if (processMoveSelection(source, moveSelection)) { isMoveGeneratedInThisIteration = true; @@ -717,8 +796,25 @@ involve, take some action in adaptWhenNearingIterationLimits() } } - checkIterationResults(isMoveGeneratedInThisIteration); - return iterationResult; + return isMoveGeneratedInThisIteration; + } + + /** + * Filters nodes down to those with a volume of the given storage type. Returns + * the list unchanged when no storage type is given. + */ + private static List nodesWithStorageType( + List nodes, @Nullable StorageType storageType) { + if (storageType == null) { + return nodes; + } + List filtered = new ArrayList<>(nodes.size()); + for (DatanodeUsageInfo node : nodes) { + if (node.getStorageTypes().contains(storageType)) { + filtered.add(node); + } + } + return filtered; } private boolean processMoveSelection(DatanodeDetails source, @@ -780,7 +876,11 @@ private void checkIterationResults(boolean isMoveGeneratedInThisIteration) { If no move was generated during this iteration then we don't need to check the move results */ - iterationResult = IterationResult.CAN_NOT_BALANCE_ANY_MORE; + if (iterationResult != IterationResult.ITERATION_INTERRUPTED) { + // Keep an interrupted result: the balancer was stopped before it could + // generate a move, which is not the same as having nothing to balance. + iterationResult = IterationResult.CAN_NOT_BALANCE_ANY_MORE; + } } else { checkIterationMoveResults(); } @@ -882,7 +982,8 @@ private long cancelMovesThatExceedTimeoutDuration() { * * @return ContainerMoveSelection containing the selected target and container */ - private ContainerMoveSelection matchSourceWithTarget(DatanodeDetails source) { + private ContainerMoveSelection matchSourceWithTarget(DatanodeDetails source, + @Nullable StorageType storageType) { Set sourceContainerIDSet = selectionCriteria.getContainerIDSet(source); @@ -903,8 +1004,13 @@ private ContainerMoveSelection matchSourceWithTarget(DatanodeDetails source) { Set toRemoveContainerIds = new HashSet<>(); for (ContainerID containerId: sourceContainerIDSet) { if (selectionCriteria.shouldBeExcluded(containerId, source, - sizeScheduledForMoveInLatestIteration)) { - toRemoveContainerIds.add(containerId); + sizeScheduledForMoveInLatestIteration, storageType)) { + // Containers excluded only because they are on another tier must stay in + // the cached set, so the pass for that tier can still consider them. + if (selectionCriteria.shouldBeExcluded(containerId, source, + sizeScheduledForMoveInLatestIteration)) { + toRemoveContainerIds.add(containerId); + } continue; } moveSelection = findTargetStrategy.findTargetForContainerMove(source, diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/MoveManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/MoveManager.java index 3e066540531e..32dd89a061a9 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/MoveManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/MoveManager.java @@ -463,11 +463,14 @@ private void sendReplicateCommand( final DatanodeDetails src) throws ContainerReplicaNotFoundException, ContainerNotFoundException, NotLeaderException { - int replicaIndex = getContainerReplicaIndex( - containerInfo.containerID(), src); + ContainerReplica sourceReplica = + getContainerReplica(containerInfo.containerID(), src); long now = clock.millis(); + // Keep the moved replica on the same storage type as the source, so a move + // between datanodes does not silently change the tier the data lives on. replicationManager.sendLowPriorityReplicateContainerCommand(containerInfo, - replicaIndex, src, tgt, now + replicationTimeout); + sourceReplica.getReplicaIndex(), src, tgt, now + replicationTimeout, + sourceReplica.getTargetStorageTypeForCopy()); pendingMoves.get(containerInfo.containerID()).setMoveStartTime(now); } @@ -493,13 +496,18 @@ private void sendDeleteCommand( private int getContainerReplicaIndex( final ContainerID id, final DatanodeDetails dn) throws ContainerNotFoundException, ContainerReplicaNotFoundException { + return getContainerReplica(id, dn).getReplicaIndex(); + } + + private ContainerReplica getContainerReplica( + final ContainerID id, final DatanodeDetails dn) + throws ContainerNotFoundException, ContainerReplicaNotFoundException { Set replicas = containerManager.getContainerReplicas(id); return replicas.stream().filter(r -> r.getDatanodeDetails().equals(dn)) //there should not be more than one replica of a container on the same //datanode. handle this if found in the future. .findFirst().orElseThrow(() -> - new ContainerReplicaNotFoundException(id, dn)) - .getReplicaIndex(); + new ContainerReplicaNotFoundException(id, dn)); } @Override diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementCapacity.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementCapacity.java index 7de1e8f48a5d..4d87c6b9ca31 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementCapacity.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementCapacity.java @@ -114,7 +114,7 @@ protected List chooseDatanodesInternal( if (healthyNodes.size() == nodesRequired) { return healthyNodes; } - return getResultSet(nodesRequired, healthyNodes); + return getResultSet(nodesRequired, healthyNodes, storageType); } /** @@ -127,6 +127,19 @@ protected List chooseDatanodesInternal( */ @Override public DatanodeDetails chooseNode(List healthyNodes) { + return chooseNode(healthyNodes, null); + } + + /** + * {@inheritDoc} + * + * When a storage type is given, the two candidates are compared on their usage + * of that storage type rather than their overall usage, so a node that is + * lightly used overall but nearly full on the requested tier is not preferred. + */ + @Override + public DatanodeDetails chooseNode(List healthyNodes, + StorageType storageType) { metrics.incrDatanodeChooseAttemptCount(); int firstNodeNdx = getRand().nextInt(healthyNodes.size()); int secondNodeNdx = getRand().nextInt(healthyNodes.size()); @@ -143,8 +156,10 @@ public DatanodeDetails chooseNode(List healthyNodes) { getNodeManager().getNodeStat(firstNodeDetails); SCMNodeMetric secondNodeMetric = getNodeManager().getNodeStat(secondNodeDetails); - datanodeDetails = !firstNodeMetric.isGreater(secondNodeMetric.get()) - ? firstNodeDetails : secondNodeDetails; + boolean firstIsFuller = storageType == null + ? firstNodeMetric.isGreater(secondNodeMetric.get()) + : firstNodeMetric.isGreater(secondNodeMetric.get(), storageType); + datanodeDetails = !firstIsFuller ? firstNodeDetails : secondNodeDetails; } healthyNodes.remove(datanodeDetails); metrics.incrDatanodeChooseSuccessCount(); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRandom.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRandom.java index 846dd5bbe4f1..d1e452f73bd7 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRandom.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/SCMContainerPlacementRandom.java @@ -91,7 +91,9 @@ protected List chooseDatanodesInternal( if (healthyNodes.size() == nodesRequired) { return healthyNodes; } - return getResultSet(nodesRequired, healthyNodes); + // Selection is random, so storageType does not affect ranking here; the + // healthy node list is already filtered to it. + return getResultSet(nodesRequired, healthyNodes, storageType); } /** diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECMisReplicationHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECMisReplicationHandler.java index c52cac57f1c3..14f3c5a6463b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECMisReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECMisReplicationHandler.java @@ -72,9 +72,10 @@ protected int sendReplicateCommands( DatanodeDetails source = replica.getDatanodeDetails(); DatanodeDetails target = targetDns.get(datanodeIdx); try { + // Place the new replica on the same storage type as the one it replaces. replicationManager.sendThrottledReplicationCommand(containerInfo, Collections.singletonList(source), target, - replica.getReplicaIndex()); + replica.getReplicaIndex(), replica.getTargetStorageTypeForCopy()); commandsSent++; } catch (CommandTargetOverloadedException e) { LOG.debug("Unable to replicate container {} and index {} from {} to {}" diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECUnderReplicationHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECUnderReplicationHandler.java index 8a3bdea82fe2..da433309dd91 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECUnderReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ECUnderReplicationHandler.java @@ -598,9 +598,10 @@ private void createReplicateCommand( throws CommandTargetOverloadedException, NotLeaderException { DatanodeDetails source = replica.getDatanodeDetails(); DatanodeDetails target = iterator.next(); + // Place the new replica on the same storage type as the one being copied. replicationManager.sendThrottledReplicationCommand( container, Collections.singletonList(source), target, - replica.getReplicaIndex()); + replica.getReplicaIndex(), replica.getTargetStorageTypeForCopy()); adjustPendingOps(replicaCount, target, replica.getReplicaIndex()); } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/QuasiClosedStuckUnderReplicationHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/QuasiClosedStuckUnderReplicationHandler.java index f2668e0a2629..d9d3fa67a2d2 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/QuasiClosedStuckUnderReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/QuasiClosedStuckUnderReplicationHandler.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Objects; import java.util.Set; import java.util.stream.Collectors; import org.apache.hadoop.fs.StorageType; @@ -118,10 +119,16 @@ public int processAndSendCommands(Set replicas, List sourceDatanodes = origin.getSources().stream() .map(ContainerReplica::getDatanodeDetails) .collect(Collectors.toList()); + // Keep the new copies on the same tier as the replicas being copied. + StorageType targetStorageType = origin.getSources().stream() + .map(ContainerReplica::getTargetStorageTypeForCopy) + .filter(Objects::nonNull) + .findFirst() + .orElse(null); for (DatanodeDetails target : targets) { try { replicationManager.sendThrottledReplicationCommand( - containerInfo, sourceDatanodes, target, 0); + containerInfo, sourceDatanodes, target, 0, targetStorageType); // Add the pending op, so we exclude the node for subsequent origins mutablePendingOps.add(new ContainerReplicaOp( ContainerReplicaOp.PendingOpType.ADD, target, 0, diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisMisReplicationHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisMisReplicationHandler.java index 985694ec5ee7..bd13e888f316 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisMisReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisMisReplicationHandler.java @@ -19,7 +19,9 @@ import java.io.IOException; import java.util.List; +import java.util.Objects; import java.util.Set; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -63,10 +65,18 @@ protected int sendReplicateCommands( throws CommandTargetOverloadedException, NotLeaderException { ReplicationManager replicationManager = getReplicationManager(); + // All Ratis replicas are interchangeable, so any of the ones being replicated + // tells us which tier the copies belong on. + StorageType targetStorageType = replicasToBeReplicated.stream() + .map(ContainerReplica::getTargetStorageTypeForCopy) + .filter(Objects::nonNull) + .findFirst() + .orElse(null); + int commandsSent = 0; for (DatanodeDetails target : targetDns) { replicationManager.sendThrottledReplicationCommand(containerInfo, - sources, target, 0); + sources, target, 0, targetStorageType); commandsSent++; } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java index 84fc2ea80f9c..f808dd19736b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java @@ -22,6 +22,7 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; +import java.util.Objects; import java.util.OptionalLong; import java.util.Set; import java.util.function.Predicate; @@ -147,8 +148,10 @@ public int processAndSendCommands( throw e; } + // Keep new replicas on the same tier as the existing ones. int commandsSent = sendReplicationCommands( - containerInfo, sourceDatanodes, targetDatanodes); + containerInfo, sourceDatanodes, targetDatanodes, + targetStorageTypeFor(replicaCount.getReplicas())); if (targetDatanodes.size() < replicaCount.additionalReplicaNeeded()) { // The placement policy failed to find enough targets to satisfy fix @@ -244,7 +247,8 @@ of other replicas, including UNHEALTHY replicas that are not pending delete (bec excludedAndUsedNodes.getExcludedNodes(), currentContainerSize, container, StorageType.DEFAULT); int count = 0; try { - count = sendReplicationCommands(container, ImmutableList.of(replica.getDatanodeDetails()), target); + count = sendReplicationCommands(container, ImmutableList.of(replica.getDatanodeDetails()), target, + replica.getTargetStorageTypeForCopy()); } catch (CommandTargetOverloadedException e) { LOG.info("Exception while replicating {} to target {} for container {}.", replica, target, container, e); if (firstException == null) { @@ -471,10 +475,37 @@ private int sendReplicationCommands( ContainerInfo containerInfo, List sources, List targets) throws CommandTargetOverloadedException, NotLeaderException { + return sendReplicationCommands(containerInfo, sources, targets, null); + } + + /** + * All Ratis replicas of a container are interchangeable, so the tier of any + * reported replica tells us where new copies belong. + * + * @return the storage type new replicas should use, or null when no replica + * reported one + */ + private static StorageType targetStorageTypeFor( + List replicas) { + return replicas.stream() + .map(ContainerReplica::getTargetStorageTypeForCopy) + .filter(Objects::nonNull) + .findFirst() + .orElse(null); + } + + /** + * @param targetStorageType storage type the new replicas should land on, so a + * re-replicated container stays on its tier, or null to allow any volume + */ + private int sendReplicationCommands( + ContainerInfo containerInfo, List sources, + List targets, StorageType targetStorageType) + throws CommandTargetOverloadedException, NotLeaderException { int commandsSent = 0; for (DatanodeDetails target : targets) { replicationManager.sendThrottledReplicationCommand( - containerInfo, sources, target, 0); + containerInfo, sources, target, 0, targetStorageType); commandsSent++; } return commandsSent; diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java index 48f041ae19e5..ff7ed0be02cd 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java @@ -27,6 +27,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.protobuf.ByteString; +import jakarta.annotation.Nullable; import java.io.IOException; import java.time.Clock; import java.time.Duration; @@ -517,6 +518,23 @@ public void sendThrottledDeleteCommand(final ContainerInfo container, public void sendThrottledReplicationCommand(ContainerInfo containerInfo, List sources, DatanodeDetails target, int replicaIndex) throws CommandTargetOverloadedException, NotLeaderException { + sendThrottledReplicationCommand(containerInfo, sources, target, replicaIndex, + null); + } + + /** + * As {@link #sendThrottledReplicationCommand(ContainerInfo, List, + * DatanodeDetails, int)}, but asks the target to place the replica on a + * specific storage type so it stays on the same tier as the replica it is + * replacing. + * + * @param targetVolumeStorageType storage type for the new replica, or null to + * let the target choose any volume + */ + public void sendThrottledReplicationCommand(ContainerInfo containerInfo, + List sources, DatanodeDetails target, int replicaIndex, + @Nullable StorageType targetVolumeStorageType) + throws CommandTargetOverloadedException, NotLeaderException { long containerID = containerInfo.getContainerID(); List> sourceWithCmds = getAvailableDatanodesForReplication(sources); @@ -532,6 +550,7 @@ public void sendThrottledReplicationCommand(ContainerInfo containerInfo, ReplicateContainerCommand cmd = ReplicateContainerCommand.toTarget(containerID, target); cmd.setReplicaIndex(replicaIndex); + cmd.setTargetVolumeStorageType(targetVolumeStorageType); sendDatanodeCommand(cmd, containerInfo, source); } @@ -629,10 +648,30 @@ public void sendLowPriorityReplicateContainerCommand( final ContainerInfo container, int replicaIndex, DatanodeDetails source, DatanodeDetails target, long scmDeadlineEpochMs) throws NotLeaderException { + sendLowPriorityReplicateContainerCommand(container, replicaIndex, source, + target, scmDeadlineEpochMs, null); + } + + /** + * As {@link #sendLowPriorityReplicateContainerCommand(ContainerInfo, int, + * DatanodeDetails, DatanodeDetails, long)}, but asks the target to place the + * replica on a specific storage type. Used by the balancer so a move stays + * within one tier. + * + * @param targetVolumeStorageType storage type for the new replica, or null to + * let the target choose any volume + */ + @SuppressWarnings("checkstyle:parameternumber") + public void sendLowPriorityReplicateContainerCommand( + final ContainerInfo container, int replicaIndex, DatanodeDetails source, + DatanodeDetails target, long scmDeadlineEpochMs, + @Nullable StorageType targetVolumeStorageType) + throws NotLeaderException { final ReplicateContainerCommand command = ReplicateContainerCommand .toTarget(container.getContainerID(), target); command.setReplicaIndex(replicaIndex); command.setPriority(ReplicationCommandPriority.LOW); + command.setTargetVolumeStorageType(targetVolumeStorageType); sendDatanodeCommand(command, container, source, scmDeadlineEpochMs); } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java index d6ac9edfaeb1..5575c25179d5 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java @@ -17,7 +17,11 @@ package org.apache.hadoop.hdds.scm.node; +import jakarta.annotation.Nullable; import java.util.Comparator; +import java.util.EnumSet; +import java.util.Set; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto; @@ -104,6 +108,54 @@ public double calculateUtilization() { return calculateUtilization(0); } + /** + * Calculates utilization of a single storage type on this datanode, so the + * balancer can compare nodes on the tier it is actually moving data within. + * + * @param plusSize the increased size + * @param storageType the storage type to measure, or null to measure the + * datanode as a whole + * @return (capacity - remaining) / capacity for that storage type, or 0 when + * the datanode has no volume of that type + */ + public double calculateUtilization(long plusSize, + @Nullable StorageType storageType) { + if (storageType == null) { + return calculateUtilization(plusSize); + } + long capacity = scmNodeStat.getCapacity(storageType).get(); + if (capacity == 0) { + return 0; + } + long numerator = + capacity - scmNodeStat.getRemaining(storageType).get() + plusSize; + return numerator / (double) capacity; + } + + /** + * Calculates current utilization of a single storage type on this datanode. + * + * @param storageType the storage type to measure, or null to measure the + * datanode as a whole + */ + public double calculateUtilization(@Nullable StorageType storageType) { + return calculateUtilization(0, storageType); + } + + /** + * Storage types this datanode actually has capacity for. Used by the balancer + * to decide which tiers are worth considering on this node. + */ + public Set getStorageTypes() { + Set types = EnumSet.noneOf(StorageType.class); + for (StorageType type : StorageType.values()) { + if (scmNodeStat.getCapacity(type).get() > 0) { + types.add(type); + } + } + return types; + } + /** * Sets DatanodeDetails of this DatanodeUsageInfo. * diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java index 52456b948bb2..36e26e28910a 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java @@ -313,7 +313,7 @@ protected List chooseDatanodesInternal( // This happens when network topology is absent or // all nodes are on the same rack. if (checkAllNodesAreEqual(nodeManager.getClusterNetworkTopologyMap())) { - return super.getResultSet(nodesRequired, healthyNodes); + return super.getResultSet(nodesRequired, healthyNodes, storageType); } else { // Since topology and rack awareness are available, picks nodes // based on them. diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerSelectionCriteria.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerSelectionCriteria.java index 5b87a13369d5..f083d0d24dc6 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerSelectionCriteria.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerSelectionCriteria.java @@ -32,9 +32,11 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Set; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -98,6 +100,65 @@ public void setup() throws Exception { containerManager, findSourceStrategy, new HashMap<>()); } + /** + * Balancing a given storage tier must only consider containers whose replica on + * the source actually sits on that tier, otherwise a move would relocate data + * across tiers. + */ + @Test + public void shouldExcludeContainerOnDifferentStorageType() throws Exception { + replicaOnSourceWithStorageType(StorageType.DISK); + + assertFalse(criteria.shouldBeExcluded(containerID, source, 0L, + StorageType.DISK)); + assertTrue(criteria.shouldBeExcluded(containerID, source, 0L, + StorageType.SSD)); + } + + /** + * A replica that reports no storage type must stay eligible, so balancing keeps + * working against datanodes from before storage type reporting. + */ + @Test + public void shouldIncludeContainerWithUnknownStorageType() throws Exception { + replicaOnSourceWithStorageType(null); + + assertFalse(criteria.shouldBeExcluded(containerID, source, 0L, + StorageType.SSD)); + assertFalse(criteria.shouldBeExcluded(containerID, source, 0L, + StorageType.ARCHIVE)); + } + + /** + * Passing no storage type must behave exactly as before, considering containers + * on any tier. + */ + @Test + public void shouldIgnoreStorageTypeWhenNotGiven() throws Exception { + replicaOnSourceWithStorageType(StorageType.ARCHIVE); + + assertFalse(criteria.shouldBeExcluded(containerID, source, 0L)); + assertFalse(criteria.shouldBeExcluded(containerID, source, 0L, null)); + } + + /** + * Replaces the container's replica set with a single replica on {@link #source} + * reporting the given volume storage type. + */ + private void replicaOnSourceWithStorageType(StorageType storageType) + throws Exception { + ContainerReplica replica = ReplicationTestUtil.createContainerReplica( + containerID, 0, IN_SERVICE, CLOSED, 1L, OzoneConsts.GB, source, + source.getID()); + if (storageType != null) { + replica = replica.toBuilder() + .setVolumeStorageType(storageType) + .build(); + } + when(containerManager.getContainerReplicas(containerID)) + .thenReturn(new HashSet<>(Collections.singletonList(replica))); + } + @Test public void shouldExcludeUnderReplicatedContainer() { when(replicationManager.getContainerReplicationHealth(eq(containerInfo), anySet())).thenReturn( diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestMoveManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestMoveManager.java index aa22a4305eeb..6424f3763f49 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestMoveManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestMoveManager.java @@ -61,6 +61,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -317,7 +318,7 @@ public void testExistingMoveScheduled() throws Exception { public void testReplicationCommandFails() throws Exception { doThrow(new RuntimeException("test")).when(replicationManager) .sendLowPriorityReplicateContainerCommand( - any(), anyInt(), any(), any(), anyLong()); + any(), anyInt(), any(), any(), anyLong(), any()); CompletableFuture res = setupSuccessfulMove(); assertEquals(FAIL_UNEXPECTED_ERROR, res.get()); } @@ -356,6 +357,37 @@ public void testSuccessfulMove() throws Exception { assertEquals(COMPLETED, finalResult); } + /** + * A move must ask the target to place the new replica on the same storage type + * the source replica lives on, otherwise balancing silently moves data between + * tiers. + */ + @Test + public void testMoveKeepsSourceStorageType() throws Exception { + setupMocks(); + + ContainerReplica srcReplica = ReplicationTestUtil.createContainerReplica( + containerInfo.containerID(), 0, IN_SERVICE, + ContainerReplicaProto.State.CLOSED); + ContainerReplica withSsd = srcReplica.toBuilder() + .setVolumeStorageType(StorageType.SSD) + .build(); + replicas.add(withSsd); + replicas.addAll(ReplicationTestUtil.createReplicas( + containerInfo.containerID(), 0, 0)); + + src = withSsd.getDatanodeDetails(); + tgt = MockDatanodeDetails.randomDatanodeDetails(); + nodes.put(src, NodeStatus.inServiceHealthy()); + nodes.put(tgt, NodeStatus.inServiceHealthy()); + + moveManager.move(containerInfo.containerID(), src, tgt); + + verify(replicationManager).sendLowPriorityReplicateContainerCommand( + eq(containerInfo), eq(0), eq(src), eq(tgt), anyLong(), + eq(StorageType.SSD)); + } + @Test public void testSuccessfulMoveNonZeroRepIndex() throws Exception { containerInfo = ReplicationTestUtil.createContainer( @@ -376,7 +408,7 @@ public void testSuccessfulMoveNonZeroRepIndex() throws Exception { verify(replicationManager).sendLowPriorityReplicateContainerCommand( eq(containerInfo), eq(srcReplica.getReplicaIndex()), eq(src), eq(tgt), - anyLong()); + anyLong(), any()); ContainerReplicaOp op = new ContainerReplicaOp( ADD, tgt, srcReplica.getReplicaIndex(), null, clock.millis() + 1000, 0, null); @@ -525,7 +557,7 @@ public void testDeleteNotSentWithExpirationTimeInPast() throws Exception { ArgumentCaptor longCaptorReplicate = ArgumentCaptor.forClass(Long.class); verify(replicationManager).sendLowPriorityReplicateContainerCommand( eq(containerInfo), eq(srcReplica.getReplicaIndex()), eq(src), - eq(tgt), longCaptorReplicate.capture()); + eq(tgt), longCaptorReplicate.capture(), any()); ContainerReplicaOp op = new ContainerReplicaOp( ADD, tgt, srcReplica.getReplicaIndex(), null, clock.millis() + 1000, 0, null); @@ -564,7 +596,7 @@ private CompletableFuture setupSuccessfulMove() moveManager.move(containerInfo.containerID(), src, tgt); verify(replicationManager).sendLowPriorityReplicateContainerCommand( - eq(containerInfo), eq(0), eq(src), eq(tgt), anyLong()); + eq(containerInfo), eq(0), eq(src), eq(tgt), anyLong(), any()); return res; } @@ -603,7 +635,7 @@ public void testMoveOverReplicatedClosedContainerWithConfigEnabled() throws Exce CompletableFuture successRes = moveManager.move(containerInfo.containerID(), src, tgt); verify(replicationManager).sendLowPriorityReplicateContainerCommand( - eq(containerInfo), eq(0), eq(src), eq(tgt), anyLong()); + eq(containerInfo), eq(0), eq(src), eq(tgt), anyLong(), any()); completeMove(containerInfo, src, tgt, successRes); } @@ -637,7 +669,7 @@ public void testMoveQuasiClosedContainerWithConfigEnabled() throws Exception { CompletableFuture successRes = moveManager.move(qcContainer.containerID(), src, tgt); verify(replicationManager).sendLowPriorityReplicateContainerCommand( - eq(qcContainer), eq(0), eq(src), eq(tgt), anyLong()); + eq(qcContainer), eq(0), eq(src), eq(tgt), anyLong(), any()); completeMove(qcContainer, src, tgt, successRes); } @@ -660,7 +692,7 @@ public void testMoveQuasiClosedContainerOverReplicatedWithConfigEnabled() throws CompletableFuture successRes = moveManager.move(qcContainer.containerID(), src, tgt); verify(replicationManager).sendLowPriorityReplicateContainerCommand( - eq(qcContainer), eq(0), eq(src), eq(tgt), anyLong()); + eq(qcContainer), eq(0), eq(src), eq(tgt), anyLong(), any()); completeMove(qcContainer, src, tgt, successRes); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java index 6bb16da585b3..e814d7518043 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java @@ -183,4 +183,102 @@ private static SCMNodeMetric createSCMNodeMetric(long capacity, long used, long singletonMap(StorageType.DEFAULT, freeSpaceToSpare), singletonMap(StorageType.DEFAULT, reserved)); } + + /** + * When a storage type is given, candidates must be ranked on their usage of + * that type. A node that is emptier overall but nearly full on the requested + * type must lose to one that has room there. + */ + @Test + public void chooseNodeRanksOnRequestedStorageType() { + OzoneConfiguration conf = new OzoneConfiguration(); + + DatanodeDetails fullOnSsd = MockDatanodeDetails.randomDatanodeDetails(); + DatanodeDetails roomOnSsd = MockDatanodeDetails.randomDatanodeDetails(); + + // fullOnSsd is emptier overall (10 of 200 used) but its SSD tier is full. + // roomOnSsd is fuller overall (100 of 200) but its SSD tier is empty. + NodeManager nodeManager = mock(NodeManager.class); + when(nodeManager.getNodeStat(fullOnSsd)).thenReturn(nodeMetric( + 100L, 95L, 100L, 5L)); + when(nodeManager.getNodeStat(roomOnSsd)).thenReturn(nodeMetric( + 100L, 5L, 100L, 95L)); + + SCMContainerPlacementCapacity policy = new SCMContainerPlacementCapacity( + nodeManager, conf, null, true, mock(SCMContainerPlacementMetrics.class)); + + // chooseNode picks two candidates at random and keeps the less used one. With + // only two nodes it sometimes draws the same index twice and returns it + // without comparing, so count outcomes over many runs rather than asserting + // on a single call. + int roomOnSsdPicked = 0; + for (int i = 0; i < 2000; i++) { + List candidates = + new ArrayList<>(Arrays.asList(fullOnSsd, roomOnSsd)); + if (roomOnSsd.equals(policy.chooseNode(candidates, StorageType.SSD))) { + roomOnSsdPicked++; + } + } + + // Whenever the two differing indices are drawn the SSD comparison must pick + // roomOnSsd, so it wins clearly more often than an even split. + assertThat(roomOnSsdPicked) + .withFailMessage("SSD-constrained choice should favour the node with " + + "free SSD capacity, but it was picked %d of 2000 times", + roomOnSsdPicked) + .isGreaterThan(1200); + } + + /** + * Without a storage type the ranking must stay on overall usage, so the node + * that is emptier overall wins even though its SSD tier is full. + */ + @Test + public void chooseNodeWithoutStorageTypeRanksOnOverallUsage() { + OzoneConfiguration conf = new OzoneConfiguration(); + + DatanodeDetails fullOnSsd = MockDatanodeDetails.randomDatanodeDetails(); + DatanodeDetails roomOnSsd = MockDatanodeDetails.randomDatanodeDetails(); + + NodeManager nodeManager = mock(NodeManager.class); + when(nodeManager.getNodeStat(fullOnSsd)).thenReturn(nodeMetric( + 100L, 95L, 100L, 5L)); + when(nodeManager.getNodeStat(roomOnSsd)).thenReturn(nodeMetric( + 100L, 5L, 100L, 95L)); + + SCMContainerPlacementCapacity policy = new SCMContainerPlacementCapacity( + nodeManager, conf, null, true, mock(SCMContainerPlacementMetrics.class)); + + int fullOnSsdPicked = 0; + for (int i = 0; i < 2000; i++) { + List candidates = + new ArrayList<>(Arrays.asList(fullOnSsd, roomOnSsd)); + if (fullOnSsd.equals(policy.chooseNode(candidates))) { + fullOnSsdPicked++; + } + } + + // Both nodes use 100 of 200 overall, so neither should dominate. + assertThat(fullOnSsdPicked).isBetween(700, 1300); + } + + /** + * Builds a metric for a node with one SSD and one DISK volume. + */ + private static SCMNodeMetric nodeMetric(long ssdCapacity, long ssdUsed, + long diskCapacity, long diskUsed) { + Map capacity = new HashMap<>(); + capacity.put(StorageType.SSD, ssdCapacity); + capacity.put(StorageType.DISK, diskCapacity); + Map used = new HashMap<>(); + used.put(StorageType.SSD, ssdUsed); + used.put(StorageType.DISK, diskUsed); + Map remaining = new HashMap<>(); + remaining.put(StorageType.SSD, ssdCapacity - ssdUsed); + remaining.put(StorageType.DISK, diskCapacity - diskUsed); + Map zeros = new HashMap<>(); + zeros.put(StorageType.SSD, 0L); + zeros.put(StorageType.DISK, 0L); + return new SCMNodeMetric(capacity, used, remaining, zeros, zeros, zeros); + } } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java index 7642e797a955..2a1958e53058 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java @@ -462,10 +462,12 @@ public static void mockRMSendThrottleReplicateCommand(ReplicationManager mock, .toTarget(containerInfo.getContainerID(), invocationOnMock.getArgument(2)); command.setReplicaIndex(invocationOnMock.getArgument(3)); + command.setTargetVolumeStorageType(invocationOnMock.getArgument(4)); commandsSent.add(Pair.of(sources.get(0), command)); return null; }).when(mock).sendThrottledReplicationCommand( - any(ContainerInfo.class), anyList(), any(DatanodeDetails.class), anyInt()); + any(ContainerInfo.class), anyList(), any(DatanodeDetails.class), anyInt(), + any()); } /** diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECMisReplicationHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECMisReplicationHandler.java index 212a1947f9b1..915652330356 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECMisReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECMisReplicationHandler.java @@ -178,7 +178,7 @@ public void testAllSourcesOverloaded() throws IOException { ReplicationManager replicationManager = getReplicationManager(); doThrow(new CommandTargetOverloadedException("Overloaded")) .when(replicationManager).sendThrottledReplicationCommand(any(), - anyList(), any(), anyInt()); + anyList(), any(), anyInt(), any()); Set availableReplicas = ReplicationTestUtil .createReplicas(Pair.of(IN_SERVICE, 1), Pair.of(IN_SERVICE, 2), Pair.of(IN_SERVICE, 3), Pair.of(IN_SERVICE, 4), diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECUnderReplicationHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECUnderReplicationHandler.java index 3e7733d6c961..1c583bde54af 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECUnderReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestECUnderReplicationHandler.java @@ -429,7 +429,7 @@ public void testUnderReplicationWithDecomNodesOverloaded() Pair.of(IN_SERVICE, 5)); doThrow(new CommandTargetOverloadedException("Overloaded")) .when(replicationManager).sendThrottledReplicationCommand( - any(), anyList(), any(), anyInt()); + any(), anyList(), any(), anyInt(), any()); assertThrows(CommandTargetOverloadedException.class, () -> testUnderReplicationWithMissingIndexes( diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisMisReplicationHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisMisReplicationHandler.java index 0d1732e2612a..5805fbd17a25 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisMisReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisMisReplicationHandler.java @@ -179,7 +179,7 @@ public void testAllSourcesOverloaded() throws IOException { ReplicationManager replicationManager = getReplicationManager(); doThrow(new CommandTargetOverloadedException("Overloaded")) .when(replicationManager).sendThrottledReplicationCommand(any(), - anyList(), any(), anyInt()); + anyList(), any(), anyInt(), any()); Set availableReplicas = ReplicationTestUtil .createReplicas(Pair.of(IN_SERVICE, 0), Pair.of(IN_SERVICE, 0), diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisUnderReplicationHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisUnderReplicationHandler.java index 404edfa1eaaf..1c28a6d3946f 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisUnderReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestRatisUnderReplicationHandler.java @@ -461,7 +461,7 @@ public void testOnlyHighestBcsidShouldBeASource() throws IOException { // Ensure that the replica with SEQ=2 is the only source sent verify(replicationManager).sendThrottledReplicationCommand(any(ContainerInfo.class), - eq(Collections.singletonList(valid.getDatanodeDetails())), any(DatanodeDetails.class), anyInt()); + eq(Collections.singletonList(valid.getDatanodeDetails())), any(DatanodeDetails.class), anyInt(), any()); } @Test diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java index 9724e7f4ff13..7f193e309e0e 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java @@ -20,7 +20,10 @@ import static java.util.Collections.singletonMap; import static org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.data.Offset.offset; +import java.util.HashMap; +import java.util.Map; import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto; @@ -74,5 +77,59 @@ void testToProtoIncludesFilesystemFieldsWhenPresent() { assertThat(proto.getFsCapacity()).isEqualTo(2000L); assertThat(proto.getFsAvailable()).isEqualTo(1500L); } + + /** + * The balancer uses this to decide which tiers a node can take part in, so it + * must report only storage types the node actually has capacity for. + */ + @Test + void testGetStorageTypesReportsOnlyTypesWithCapacity() { + DatanodeUsageInfo info = new DatanodeUsageInfo(randomDatanodeDetails(), + twoTierStat()); + + assertThat(info.getStorageTypes()) + .containsExactlyInAnyOrder(StorageType.SSD, StorageType.DISK); + } + + /** + * Utilization must be measured per storage type, otherwise the balancer cannot + * tell a node that is full on one tier from one that is full overall. + */ + @Test + void testCalculateUtilizationPerStorageType() { + DatanodeUsageInfo info = new DatanodeUsageInfo(randomDatanodeDetails(), + twoTierStat()); + + // SSD: 100 capacity, 10 remaining -> 90% used. + assertThat(info.calculateUtilization(StorageType.SSD)) + .isEqualTo(0.9, offset(0.0001)); + // DISK: 100 capacity, 80 remaining -> 20% used. + assertThat(info.calculateUtilization(StorageType.DISK)) + .isEqualTo(0.2, offset(0.0001)); + // Whole node: 200 capacity, 90 remaining -> 55% used. + assertThat(info.calculateUtilization(null)) + .isEqualTo(info.calculateUtilization()); + // A tier the node does not have reports no usage rather than failing. + assertThat(info.calculateUtilization(StorageType.ARCHIVE)).isEqualTo(0.0); + } + + /** + * A node with one unevenly used tier: SSD nearly full, DISK mostly free. + */ + private static SCMNodeStat twoTierStat() { + Map capacity = new HashMap<>(); + capacity.put(StorageType.SSD, 100L); + capacity.put(StorageType.DISK, 100L); + Map used = new HashMap<>(); + used.put(StorageType.SSD, 90L); + used.put(StorageType.DISK, 20L); + Map remaining = new HashMap<>(); + remaining.put(StorageType.SSD, 10L); + remaining.put(StorageType.DISK, 80L); + Map zeros = new HashMap<>(); + zeros.put(StorageType.SSD, 0L); + zeros.put(StorageType.DISK, 0L); + return new SCMNodeStat(capacity, used, remaining, zeros, zeros, zeros); + } } From 966d2934552f9a872e4414aaf4d2382a25b514b1 Mon Sep 17 00:00:00 2001 From: Devesh Singh Date: Sat, 3 Oct 2026 19:55:59 +0530 Subject: [PATCH 2/2] HDDS-16653. Findbugs and PMD errors. --- .../container/replication/ContainerImporter.java | 12 ++++++------ .../container/replication/TestContainerImporter.java | 8 ++++---- .../replication/RatisUnderReplicationHandler.java | 7 ------- 3 files changed, 10 insertions(+), 17 deletions(-) diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java index de1144f3083a..074d4bc951f4 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ContainerImporter.java @@ -26,6 +26,7 @@ import java.nio.file.Paths; import java.util.Collections; import java.util.HashSet; +import java.util.Objects; import java.util.Set; import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.conf.ConfigurationSource; @@ -92,6 +93,7 @@ public boolean isAllowedContainerImport(long containerID) { public void importContainer(long containerID, Path tarFilePath, HddsVolume targetVolume, CopyContainerCompression compression) throws IOException { + Objects.requireNonNull(targetVolume, "targetVolume == null"); if (!importContainerProgress.add(containerID)) { deleteFileQuietely(tarFilePath); String log = "Container import in progress with container Id " + containerID; @@ -118,12 +120,10 @@ public void importContainer(long containerID, Path tarFilePath, } ContainerUtils.verifyContainerFileChecksum(containerData, conf); containerData.setVolume(targetVolume); - if (targetVolume != null) { - // The descriptor carries the source volume's storage type. Record the - // type of the volume actually chosen here, so the replica reports where - // it really lives rather than where its source lived. - containerData.setStorageType(targetVolume.getStorageType()); - } + // The descriptor carries the source volume's storage type. Record the + // type of the volume actually chosen here, so the replica reports where + // it really lives rather than where its source lived. + containerData.setStorageType(targetVolume.getStorageType()); // lastDataScanTime should be cleared for an imported container containerData.setDataScanTimestamp(null); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java index c299163b1bb9..e7e3b0d62933 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerImporter.java @@ -128,7 +128,7 @@ void importSameContainerWhenAlreadyImport() throws Exception { // second import should fail immediately StorageContainerException ex = assertThrows(StorageContainerException.class, () -> containerImporter.importContainer(containerId, tarFile.toPath(), - null, NO_COMPRESSION)); + mock(HddsVolume.class), NO_COMPRESSION)); assertEquals(ContainerProtos.Result.CONTAINER_EXISTS, ex.getResult()); assertThat(ex.getMessage()).contains("Container already exists"); } @@ -147,7 +147,7 @@ void importSameContainerWhenFirstInProgress() throws Exception { CompletableFuture.runAsync(() -> { try { containerImporter.importContainer(containerId, tarFile.toPath(), - null, NO_COMPRESSION); + mock(HddsVolume.class), NO_COMPRESSION); } catch (Exception ex) { // do nothing } @@ -158,7 +158,7 @@ void importSameContainerWhenFirstInProgress() throws Exception { StorageContainerException ex = assertThrows( StorageContainerException.class, () -> containerImporter.importContainer(containerId, tarFile.toPath(), - null, NO_COMPRESSION)); + mock(HddsVolume.class), NO_COMPRESSION)); assertEquals(ContainerProtos.Result.CONTAINER_EXISTS, ex.getResult()); assertThat(ex.getMessage()).contains("import in progress"); @@ -190,7 +190,7 @@ public void testInconsistentChecksumContainerShouldThrowError() throws Exception StorageContainerException scException = assertThrows(StorageContainerException.class, () -> importer.importContainer(containerId, - tarFile.toPath(), null, NO_COMPRESSION)); + tarFile.toPath(), mock(HddsVolume.class), NO_COMPRESSION)); Assertions.assertTrue(scException.getMessage(). contains("Container checksum error")); } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java index f808dd19736b..59ee00cef691 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/RatisUnderReplicationHandler.java @@ -471,13 +471,6 @@ private List getTargets( currentContainerSize, replicaCount.getContainer(), StorageType.DEFAULT); } - private int sendReplicationCommands( - ContainerInfo containerInfo, List sources, - List targets) throws CommandTargetOverloadedException, - NotLeaderException { - return sendReplicationCommands(containerInfo, sources, targets, null); - } - /** * All Ratis replicas of a container are interchangeable, so the tier of any * reported replica tells us where new copies belong.