This is an automated email from the ASF dual-hosted git repository. devmadhuu pushed a commit to branch HDDS-11233 in repository https://gitbox.apache.org/repos/asf/ozone.git
commit 5a16f67a1c64bb45adcdf5f420bd86d713b01fe8 Author: Devesh Kumar Singh <[email protected]> AuthorDate: Tue Aug 4 07:14:48 2026 +0530 HDDS-15408. SCM ContainerInfo Support Create And Get With StorageTier (#10810) Co-authored-by: XiChen <[email protected]> --- .../org/apache/hadoop/hdds/client/StorageTier.java | 7 + .../hadoop/hdds/scm/container/ContainerInfo.java | 24 ++ .../hdds/scm/container/TestContainerInfo.java | 13 ++ .../protocol/StorageContainerLocationProtocol.java | 25 +- ...inerLocationProtocolClientSideTranslatorPB.java | 10 +- .../src/main/proto/ScmAdminProtocol.proto | 1 + .../interface-client/src/main/proto/hdds.proto | 1 + .../hdds/scm/container/ContainerManager.java | 13 +- .../hdds/scm/container/ContainerManagerImpl.java | 71 +++--- .../hdds/scm/container/ContainerStateManager.java | 6 +- .../scm/container/ContainerStateManagerImpl.java | 23 +- .../ha/invoker/ContainerStateManagerInvoker.java | 19 +- .../hdds/scm/pipeline/PipelinePlacementPolicy.java | 27 +-- .../hadoop/hdds/scm/pipeline/PipelineStateMap.java | 9 +- .../hdds/scm/pipeline/SimplePipelineProvider.java | 2 +- .../scm/pipeline/WritableECContainerProvider.java | 12 +- .../pipeline/WritableRatisContainerProvider.java | 3 +- ...inerLocationProtocolServerSideTranslatorPB.java | 8 +- .../hdds/scm/server/SCMClientProtocolServer.java | 12 +- .../hdds/scm/server/StorageContainerManager.java | 14 ++ .../org/apache/hadoop/hdds/scm/HddsTestUtils.java | 3 +- .../scm/container/TestContainerManagerImpl.java | 37 +-- .../scm/container/TestContainerStateManager.java | 48 +++- .../hdds/scm/node/TestContainerPlacement.java | 2 +- .../hdds/scm/node/TestNodeDecommissionManager.java | 47 ++-- .../hdds/scm/pipeline/MockPipelineManager.java | 5 +- .../hdds/scm/pipeline/TestPipelineManagerImpl.java | 2 +- .../scm/pipeline/TestPipelinePlacementPolicy.java | 23 ++ .../hdds/scm/pipeline/TestPipelineStateMap.java | 31 +++ .../pipeline/TestWritableECContainerProvider.java | 2 +- .../TestWritableRatisContainerProvider.java | 3 +- .../server/TestStorageContainerManagerStarter.java | 11 + .../hadoop/ozone/recon/TestReconAsPassiveScm.java | 6 +- .../hadoop/ozone/recon/TestReconScmSnapshot.java | 4 +- .../apache/hadoop/ozone/recon/TestReconTasks.java | 15 +- .../hadoop/hdds/scm/TestSCMInstallSnapshot.java | 5 +- .../hdds/scm/TestSCMInstallSnapshotWithHA.java | 4 +- .../org/apache/hadoop/hdds/scm/TestSCMMXBean.java | 3 +- .../apache/hadoop/hdds/scm/TestSCMSnapshot.java | 6 +- .../hadoop/hdds/scm/TestSecretKeySnapshot.java | 3 +- .../TestAllocateContainerWithStorageTier.java | 254 +++++++++++++++++++++ .../TestContainerStateManagerIntegration.java | 14 +- .../container/TestScmApplyTransactionFailure.java | 4 +- .../metrics/TestSCMContainerManagerMetrics.java | 7 +- .../TestReplicationManagerIntegration.java | 9 +- .../hdds/scm/pipeline/TestNode2PipelineMap.java | 3 +- .../hdds/scm/pipeline/TestPipelineClose.java | 7 +- .../hadoop/hdds/scm/pipeline/TestSCMRestart.java | 9 +- .../hdds/scm/storage/TestContainerCommandsEC.java | 7 +- .../hadoop/hdds/upgrade/TestHDDSUpgrade.java | 4 +- .../apache/hadoop/ozone/TestMiniOzoneCluster.java | 15 ++ .../org/apache/hadoop/ozone/MiniOzoneCluster.java | 35 +++ .../hadoop/ozone/UniformDatanodesFactory.java | 49 +++- 53 files changed, 803 insertions(+), 174 deletions(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTier.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTier.java index 24b783f1b14..9e4ab78b767 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTier.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTier.java @@ -68,6 +68,13 @@ public static StorageTier getDefaultTier() { return defaultTier; } + public static void setDefault(StorageTier tier) { + if (tier == null) { + throw new IllegalArgumentException("Default StorageTier cannot be null."); + } + defaultTier = tier; + } + public StorageTierProto toProto() { switch (this) { case SSD: diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerInfo.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerInfo.java index 29f937bc0e6..3b2c8ccf1a8 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerInfo.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerInfo.java @@ -22,6 +22,7 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.annotation.JsonProperty; +import jakarta.annotation.Nullable; import java.time.Clock; import java.time.Instant; import java.util.Comparator; @@ -29,6 +30,7 @@ import org.apache.commons.lang3.builder.HashCodeBuilder; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.pipeline.PipelineID; import org.apache.hadoop.hdds.utils.db.Codec; @@ -86,6 +88,7 @@ public final class ContainerInfo implements Comparable<ContainerInfo> { // Health state of the container (determined by ReplicationManager) private ContainerHealthState healthState; private boolean suppressed; + private final StorageTier storageTier; private ContainerInfo(Builder b) { containerID = ContainerID.valueOf(b.containerID); @@ -102,6 +105,7 @@ private ContainerInfo(Builder b) { clock = b.clock; healthState = b.healthState != null ? b.healthState : ContainerHealthState.HEALTHY; suppressed = b.suppressed; + storageTier = b.storageTier; } public static Codec<ContainerInfo> getCodec() { @@ -127,6 +131,10 @@ public static ContainerInfo fromProtobuf(HddsProtos.ContainerInfoProto info) { builder.setSuppressed(info.getSuppressed()); } + if (info.hasStorageTier()) { + builder.setStorageTier(StorageTier.fromProto(info.getStorageTier())); + } + if (info.hasPipelineID()) { builder.setPipelineID(PipelineID.getFromProtobuf(info.getPipelineID())); } @@ -290,6 +298,11 @@ public void setSuppressed(boolean suppressed) { this.suppressed = suppressed; } + @Nullable + public StorageTier getStorageTier() { + return storageTier; + } + @JsonIgnore public HddsProtos.ContainerInfoProto getProtobuf() { HddsProtos.ContainerInfoProto.Builder builder = @@ -319,6 +332,10 @@ public HddsProtos.ContainerInfoProto getProtobuf() { builder.setSuppressed(true); } + if (storageTier != null && storageTier != StorageTier.EMPTY) { + builder.setStorageTier(storageTier.toProto()); + } + return builder.build(); } @@ -338,6 +355,7 @@ public String toString() { + ", stateEnterTime=" + stateEnterTime + ", pipelineID=" + pipelineID + ", owner=" + owner + + ", storageTier=" + storageTier + '}'; } @@ -422,6 +440,7 @@ public static class Builder { private ReplicationConfig replicationConfig; private ContainerHealthState healthState; private boolean suppressed; + private StorageTier storageTier; public Builder setPipelineID(PipelineID pipelineId) { this.pipelineID = pipelineId; @@ -485,6 +504,11 @@ public Builder setSuppressed(boolean suppressed) { return this; } + public Builder setStorageTier(StorageTier storageTier) { + this.storageTier = storageTier; + return this; + } + /** * Also resets {@code stateEnterTime}, so make sure to set clock first. */ diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerInfo.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerInfo.java index d7c14fd00c8..84e39eddfd2 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerInfo.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerInfo.java @@ -30,8 +30,10 @@ import java.time.Instant; import java.util.concurrent.ThreadLocalRandom; import org.apache.commons.lang3.builder.HashCodeBuilder; +import org.apache.hadoop.hdds.JsonTestUtils; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.pipeline.PipelineID; import org.apache.ozone.test.MockClock; @@ -131,6 +133,17 @@ void restoreState() { assertThrows(IllegalStateException.class, subject::revertState); } + @Test + void storageTierIsIncludedInJson() { + ContainerInfo container = newBuilderForTest() + .setReplicationConfig(RatisReplicationConfig.getInstance(THREE)) + .setStorageTier(StorageTier.ARCHIVE) + .build(); + + assertEquals("ARCHIVE", JsonTestUtils.valueToJsonNode(container) + .get("storageTier").asText()); + } + public static ContainerInfo.Builder newBuilderForTest() { return new ContainerInfo.Builder() .setContainerID(1234) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java index 98f8efa9ae3..24d8d34ee04 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java @@ -29,6 +29,7 @@ import java.util.UUID; import org.apache.commons.lang3.tuple.Pair; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DeletedBlocksTransactionInfo; @@ -86,12 +87,32 @@ public interface StorageContainerLocationProtocol extends Closeable { * set of datanodes that should be used creating this container. * */ - ContainerWithPipeline allocateContainer( + default ContainerWithPipeline allocateContainer( HddsProtos.ReplicationType replicationType, HddsProtos.ReplicationFactor factor, String owner) + throws IOException { + return allocateContainer(replicationType, factor, owner, + StorageTier.getDefaultTier().toProto()); + } + + /** + * Asks SCM where a container should be allocated. SCM responds with the + * set of datanodes that should be used creating this container. + * + */ + ContainerWithPipeline allocateContainer( + HddsProtos.ReplicationType replicationType, + HddsProtos.ReplicationFactor factor, String owner, + HddsProtos.StorageTierProto storageTier) throws IOException; - ContainerWithPipeline allocateContainer(ReplicationConfig replicationConfig, String owner) throws IOException; + default ContainerWithPipeline allocateContainer(ReplicationConfig replicationConfig, String owner) + throws IOException { + return allocateContainer(replicationConfig, owner, StorageTier.getDefaultTier().toProto()); + } + + ContainerWithPipeline allocateContainer(ReplicationConfig replicationConfig, String owner, + HddsProtos.StorageTierProto storageTier) throws IOException; /** * Ask SCM the location of the container. SCM responds with a group of diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java index 7808cb286a2..5209dbadcea 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java @@ -259,20 +259,22 @@ private ScmContainerLocationResponse submitRpcRequest( @Override public ContainerWithPipeline allocateContainer( HddsProtos.ReplicationType type, HddsProtos.ReplicationFactor factor, - String owner) throws IOException { + String owner, HddsProtos.StorageTierProto storageTier) throws IOException { ReplicationConfig replicationConfig = ReplicationConfig.fromProtoTypeAndFactor(type, factor); - return allocateContainer(replicationConfig, owner); + return allocateContainer(replicationConfig, owner, storageTier); } @Override public ContainerWithPipeline allocateContainer( - ReplicationConfig replicationConfig, String owner) throws IOException { + ReplicationConfig replicationConfig, String owner, + HddsProtos.StorageTierProto storageTier) throws IOException { ContainerRequestProto.Builder request = ContainerRequestProto.newBuilder() .setTraceID(TracingUtil.exportCurrentSpan()) .setReplicationType(replicationConfig.getReplicationType()) - .setOwner(owner); + .setOwner(owner) + .setStorageTier(storageTier); if (replicationConfig.getReplicationType() == HddsProtos.ReplicationType.EC) { HddsProtos.ECReplicationConfig ecProto = diff --git a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto index 6c39fc22cc4..1598897c85c 100644 --- a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto +++ b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto @@ -220,6 +220,7 @@ message ContainerRequestProto { required string owner = 4; optional string traceID = 5; optional ECReplicationConfig ecReplicationConfig = 6; + optional StorageTierProto storageTier = 7; } /** diff --git a/hadoop-hdds/interface-client/src/main/proto/hdds.proto b/hadoop-hdds/interface-client/src/main/proto/hdds.proto index f14c888c644..b396af1966f 100644 --- a/hadoop-hdds/interface-client/src/main/proto/hdds.proto +++ b/hadoop-hdds/interface-client/src/main/proto/hdds.proto @@ -276,6 +276,7 @@ message ContainerInfoProto { required ReplicationType replicationType = 11; optional ECReplicationConfig ecReplicationConfig = 12; optional bool suppressed = 13; + optional StorageTierProto storageTier = 14; } message ContainerWithPipeline { diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManager.java index e7d2070d8fa..21707b04044 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManager.java @@ -19,10 +19,10 @@ import jakarta.annotation.Nullable; import java.io.IOException; -import java.util.Collections; import java.util.List; import java.util.Set; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ContainerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleEvent; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState; @@ -155,7 +155,7 @@ default long getTotalContainerCount() { * @throws IOException */ ContainerInfo allocateContainer(ReplicationConfig replicationConfig, - String owner) + String owner, StorageTier storageTier) throws IOException; /** @@ -204,24 +204,21 @@ void updateContainerReplica(ContainerID containerID, ContainerReplica replica) void removeContainerReplica(ContainerID containerID, ContainerReplica replica) throws ContainerNotFoundException, ContainerReplicaNotFoundException; - default ContainerInfo getMatchingContainer(long size, String owner, - Pipeline pipeline) { - return getMatchingContainer(size, owner, pipeline, Collections.emptySet()); - } - /** * Returns ContainerInfo which matches the requirements. * @param size - the amount of space required in the container * @param owner - the user which requires space in its owned container * @param pipeline - pipeline to which the container should belong. * @param excludedContainerIDS - containerIds to be excluded. + * @param storageTier - Required storageTier for Container. * @return ContainerInfo for the matching container, or null if a container could not be found and could not be * allocated */ @Nullable ContainerInfo getMatchingContainer(long size, String owner, Pipeline pipeline, - Set<ContainerID> excludedContainerIDS); + Set<ContainerID> excludedContainerIDS, + StorageTier storageTier); /** * Once after report processor handler completes, call this to notify diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java index ad965640376..588c979c052 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java @@ -23,6 +23,7 @@ import java.util.Iterator; import java.util.List; import java.util.NavigableSet; +import java.util.Objects; import java.util.Random; import java.util.Set; import java.util.concurrent.locks.Lock; @@ -169,7 +170,8 @@ public int getContainerStateCount(final LifeCycleState state) { @Override public ContainerInfo allocateContainer( - final ReplicationConfig replicationConfig, final String owner) + final ReplicationConfig replicationConfig, final String owner, + StorageTier storageTier) throws IOException { // Acquire pipeline manager lock, to avoid any updates to pipeline // while allocate container happens. This is to avoid scenario like @@ -181,10 +183,10 @@ public ContainerInfo allocateContainer( ContainerInfo containerInfo = null; try { pipelines = pipelineManager - .getPipelines(replicationConfig, Pipeline.PipelineState.OPEN); + .getPipelines(replicationConfig, Pipeline.PipelineState.OPEN, storageTier); if (!pipelines.isEmpty()) { pipeline = pipelines.get(random.nextInt(pipelines.size())); - containerInfo = createContainer(pipeline, owner); + containerInfo = createContainer(pipeline, owner, storageTier); } } finally { lock.unlock(); @@ -193,8 +195,7 @@ public ContainerInfo allocateContainer( if (pipelines.isEmpty()) { try { - pipeline = pipelineManager.createPipeline(replicationConfig, - StorageTier.getDefaultTier()); + pipeline = pipelineManager.createPipeline(replicationConfig, storageTier); if (replicationConfig.getReplicationType() == HddsProtos.ReplicationType.EC) { pipelineManager.openPipeline(pipeline.getId()); } @@ -203,20 +204,20 @@ public ContainerInfo allocateContainer( scmContainerManagerMetrics.incNumFailureCreateContainers(); throw new IOException("Could not allocate container. Cannot get any" + " matching pipeline for replicationConfig: " + replicationConfig - + ", State:PipelineState.OPEN", e); + + ", State:PipelineState.OPEN, storageTier " + storageTier, e); } pipelineManager.acquireReadLock(); lock.lock(); try { pipelines = pipelineManager - .getPipelines(replicationConfig, Pipeline.PipelineState.OPEN); + .getPipelines(replicationConfig, Pipeline.PipelineState.OPEN, storageTier); if (!pipelines.isEmpty()) { pipeline = pipelines.get(random.nextInt(pipelines.size())); - containerInfo = createContainer(pipeline, owner); + containerInfo = createContainer(pipeline, owner, storageTier); } else { throw new IOException("Could not allocate container. Cannot get any" + " matching pipeline for replicationConfig: " + replicationConfig - + ", State:PipelineState.OPEN"); + + ", State:PipelineState.OPEN, storageTier " + storageTier); } } finally { lock.unlock(); @@ -226,9 +227,10 @@ public ContainerInfo allocateContainer( return containerInfo; } - private ContainerInfo createContainer(Pipeline pipeline, String owner) + private ContainerInfo createContainer(Pipeline pipeline, String owner, + StorageTier storageTier) throws IOException { - final ContainerInfo containerInfo = allocateContainer(pipeline, owner); + final ContainerInfo containerInfo = allocateContainer(pipeline, owner, storageTier); if (LOG.isTraceEnabled()) { LOG.trace("New container allocated: {}", containerInfo); } @@ -236,12 +238,15 @@ private ContainerInfo createContainer(Pipeline pipeline, String owner) } private ContainerInfo allocateContainer(final Pipeline pipeline, - final String owner) + final String owner, + StorageTier storageTier) throws IOException { final long uniqueId = sequenceIdGen.getNextId(SequenceIdType.containerId); Preconditions.checkState(uniqueId > 0, "Cannot allocate container, negative container id" + " generated. %s.", uniqueId); + Objects.requireNonNull(storageTier, + "Cannot allocate container, StorageTier cannot be null."); final ContainerID containerID = ContainerID.valueOf(uniqueId); if (!pipelineManager.checkSpaceAndRecordAllocation(pipeline, containerID)) { @@ -258,7 +263,8 @@ private ContainerInfo allocateContainer(final Pipeline pipeline, .setStateEnterTime(Time.now()) .setOwner(owner) .setContainerID(containerID.getId()) - .setReplicationType(pipeline.getType()); + .setReplicationType(pipeline.getType()) + .setStorageTier(storageTier.toProto()); if (pipeline.getReplicationConfig() instanceof ECReplicationConfig) { containerInfoBuilder.setEcReplicationConfig( @@ -358,24 +364,25 @@ public void removeContainerReplica(final ContainerID cid, @Override public ContainerInfo getMatchingContainer(final long size, final String owner, - final Pipeline pipeline, final Set<ContainerID> excludedContainerIDs) { + final Pipeline pipeline, final Set<ContainerID> excludedContainerIDs, + StorageTier storageTier) { NavigableSet<ContainerID> containerIDs; ContainerInfo containerInfo; try { synchronized (pipeline.getId()) { - containerIDs = getContainersForOwner(pipeline, owner); + containerIDs = getContainersForOwnerAndStorageTier(pipeline, owner, storageTier); if (containerIDs.size() < pipelineManager.openContainerLimit(pipeline.getNodes())) { - ContainerInfo allocated = allocateContainer(pipeline, owner); + ContainerInfo allocated = allocateContainer(pipeline, owner, storageTier); if (allocated != null) { // New container was created, refresh IDs so it becomes eligible. - containerIDs = getContainersForOwner(pipeline, owner); + containerIDs = getContainersForOwnerAndStorageTier(pipeline, owner, storageTier); } } containerIDs.removeAll(excludedContainerIDs); - containerInfo = containerStateManager.getMatchingContainer( - size, owner, pipeline.getId(), containerIDs); + containerInfo = containerStateManager.getMatchingContainerAndStorageTier( + size, owner, pipeline.getId(), containerIDs, storageTier); if (containerInfo == null) { - containerInfo = allocateContainer(pipeline, owner); + containerInfo = allocateContainer(pipeline, owner, storageTier); } return containerInfo; } @@ -386,20 +393,32 @@ public ContainerInfo getMatchingContainer(final long size, final String owner, } /** - * Returns the container ID's matching with specified owner. - * @param pipeline - * @param owner + * Returns the container ID's matching with specified owner and storage tier. + * A stored container with a null storage tier is treated as matching any + * tier (upgrade-compat with older containers that predate the storageTier + * field). The requested {@code storageTier} argument itself must not be + * null; callers must supply a concrete tier (defaulting to + * {@link org.apache.hadoop.hdds.client.StorageTier#getDefaultTier()} if + * the client did not request one). + * @param pipeline pipeline + * @param owner owner + * @param storageTier requested storageTier, must not be null * @return NavigableSet<ContainerID> */ - private NavigableSet<ContainerID> getContainersForOwner( - Pipeline pipeline, String owner) throws IOException { + private NavigableSet<ContainerID> getContainersForOwnerAndStorageTier( + Pipeline pipeline, String owner, StorageTier storageTier) throws IOException { + Objects.requireNonNull(storageTier, + "storageTier is required for container matching"); NavigableSet<ContainerID> containerIDs = pipelineManager.getContainersInPipeline(pipeline.getId()); Iterator<ContainerID> containerIDIterator = containerIDs.iterator(); while (containerIDIterator.hasNext()) { ContainerID cid = containerIDIterator.next(); try { - if (!getContainer(cid).getOwner().equals(owner)) { + ContainerInfo info = getContainer(cid); + if (!info.getOwner().equals(owner) || + (info.getStorageTier() != null && + !info.getStorageTier().equals(storageTier))) { containerIDIterator.remove(); } } catch (ContainerNotFoundException e) { diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManager.java index 913e1ce84a4..d6c3790eb7d 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManager.java @@ -21,6 +21,7 @@ import java.util.List; import java.util.NavigableSet; import java.util.Set; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ContainerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState; @@ -195,9 +196,10 @@ void transitionDeletingOrDeletedToTargetState(HddsProtos.ContainerID id, LifeCyc /** * */ - ContainerInfo getMatchingContainer(long size, String owner, + ContainerInfo getMatchingContainerAndStorageTier(long size, String owner, PipelineID pipelineID, - NavigableSet<ContainerID> containerIDs); + NavigableSet<ContainerID> containerIDs, + StorageTier storageTier); /** * diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManagerImpl.java index 65060c0a171..784d3a2f40f 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerStateManagerImpl.java @@ -33,6 +33,7 @@ import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CONTAINER_LOCK_STRIPE_SIZE_DEFAULT; import com.google.common.util.concurrent.Striped; +import jakarta.annotation.Nonnull; import java.io.IOException; import java.util.EnumMap; import java.util.HashSet; @@ -46,6 +47,7 @@ import java.util.concurrent.locks.ReentrantReadWriteLock; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.conf.StorageUnit; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ContainerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleEvent; @@ -490,8 +492,9 @@ public void removeContainerReplica(final ContainerReplica replica) { } @Override - public ContainerInfo getMatchingContainer(final long size, String owner, - PipelineID pipelineID, NavigableSet<ContainerID> containerIDs) { + public ContainerInfo getMatchingContainerAndStorageTier(final long size, String owner, + PipelineID pipelineID, NavigableSet<ContainerID> containerIDs, + StorageTier storageTier) { if (containerIDs.isEmpty()) { return null; } @@ -509,7 +512,8 @@ public ContainerInfo getMatchingContainer(final long size, String owner, if (resultSet.isEmpty()) { resultSet = containerIDs; } - ContainerInfo selectedContainer = findContainerWithSpace(size, resultSet); + ContainerInfo selectedContainer = + findContainerWithSpaceAndStorageTier(size, resultSet, storageTier); if (selectedContainer == null) { // If we did not find any space in the tailSet, we need to look for @@ -521,7 +525,7 @@ public ContainerInfo getMatchingContainer(final long size, String owner, // last element in the sorted set. resultSet = containerIDs.headSet(lastID, true); - selectedContainer = findContainerWithSpace(size, resultSet); + selectedContainer = findContainerWithSpaceAndStorageTier(size, resultSet, storageTier); } // TODO: cleanup entries in lastUsedMap @@ -531,14 +535,15 @@ public ContainerInfo getMatchingContainer(final long size, String owner, return selectedContainer; } - private ContainerInfo findContainerWithSpace(final long size, - final NavigableSet<ContainerID> - searchSet) { - // Get the container with space to meet our request. + private ContainerInfo findContainerWithSpaceAndStorageTier(final long size, + final NavigableSet<ContainerID> searchSet, @Nonnull StorageTier storageTier) { + // Get the container with space to meet our request and exact tier. for (ContainerID id : searchSet) { try (AutoCloseableLock ignored = readLock(id)) { final ContainerInfo containerInfo = containers.getContainerInfo(id); - if (containerInfo.getUsedBytes() + size <= this.containerSize) { + if (containerInfo.getUsedBytes() + size <= this.containerSize && + containerInfo.getStorageTier() != null && + containerInfo.getStorageTier().equals(storageTier)) { containerInfo.updateLastUsedTime(); return containerInfo; } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/ContainerStateManagerInvoker.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/ContainerStateManagerInvoker.java index d5c7be1c8fe..4e719fd5b90 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/ContainerStateManagerInvoker.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/ContainerStateManagerInvoker.java @@ -21,6 +21,7 @@ import java.util.List; import java.util.NavigableSet; import java.util.Set; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ContainerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleEvent; @@ -141,8 +142,9 @@ public Set<ContainerReplica> getContainerReplicas(ContainerID arg0) { } @Override - public ContainerInfo getMatchingContainer(long arg0, String arg1, PipelineID arg2, NavigableSet arg3) { - return invoker.getImpl().getMatchingContainer(arg0, arg1, arg2, arg3); + public ContainerInfo getMatchingContainerAndStorageTier(long arg0, String arg1, PipelineID arg2, NavigableSet + arg3, StorageTier arg4) { + return invoker.getImpl().getMatchingContainerAndStorageTier(arg0, arg1, arg2, arg3, arg4); } @Override @@ -263,13 +265,14 @@ public Message invokeLocal(String methodName, Object[] p) throws Exception { returnValue = getImpl().getContainerReplicas(arg15); break; - case "getMatchingContainer": - final long arg16 = p.length > 0 ? (long) p[0] : 0L; - final String arg17 = p.length > 1 ? (String) p[1] : null; - final PipelineID arg18 = p.length > 2 ? (PipelineID) p[2] : null; - final NavigableSet arg19 = p.length > 3 ? (NavigableSet) p[3] : null; + case "getMatchingContainerAndStorageTier": + final long arg15 = p.length > 0 ? (long) p[0] : 0L; + final String arg16 = p.length > 1 ? (String) p[1] : null; + final PipelineID arg17 = p.length > 2 ? (PipelineID) p[2] : null; + final NavigableSet arg18 = p.length > 3 ? (NavigableSet) p[3] : null; + final StorageTier arg19 = p.length > 4 ? (StorageTier) p[4] : null; returnType = ContainerInfo.class; - returnValue = getImpl().getMatchingContainer(arg16, arg17, arg18, arg19); + returnValue = getImpl().getMatchingContainerAndStorageTier(arg15, arg16, arg17, arg18, arg19); break; case "reinitialize": 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 8dbdafea23b..f18b0346a96 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 @@ -115,20 +115,21 @@ List<DatanodeDetails> filterPipelineLimit(Iterable<DatanodeDetails> datanodes, S } private static boolean isNonClosedRatisThreePipeline(Pipeline p, StorageType storageType) { - boolean matchedTier = false; - if (storageType != null) { - try { - // TODO: Do we need to return true if getSupportedStorageTier() is empty - // otherwise this might cause more pipelines to be created than necessary - matchedTier = p.getSupportedStorageTier().getUniformStorageType().equals(storageType); - } catch (IllegalArgumentException e) { - LOG.debug("Cannot convert pipeline storage tier {} to storage type.", p.getSupportedStorageTier(), e); - return false; - } + if (p == null || storageType == null || + p.getSupportedStorageTier() == null || + !p.getReplicationConfig().equals( + RatisReplicationConfig.getInstance(ReplicationFactor.THREE)) || + p.isClosed()) { + return false; + } + try { + return p.getSupportedStorageTier().getUniformStorageType() + .equals(storageType); + } catch (IllegalArgumentException e) { + LOG.debug("Cannot convert pipeline storage tier {} to storage type.", + p.getSupportedStorageTier(), e); + return false; } - return p != null && p.getReplicationConfig() - .equals(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)) - && !p.isClosed() && matchedTier; } @Override diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java index 82a34c54e16..cc22b3fdeed 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java @@ -233,15 +233,12 @@ List<Pipeline> getPipelines(ReplicationConfig replicationConfig, } /** - * A pipeline with no supportedStorageTier (e.g. one created by an older - * SCM before this field was introduced, or restored from a legacy DB - * snapshot) is treated as matching any tier. This preserves backward - * compatibility during rolling upgrade: legacy pipelines remain usable - * until they are scrubbed and replaced by tier-tagged ones. + * A pipeline matches only when it explicitly supports the requested tier. + * An untyped legacy pipeline cannot safely satisfy an explicit tier request. */ static boolean matchesStorageTier(Pipeline pipeline, StorageTier storageTier) { final StorageTier pipelineTier = pipeline.getSupportedStorageTier(); - return pipelineTier == null || Objects.equals(pipelineTier, storageTier); + return Objects.equals(pipelineTier, storageTier); } /** diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/SimplePipelineProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/SimplePipelineProvider.java index 9e4d98feb8f..85f6f0d1ddd 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/SimplePipelineProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/SimplePipelineProvider.java @@ -60,7 +60,7 @@ public Pipeline create(StandaloneReplicationConfig replicationConfig, List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, StorageTier storageTier) throws IOException { StorageType storageType = storageTier.getUniformStorageType(); - List<DatanodeDetails> dns = pickNodesNotUsed(replicationConfig); + List<DatanodeDetails> dns = pickAllNodesNotUsed(replicationConfig); dns = dns.stream().filter(dn -> ((DatanodeInfo) dn).getStorageReports().stream() .anyMatch(reportProto -> diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java index 8847a950ac3..91574102dd2 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java @@ -104,7 +104,7 @@ public ContainerInfo getContainer(final long size, Pipeline.PipelineState.OPEN); if (openPipelineCount < maximumPipelines) { try { - return allocateContainer(repConfig, size, owner, excludeList); + return allocateContainer(repConfig, size, owner, excludeList, storageTier); } catch (IOException e) { LOG.warn("Unable to allocate a container with {} existing ones; " + "requested size={}, replication={}, owner={}, {}", @@ -173,7 +173,7 @@ public ContainerInfo getContainer(final long size, } if (openPipelineCount < maximumPipelines) { synchronized (this) { - return allocateContainer(repConfig, size, owner, excludeList); + return allocateContainer(repConfig, size, owner, excludeList, storageTier); } } throw new IOException("Pipeline limit (" + maximumPipelines @@ -197,7 +197,8 @@ private int getMaximumPipelines(ECReplicationConfig repConfig) { } private ContainerInfo allocateContainer(ReplicationConfig repConfig, - long size, String owner, ExcludeList excludeList) + long size, String owner, ExcludeList excludeList, + @Nonnull StorageTier storageTier) throws IOException { List<DatanodeDetails> excludedNodes = Collections.emptyList(); @@ -205,13 +206,12 @@ private ContainerInfo allocateContainer(ReplicationConfig repConfig, excludedNodes = new ArrayList<>(excludeList.getDatanodes()); } - // TODO StoragePolicy Support EC Pipeline newPipeline = pipelineManager.createPipeline(repConfig, - excludedNodes, Collections.emptyList(), StorageTier.getDefaultTier()); + excludedNodes, Collections.emptyList(), storageTier); // the returned ContainerInfo should not be null (due to not enough space in the Datanodes specifically) because // this is a new pipeline and pipeline creation checks for sufficient space in the Datanodes ContainerInfo container = - containerManager.getMatchingContainer(size, owner, newPipeline); + containerManager.getMatchingContainer(size, owner, newPipeline, Collections.emptySet(), storageTier); if (container == null) { // defensive null handling throw new IOException("Could not allocate a new container"); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java index edfddfc2b81..1c935449777 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java @@ -190,9 +190,8 @@ private List<Pipeline> findPipelinesByState( availablePipelines, req); // look for OPEN containers that match the criteria. - // TODO StoragePolicy replace this StorageType with container actual StorageType final ContainerInfo containerInfo = containerManager.getMatchingContainer( - req.getSize(), owner, pipeline, excludeList.getContainerIds()); + req.getSize(), owner, pipeline, excludeList.getContainerIds(), storageTier); if (containerInfo != null) { return containerInfo; diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java index 00a9b6b3a0c..a05001f6c0b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java @@ -42,6 +42,7 @@ import org.apache.hadoop.hdds.annotation.InterfaceAudience; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.TransferLeadershipRequestProto; @@ -790,11 +791,14 @@ public GetContainerReplicasResponseProto getContainerReplicas( public ContainerResponseProto allocateContainer(ContainerRequestProto request, int clientVersion) throws IOException { - ReplicationConfig replicationConfig = ReplicationConfig.fromProto(request.getReplicationType(), + ReplicationConfig replicationConfig = ReplicationConfig.fromProto(request.getReplicationType(), request.getReplicationFactor(), request.getEcReplicationConfig() ); - ContainerWithPipeline cp = impl.allocateContainer(replicationConfig, request.getOwner()); + HddsProtos.StorageTierProto storageTier = request.hasStorageTier() + ? request.getStorageTier() + : StorageTier.getDefaultTier().toProto(); + ContainerWithPipeline cp = impl.allocateContainer(replicationConfig, request.getOwner(), storageTier); return ContainerResponseProto.newBuilder() .setContainerWithPipeline(cp.getProtobuf(clientVersion)) .setErrorCode(ContainerResponseProto.Error.success) diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java index a595160e8a6..1b0aea5f78e 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java @@ -235,18 +235,22 @@ public void join() throws InterruptedException { @Override public ContainerWithPipeline allocateContainer(HddsProtos.ReplicationType replicationType, HddsProtos.ReplicationFactor factor, - String owner) throws IOException { + String owner, HddsProtos.StorageTierProto storageTier) throws IOException { ReplicationConfig replicationConfig = ReplicationConfig.fromProtoTypeAndFactor(replicationType, factor); - return allocateContainer(replicationConfig, owner); + return allocateContainer(replicationConfig, owner, storageTier); } @Override - public ContainerWithPipeline allocateContainer(ReplicationConfig replicationConfig, String owner) throws IOException { + public ContainerWithPipeline allocateContainer(ReplicationConfig replicationConfig, String owner, + HddsProtos.StorageTierProto storageTier) throws IOException { + StorageTier tier = storageTier != null + ? StorageTier.fromProto(storageTier) : StorageTier.getDefaultTier(); Map<String, String> auditMap = Maps.newHashMap(); auditMap.put("replicationType", String.valueOf(replicationConfig.getReplicationType())); auditMap.put("replication", String.valueOf(replicationConfig.getReplication())); auditMap.put("owner", String.valueOf(owner)); + auditMap.put("storageTier", String.valueOf(tier)); try { if (scm.getScmContext().isInSafeMode()) { @@ -255,7 +259,7 @@ public ContainerWithPipeline allocateContainer(ReplicationConfig replicationConf } getScm().checkAdminAccess(getRemoteUser(), false); final ContainerInfo container = scm.getContainerManager() - .allocateContainer(replicationConfig, owner); + .allocateContainer(replicationConfig, owner, tier); if (container == null) { throw new SCMException( "Could not allocate container for replication " + replicationConfig diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java index 100c9feb9b6..7402d20f637 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java @@ -24,6 +24,8 @@ import static org.apache.hadoop.hdds.utils.HddsServerUtil.getRemoteUser; import static org.apache.hadoop.hdds.utils.HddsServerUtil.getScmSecurityClientWithMaxRetry; import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_ADMINISTRATORS; +import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_DEFAULT_STORAGE_TIER_DEFAULT; +import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_DEFAULT_STORAGE_TIER_KEY; import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_READONLY_ADMINISTRATORS; import static org.apache.hadoop.ozone.OzoneConsts.SCM_ROOT_CA_COMPONENT_NAME; import static org.apache.hadoop.ozone.OzoneConsts.SCM_SUB_CA_PREFIX; @@ -46,6 +48,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; +import java.util.Locale; import java.util.Map; import java.util.Objects; import java.util.UUID; @@ -60,6 +63,7 @@ import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.HddsUtils; import org.apache.hadoop.hdds.annotation.InterfaceAudience; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.conf.ReconfigurationHandler; @@ -338,6 +342,15 @@ public StorageContainerManager(OzoneConfiguration conf) this(conf, new SCMConfigurator()); } + @VisibleForTesting + static StorageTier getConfiguredDefaultStorageTier( + ConfigurationSource conf) { + String configuredTier = conf.get( + OZONE_DEFAULT_STORAGE_TIER_KEY, OZONE_DEFAULT_STORAGE_TIER_DEFAULT); + return StorageTier.valueOf( + configuredTier.trim().toUpperCase(Locale.ROOT)); + } + /** * This constructor offers finer control over how SCM comes up. * To use this, user needs to create a SCMConfigurator and set various @@ -459,6 +472,7 @@ private StorageContainerManager(OzoneConfiguration conf, scmAdmins = OzoneAdmins.getOzoneAdmins(scmStarterUser, conf); scmReadOnlyAdmins = OzoneAdmins.getReadonlyAdmins(conf); LOG.info("SCM start with adminUsers: {}", scmAdmins.getAdminUsernames()); + StorageTier.setDefault(getConfiguredDefaultStorageTier(conf)); datanodeProtocolServer = new SCMDatanodeProtocolServer(conf, this, eventQueue, scmContext); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/HddsTestUtils.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/HddsTestUtils.java index 71f60032fe1..81d91057289 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/HddsTestUtils.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/HddsTestUtils.java @@ -35,6 +35,7 @@ import org.apache.commons.lang3.RandomUtils; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; @@ -525,7 +526,7 @@ public static CommandStatusReportsProto createCommandStatusReport( return containerManager .allocateContainer(RatisReplicationConfig .getInstance(ReplicationFactor.THREE), - "root"); + "root", StorageTier.getDefaultTier()); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java index d5c5604d56a..7b9fe108a0e 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java @@ -38,6 +38,7 @@ import java.util.Collections; import java.util.List; import java.util.concurrent.TimeoutException; +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.client.StorageTier; @@ -95,7 +96,7 @@ void setUp() throws Exception { final OzoneConfiguration conf = SCMTestUtils.getConf(testDir); dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get()); scmhaManager = SCMHAManagerStub.getInstance(true); - NodeManager nodeManager = new MockNodeManager(true, 10); + NodeManager nodeManager = new MockNodeManager(true, 10, StorageType.DISK); sequenceIdGen = new SequenceIdGenerator( conf, scmhaManager, SCMDBDefinition.SEQUENCE_ID.getTable(dbStore)); PipelineManager base = new MockPipelineManager(dbStore, scmhaManager, nodeManager); @@ -128,7 +129,7 @@ void testAllocateContainer() throws Exception { containerManager.getContainers().isEmpty()); final ContainerInfo container = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "admin"); + ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); assertNotNull(container); assertEquals(1, containerManager.getContainers().size()); @@ -150,14 +151,15 @@ public void testGetMatchingContainerReturnsNullWhenNotEnoughSpaceInDatanodes() t // MockPipelineManager#checkSpaceAndRecordAllocation always returns false // the pipeline has no existing containers, so a new container gets allocated in getMatchingContainer ContainerInfo container = containerManager - .getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet()); + .getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet(), StorageTier.getDefaultTier()); assertNull(container); // create an EC pipeline to test for EC containers ECReplicationConfig ecReplicationConfig = new ECReplicationConfig(3, 2); pipelineManager.createPipeline(ecReplicationConfig, StorageTier.getDefaultTier()); pipeline = pipelineManager.getPipelines(ecReplicationConfig).iterator().next(); - container = containerManager.getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet()); + container = containerManager.getMatchingContainer(sizeRequired, "test", + pipeline, Collections.emptySet(), StorageTier.getDefaultTier()); assertNull(container); } @@ -180,14 +182,15 @@ public void testGetMatchingContainerReturnsContainerWhenEnoughSpaceInDatanodes() Pipeline pipeline = spyPipelineManager.getPipelines().iterator().next(); // the pipeline has no existing containers, so a new container gets allocated in getMatchingContainer ContainerInfo container = manager - .getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet()); + .getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet(), StorageTier.getDefaultTier()); assertNotNull(container); // create an EC pipeline to test for EC containers ECReplicationConfig ecReplicationConfig = new ECReplicationConfig(3, 2); spyPipelineManager.createPipeline(ecReplicationConfig, StorageTier.getDefaultTier()); pipeline = spyPipelineManager.getPipelines(ecReplicationConfig).iterator().next(); - container = manager.getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet()); + container = manager.getMatchingContainer(sizeRequired, "test", pipeline, + Collections.emptySet(), StorageTier.getDefaultTier()); assertNotNull(container); } @@ -195,7 +198,7 @@ public void testGetMatchingContainerReturnsContainerWhenEnoughSpaceInDatanodes() void testUpdateContainerState() throws Exception { final ContainerInfo container = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "admin"); + ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); final ContainerID cid = container.containerID(); assertEquals(LifeCycleState.OPEN, containerManager.getContainer(cid).getState()); containerManager.updateContainerState(cid, @@ -217,8 +220,10 @@ void testTransitionDeletingOrDeletedToTargetState(HddsProtos.LifeCycleState desi // Allocate OPEN Ratis and Ec containers, and do a series of state changes to transition them to DELETING / DELETED final ContainerInfo container = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "admin"); - ContainerInfo ecContainer = containerManager.allocateContainer(new ECReplicationConfig(3, 2), "admin"); + ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); + ContainerInfo ecContainer = containerManager.allocateContainer( + new ECReplicationConfig(3, 2), "admin", + StorageTier.getDefaultTier()); final ContainerID cid = container.containerID(); final ContainerID ecCid = ecContainer.containerID(); assertEquals(LifeCycleState.OPEN, containerManager.getContainer(cid).getState()); @@ -266,14 +271,16 @@ void testTransitionContainerToClosedStateAllowOnlyDeletingOrDeletedContainers() // test for RATIS container final ContainerInfo container = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "admin"); + ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); final ContainerID cid = container.containerID(); assertEquals(LifeCycleState.OPEN, containerManager.getContainer(cid).getState()); assertThrows(IOException.class, () -> containerManager.transitionDeletingOrDeletedToTargetState(cid, LifeCycleState.CLOSED)); // test for EC container - final ContainerInfo ecContainer = containerManager.allocateContainer(new ECReplicationConfig(3, 2), "admin"); + final ContainerInfo ecContainer = containerManager.allocateContainer( + new ECReplicationConfig(3, 2), "admin", + StorageTier.getDefaultTier()); final ContainerID ecCid = ecContainer.containerID(); assertEquals(LifeCycleState.OPEN, containerManager.getContainer(ecCid).getState()); assertThrows(IOException.class, () -> @@ -287,7 +294,7 @@ void testGetContainers() throws Exception { List<ContainerID> ids = new ArrayList<>(); for (int i = 0; i < 10; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); ids.add(container.containerID()); } @@ -345,7 +352,7 @@ private static void assertIds( @Test void testAllocateContainersWithECReplicationConfig() throws Exception { final ContainerInfo admin = containerManager - .allocateContainer(new ECReplicationConfig(3, 2), "admin"); + .allocateContainer(new ECReplicationConfig(3, 2), "admin", StorageTier.getDefaultTier()); assertEquals(1, containerManager.getContainers().size()); assertNotNull( containerManager.getContainer(admin.containerID())); @@ -356,7 +363,7 @@ void testUpdateContainerReplicaInvokesPendingOp() throws IOException, TimeoutException { final ContainerInfo container = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "admin"); + ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); DatanodeDetails dn = MockDatanodeDetails.randomDatanodeDetails(); containerManager.updateContainerReplica(container.containerID(), ContainerReplica.newBuilder() @@ -376,7 +383,7 @@ void testRemoveContainerReplicaInvokesPendingOp() throws IOException, TimeoutException { final ContainerInfo container = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "admin"); + ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); DatanodeDetails dn = MockDatanodeDetails.randomDatanodeDetails(); containerManager.removeContainerReplica(container.containerID(), ContainerReplica.newBuilder() diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java index 525bbb556d4..335116f78f3 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java @@ -22,6 +22,7 @@ import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.fail; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.isNull; @@ -37,7 +38,9 @@ import java.time.Clock; import java.time.ZoneId; import java.util.ArrayList; +import java.util.NavigableSet; import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.TimeoutException; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.ECReplicationConfig; @@ -214,6 +217,29 @@ public void checkReplicationStateMissingReplica() assertEquals(3, c1.getReplicationConfig().getRequiredNodes()); } + @Test + public void testMatchingContainerRequiresExactStorageTier() + throws IOException { + ContainerInfo legacyContainer = createContainer(1, null); + containerStateManager.addContainer(legacyContainer.getProtobuf()); + NavigableSet<ContainerID> containerIDs = new TreeSet<>(); + containerIDs.add(legacyContainer.containerID()); + + assertNull(containerStateManager.getMatchingContainerAndStorageTier( + 0, "root", pipeline.getId(), containerIDs, StorageTier.DISK)); + + ContainerInfo diskContainer = createContainer(2, StorageTier.DISK); + containerStateManager.addContainer(diskContainer.getProtobuf()); + containerIDs.add(diskContainer.containerID()); + + assertEquals(diskContainer.containerID(), + containerStateManager.getMatchingContainerAndStorageTier( + 0, "root", pipeline.getId(), containerIDs, StorageTier.DISK) + .containerID()); + assertNull(containerStateManager.getMatchingContainerAndStorageTier( + 0, "root", pipeline.getId(), containerIDs, StorageTier.ARCHIVE)); + } + @ParameterizedTest @EnumSource(value = HddsProtos.LifeCycleState.class, names = {"DELETING", "DELETED"}) @@ -486,19 +512,27 @@ private void addReplica(ContainerInfo cont, DatanodeDetails node) { private ContainerInfo allocateContainer() throws IOException, TimeoutException { - final ContainerInfo containerInfo = new ContainerInfo.Builder() + final ContainerInfo containerInfo = createContainer(1, null); + + containerStateManager.addContainer(containerInfo.getProtobuf()); + return containerInfo; + } + + private ContainerInfo createContainer( + long containerID, StorageTier storageTier) { + ContainerInfo.Builder builder = new ContainerInfo.Builder() .setState(HddsProtos.LifeCycleState.OPEN) .setPipelineID(pipeline.getId()) .setUsedBytes(0) .setNumberOfKeys(0) .setOwner("root") - .setContainerID(1) + .setContainerID(containerID) .setDeleteTransactionId(0) - .setReplicationConfig(pipeline.getReplicationConfig()) - .build(); - - containerStateManager.addContainer(containerInfo.getProtobuf()); - return containerInfo; + .setReplicationConfig(pipeline.getReplicationConfig()); + if (storageTier != null) { + builder.setStorageTier(storageTier); + } + return builder.build(); } private static StorageContainerDatanodeProtocolProtos.ContainerReportsProto getContainerReportsProto( diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java index d3015aca234..09e5320eca5 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java @@ -236,7 +236,7 @@ public void testContainerPlacementCapacity() throws IOException, ReplicationConfig.fromProtoTypeAndFactor( SCMTestUtils.getReplicationType(conf), SCMTestUtils.getReplicationFactor(conf)), - OzoneConsts.OZONE); + OzoneConsts.OZONE, StorageTier.getDefaultTier()); assertNotNull(container, "allocateContainer returned null (unexpected in this placement test)"); int replicaCount = 0; diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestNodeDecommissionManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestNodeDecommissionManager.java index 4a3ccecd725..271cabf6974 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestNodeDecommissionManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestNodeDecommissionManager.java @@ -42,6 +42,7 @@ import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; @@ -81,7 +82,7 @@ void setup(@TempDir File dir) throws Exception { containerManager = mock(ContainerManager.class); decom = new NodeDecommissionManager(conf, nodeManager, containerManager, SCMContext.emptyContext(), new EventQueue(), null); - when(containerManager.allocateContainer(any(ReplicationConfig.class), anyString())) + when(containerManager.allocateContainer(any(ReplicationConfig.class), anyString(), any(StorageTier.class))) .thenAnswer(invocation -> createMockContainer((ReplicationConfig)invocation.getArguments()[0], (String) invocation.getArguments()[1])); } @@ -423,7 +424,9 @@ public void testInsufficientNodeDecommissionThrowsExceptionForRatis() throws Set<ContainerID> idsRatis = new HashSet<>(); for (int i = 0; i < 5; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), + "admin", StorageTier.getDefaultTier()); idsRatis.add(container.containerID()); } @@ -477,7 +480,9 @@ public void testInsufficientNodeDecommissionThrowsExceptionForEc() throws Set<ContainerID> idsEC = new HashSet<>(); for (int i = 0; i < 5; i++) { - ContainerInfo container = containerManager.allocateContainer(new ECReplicationConfig(3, 2), "admin"); + ContainerInfo container = containerManager.allocateContainer( + new ECReplicationConfig(3, 2), "admin", + StorageTier.getDefaultTier()); idsEC.add(container.containerID()); } @@ -513,10 +518,12 @@ public void testInsufficientNodeDecommissionThrowsExceptionRatisAndEc() throws Set<ContainerID> idsRatis = new HashSet<>(); ContainerInfo containerRatis = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); idsRatis.add(containerRatis.containerID()); Set<ContainerID> idsEC = new HashSet<>(); - ContainerInfo containerEC = containerManager.allocateContainer(new ECReplicationConfig(3, 2), "admin"); + ContainerInfo containerEC = containerManager.allocateContainer( + new ECReplicationConfig(3, 2), "admin", + StorageTier.getDefaultTier()); idsEC.add(containerEC.containerID()); when(containerManager.getContainer(any(ContainerID.class))) @@ -570,7 +577,9 @@ public void testInsufficientNodeDecommissionChecksNotInService() throws Set<ContainerID> idsRatis = new HashSet<>(); for (int i = 0; i < 5; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), + "admin", StorageTier.getDefaultTier()); idsRatis.add(container.containerID()); } @@ -606,7 +615,9 @@ public void testInsufficientNodeDecommissionChecksForNNF() throws Set<ContainerID> idsRatis = new HashSet<>(); for (int i = 0; i < 3; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), + "admin", StorageTier.getDefaultTier()); idsRatis.add(container.containerID()); } @@ -667,7 +678,9 @@ public void testInsufficientNodeMaintenanceThrowsExceptionForRatis() throws Set<ContainerID> idsRatis = new HashSet<>(); for (int i = 0; i < 5; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), + "admin", StorageTier.getDefaultTier()); idsRatis.add(container.containerID()); } for (DatanodeDetails dn : nodeManager.getAllNodes().subList(0, 3)) { @@ -767,7 +780,9 @@ public void testInsufficientNodeMaintenanceThrowsExceptionForEc() throws } Set<ContainerID> idsEC = new HashSet<>(); for (int i = 0; i < 5; i++) { - ContainerInfo container = containerManager.allocateContainer(new ECReplicationConfig(3, 2), "admin"); + ContainerInfo container = containerManager.allocateContainer( + new ECReplicationConfig(3, 2), "admin", + StorageTier.getDefaultTier()); idsEC.add(container.containerID()); } for (DatanodeDetails dn : nodeManager.getAllNodes()) { @@ -837,10 +852,12 @@ public void testInsufficientNodeMaintenanceThrowsExceptionForRatisAndEc() throws } Set<ContainerID> idsRatis = new HashSet<>(); ContainerInfo containerRatis = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin", StorageTier.getDefaultTier()); idsRatis.add(containerRatis.containerID()); Set<ContainerID> idsEC = new HashSet<>(); - ContainerInfo containerEC = containerManager.allocateContainer(new ECReplicationConfig(3, 2), "admin"); + ContainerInfo containerEC = containerManager.allocateContainer( + new ECReplicationConfig(3, 2), "admin", + StorageTier.getDefaultTier()); idsEC.add(containerEC.containerID()); when(containerManager.getContainer(any(ContainerID.class))) @@ -924,7 +941,9 @@ public void testInsufficientNodeMaintenanceChecksNotInService() throws Set<ContainerID> idsRatis = new HashSet<>(); for (int i = 0; i < 5; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), + "admin", StorageTier.getDefaultTier()); idsRatis.add(container.containerID()); } for (DatanodeDetails dn : nodeManager.getAllNodes().subList(0, 3)) { @@ -964,7 +983,9 @@ public void testInsufficientNodeMaintenanceChecksForNNF() throws Set<ContainerID> idsRatis = new HashSet<>(); for (int i = 0; i < 3; i++) { ContainerInfo container = containerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), "admin"); + RatisReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), + "admin", StorageTier.getDefaultTier()); idsRatis.add(container.containerID()); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java index 39d78ab1229..3120ff19e3e 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java @@ -81,7 +81,10 @@ public Pipeline createPipeline(ReplicationConfig replicationConfig, if (replicationConfig.getReplicationType() == HddsProtos.ReplicationType.EC) { pipeline = buildECPipeline( - replicationConfig, excludedNodes, favoredNodes); + replicationConfig, excludedNodes, favoredNodes) + .toBuilder() + .setSupportedStorageTier(storageTier) + .build(); } else { pipeline = createPipeline(replicationConfig, ImmutableList.of(MockDatanodeDetails.randomDatanodeDetails(), diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java index bd542581ae0..10f32987285 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java @@ -953,7 +953,7 @@ public void testWaitForAllocatedPipeline() throws IOException { pipelineManager.addContainerToPipeline( allocatedPipeline.getId(), container.containerID()); doReturn(container).when(containerManager).getMatchingContainer(anyLong(), - anyString(), eq(allocatedPipeline), any()); + anyString(), eq(allocatedPipeline), any(), any(StorageTier.class)); assertTrue(pipelineManager.getPipelines(repConfig, OPEN) diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java index 6b395aab0ce..7d59d5f0e04 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java @@ -729,6 +729,29 @@ public void testCurrentRatisThreePipelineCount() assertEquals(pipelineCount, 2); } + @Test + public void testCurrentRatisThreePipelineCountIgnoresLegacyPipeline() + throws IOException { + List<DatanodeDetails> healthyNodes = nodeManager + .getNodes(NodeStatus.inServiceHealthy()); + List<DatanodeDetails> pipelineNodes = healthyNodes.subList(0, 3); + Pipeline legacyPipeline = Pipeline.newBuilder() + .setId(PipelineID.randomId()) + .setState(Pipeline.PipelineState.OPEN) + .setReplicationConfig(RatisReplicationConfig.getInstance( + ReplicationFactor.THREE)) + .setNodes(pipelineNodes) + .build(); + + nodeManager.addPipeline(legacyPipeline); + stateManager.addPipeline(legacyPipeline.getProtobufMessage( + ClientVersion.CURRENT_VERSION)); + + assertEquals(0, PipelinePlacementPolicy.currentRatisThreePipelineCount( + nodeManager, stateManager, pipelineNodes.get(0), + StorageType.DEFAULT)); + } + @Test public void testPipelinePlacementPolicyDefaultLimitFiltersNodeAtLimit() throws IOException, TimeoutException { diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateMap.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateMap.java index 446d6b42608..2f0e3efc780 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateMap.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateMap.java @@ -25,6 +25,7 @@ import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -93,4 +94,34 @@ public void testCountPipelines() throws IOException { assertEquals(1, map.getPipelineCount(new ECReplicationConfig(3, 2), Pipeline.PipelineState.CLOSED)); } + + @Test + public void testGetPipelinesRequiresExactStorageTier() throws IOException { + Pipeline legacyPipeline = MockPipeline.createPipeline(1); + Pipeline diskPipeline = withStorageTier( + MockPipeline.createPipeline(1), StorageTier.DISK); + map.addPipeline(legacyPipeline); + map.addPipeline(diskPipeline); + + assertEquals(1, map.getPipelines( + StandaloneReplicationConfig.getInstance(ONE), + Pipeline.PipelineState.OPEN, StorageTier.DISK).size()); + assertEquals(diskPipeline, map.getPipelines( + StandaloneReplicationConfig.getInstance(ONE), + Pipeline.PipelineState.OPEN, StorageTier.DISK).get(0)); + assertEquals(0, map.getPipelines( + StandaloneReplicationConfig.getInstance(ONE), + Pipeline.PipelineState.OPEN, StorageTier.ARCHIVE).size()); + } + + private static Pipeline withStorageTier( + Pipeline pipeline, StorageTier storageTier) { + return Pipeline.newBuilder() + .setState(pipeline.getPipelineState()) + .setId(pipeline.getId()) + .setReplicationConfig(pipeline.getReplicationConfig()) + .setNodes(pipeline.getNodes()) + .setSupportedStorageTier(storageTier) + .build(); + } } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java index 29de71bd6b5..595cc2a3445 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java @@ -145,7 +145,7 @@ void setup(@TempDir File testDir) throws IOException { containers.put(container.containerID(), container); return container; }).when(containerManager).getMatchingContainer(anyLong(), - anyString(), any(Pipeline.class)); + anyString(), any(Pipeline.class), any(), any(StorageTier.class)); doAnswer(call -> containers.get((ContainerID)call.getArguments()[0])) diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java index 6c43a0d1759..9bc26609965 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java @@ -141,7 +141,8 @@ private ContainerInfo pipelineHasContainer(Pipeline pipeline) { .setPipelineID(pipeline.getId()) .build(); - when(containerManager.getMatchingContainer(CONTAINER_SIZE, OWNER, pipeline, emptySet())) + when(containerManager.getMatchingContainer( + CONTAINER_SIZE, OWNER, pipeline, emptySet(), StorageTier.getDefaultTier())) .thenReturn(container); return container; diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestStorageContainerManagerStarter.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestStorageContainerManagerStarter.java index 71a4d2a084d..e28f5310af5 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestStorageContainerManagerStarter.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestStorageContainerManagerStarter.java @@ -17,6 +17,7 @@ package org.apache.hadoop.hdds.scm.server; +import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_DEFAULT_STORAGE_TIER_KEY; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -29,6 +30,7 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import org.apache.hadoop.hdds.cli.GenericCli; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -138,6 +140,15 @@ public void testGenClusterIdWithInvalidParamDoesNotRun() { assertFalse(mock.generateCalled); } + @Test + public void testConfiguredDefaultStorageTierIgnoresCaseAndWhitespace() { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(OZONE_DEFAULT_STORAGE_TIER_KEY, " aRcHiVe "); + + assertEquals(StorageTier.ARCHIVE, + StorageContainerManager.getConfiguredDefaultStorageTier(conf)); + } + @Test public void testUsagePrintedOnInvalidInput() throws UnsupportedEncodingException { diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java index 64738d9c7a8..3f396805aeb 100644 --- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java +++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java @@ -126,7 +126,8 @@ public void testDatanodeRegistrationAndReports() throws Exception { ContainerManager reconContainerManager = reconScm.getContainerManager(); ContainerInfo containerInfo = scmContainerManager - .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test"); + .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test", + StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); Pipeline pipeline = scmPipelineManager.getPipeline(containerInfo.getPipelineID()); @@ -166,7 +167,8 @@ public void testReconRestart() throws Exception { // Create container in SCM. ContainerInfo containerInfo = scmContainerManager - .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test"); + .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test", + StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); PipelineManager scmPipelineManager = scm.getPipelineManager(); Pipeline pipeline = diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconScmSnapshot.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconScmSnapshot.java index 9414d92f141..08ef192c50a 100644 --- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconScmSnapshot.java +++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconScmSnapshot.java @@ -23,6 +23,7 @@ import java.util.List; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -88,7 +89,8 @@ public void testScmSnapshot() throws Exception { for (int i = 0; i < 10; i++) { containerManager.allocateContainer(RatisReplicationConfig.getInstance( - HddsProtos.ReplicationFactor.ONE), "testOwner"); + HddsProtos.ReplicationFactor.ONE), "testOwner", + StorageTier.getDefaultTier()); } recon.start(conf); diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java index f847812edcd..d98e9649afa 100644 --- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java +++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java @@ -31,6 +31,7 @@ import java.util.List; import java.util.Set; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -170,9 +171,9 @@ public void testSyncSCMContainerInfo() throws Exception { ContainerManager reconCm = reconScm.getContainerManager(); final ContainerInfo container1 = scmContainerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE), "admin"); + RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE), "admin", StorageTier.getDefaultTier()); final ContainerInfo container2 = scmContainerManager.allocateContainer( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE), "admin"); + RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE), "admin", StorageTier.getDefaultTier()); scmContainerManager.updateContainerState(container1.containerID(), HddsProtos.LifeCycleEvent.FINALIZE); scmContainerManager.updateContainerState(container2.containerID(), @@ -222,7 +223,7 @@ public void testContainerHealthTaskDetectsUnderReplicatedAfterNodeFailure() ContainerInfo containerInfo = scmContainerManager.allocateContainer( RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), - "test"); + "test", StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); Pipeline pipeline = scmPipelineManager.getPipeline(containerInfo.getPipelineID()); @@ -361,7 +362,7 @@ public void testContainerHealthTaskDetectsEmptyMissingWhenAllReplicasLost() ReconContainerManager reconCm = (ReconContainerManager) reconScm.getContainerManager(); ContainerInfo containerInfo = scm.getContainerManager() - .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test"); + .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test", StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); Pipeline pipeline = scmPipelineManager.getPipeline(containerInfo.getPipelineID()); @@ -464,7 +465,7 @@ public void testContainerHealthTaskDetectsMissingForContainerWithKeys() (ReconContainerManager) reconScm.getContainerManager(); ContainerInfo containerInfo = scm.getContainerManager() - .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test"); + .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test", StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); Pipeline pipeline = scmPipelineManager.getPipeline(containerInfo.getPipelineID()); @@ -569,7 +570,7 @@ public void testContainerHealthTaskDetectsOverReplicatedAndNegativeSize() (ReconContainerManager) reconScm.getContainerManager(); ContainerInfo containerInfo = scm.getContainerManager() - .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test"); + .allocateContainer(RatisReplicationConfig.getInstance(ONE), "test", StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); Pipeline pipeline = scmPipelineManager.getPipeline(containerInfo.getPipelineID()); @@ -692,7 +693,7 @@ public void testContainerHealthTaskDetectsReplicaMismatch() throws Exception { ContainerInfo containerInfo = scm.getContainerManager().allocateContainer( RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), - "test"); + "test", StorageTier.getDefaultTier()); long containerID = containerInfo.getContainerID(); Pipeline pipeline = scmPipelineManager.getPipeline(containerInfo.getPipelineID()); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshot.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshot.java index 53a2f75af22..7b82c9ae2d9 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshot.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshot.java @@ -32,6 +32,7 @@ import java.util.Map; import org.apache.hadoop.fs.FileUtil; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.ContainerManager; @@ -93,12 +94,12 @@ private DBCheckpoint downloadSnapshot() throws Exception { PipelineManager pipelineManager = scm.getPipelineManager(); Pipeline ratisPipeline1 = pipelineManager.getPipeline( containerManager.allocateContainer( - RatisReplicationConfig.getInstance(THREE), "Owner1") + RatisReplicationConfig.getInstance(THREE), "Owner1", StorageTier.getDefaultTier()) .getPipelineID()); pipelineManager.openPipeline(ratisPipeline1.getId()); Pipeline ratisPipeline2 = pipelineManager.getPipeline( containerManager.allocateContainer( - RatisReplicationConfig.getInstance(ONE), "Owner2") + RatisReplicationConfig.getInstance(ONE), "Owner2", StorageTier.getDefaultTier()) .getPipelineID()); pipelineManager.openPipeline(ratisPipeline2.getId()); SCMNodeDetails scmNodeDetails = new SCMNodeDetails.Builder() diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshotWithHA.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshotWithHA.java index 4412c23402b..b47c33f8aa9 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshotWithHA.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMInstallSnapshotWithHA.java @@ -34,6 +34,7 @@ import org.apache.commons.io.FileUtils; import org.apache.hadoop.hdds.ExitManager; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; import org.apache.hadoop.hdds.scm.container.ContainerInfo; @@ -277,7 +278,8 @@ private List<ContainerInfo> writeToIncreaseLogIndex( containers.add(scm.getContainerManager() .allocateContainer( RatisReplicationConfig.getInstance(ReplicationFactor.THREE), - TestSCMInstallSnapshotWithHA.class.getName())); + TestSCMInstallSnapshotWithHA.class.getName(), + StorageTier.getDefaultTier())); Thread.sleep(100); logIndex = stateMachine.getLastAppliedTermIndex().getIndex(); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMMXBean.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMMXBean.java index f2ae925dcbd..417ea514116 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMMXBean.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMMXBean.java @@ -33,6 +33,7 @@ import javax.management.openmbean.CompositeData; import javax.management.openmbean.TabularData; import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; import org.apache.hadoop.hdds.scm.container.ContainerID; @@ -111,7 +112,7 @@ public void testSCMContainerStateCount() throws Exception { containerInfoList.add( scmContainerManager.allocateContainer( StandaloneReplicationConfig.getInstance(ReplicationFactor.ONE), - UUID.randomUUID().toString())); + UUID.randomUUID().toString(), StorageTier.getDefaultTier())); } long containerID; for (int i = 0; i < 10; i++) { diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMSnapshot.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMSnapshot.java index 703e6bf30cd..62b288b941d 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMSnapshot.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMSnapshot.java @@ -23,6 +23,7 @@ import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.scm.container.ContainerManager; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; @@ -61,12 +62,13 @@ public void testSnapshot() throws Exception { PipelineManager pipelineManager = scm.getPipelineManager(); Pipeline ratisPipeline1 = pipelineManager.getPipeline( containerManager.allocateContainer( - RatisReplicationConfig.getInstance(THREE), "Owner1") + RatisReplicationConfig.getInstance(THREE), "Owner1", StorageTier.getDefaultTier()) .getPipelineID()); pipelineManager.openPipeline(ratisPipeline1.getId()); Pipeline ratisPipeline2 = pipelineManager.getPipeline( containerManager.allocateContainer( - RatisReplicationConfig.getInstance(ONE), "Owner2").getPipelineID()); + RatisReplicationConfig.getInstance(ONE), "Owner2", + StorageTier.getDefaultTier()).getPipelineID()); pipelineManager.openPipeline(ratisPipeline2.getId()); long snapshotInfo2 = scm.getScmHAManager().asSCMHADBTransactionBuffer() .getLatestTrxInfo().getTransactionIndex(); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSecretKeySnapshot.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSecretKeySnapshot.java index bd379a8d86e..ba026d25be6 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSecretKeySnapshot.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSecretKeySnapshot.java @@ -52,6 +52,7 @@ import java.util.Properties; import java.util.concurrent.TimeoutException; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; import org.apache.hadoop.hdds.scm.container.ContainerInfo; @@ -272,7 +273,7 @@ private List<ContainerInfo> writeToIncreaseLogIndex( containers.add(scm.getContainerManager() .allocateContainer( RatisReplicationConfig.getInstance(ReplicationFactor.ONE), - this.getClass().getName())); + this.getClass().getName(), StorageTier.getDefaultTier())); Thread.sleep(100); logIndex = stateMachine.getLastAppliedTermIndex().getIndex(); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestAllocateContainerWithStorageTier.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestAllocateContainerWithStorageTier.java new file mode 100644 index 00000000000..ac54c014cc4 --- /dev/null +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestAllocateContainerWithStorageTier.java @@ -0,0 +1,254 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.hdds.scm.container; + +import static java.util.Collections.emptyList; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.File; +import java.io.IOException; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +/** + * End-to-end tests that a container is allocated on a pipeline whose + * datanodes advertise the requested {@link StorageTier}, and that + * {@link PipelineManager#getPipelines(ReplicationConfig, PipelineState, + * java.util.Collection, java.util.Collection, StorageTier)} filters by tier. + */ +public class TestAllocateContainerWithStorageTier { + + @TempDir + private File dir; + private MiniOzoneCluster cluster; + private StorageContainerManager scm; + private ContainerManager containerManager; + private final ReplicationConfig ratisThree = + RatisReplicationConfig.getInstance(ReplicationFactor.THREE); + private final ReplicationConfig standaloneOne = + StandaloneReplicationConfig.getInstance(ReplicationFactor.ONE); + + private void createCluster(List<List<StorageType>> storageTypeList) throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setInt(OZONE_DATANODE_PIPELINE_LIMIT, 1); + cluster = MiniOzoneCluster.newBuilder(conf) + .setNumDatanodes(storageTypeList.size()) + .setNumDataVolumes(storageTypeList.get(0).size()) + .setDatanodeStorageType(storageTypeList) + .build(); + cluster.waitForClusterToBeReady(); + cluster.waitTobeOutOfSafeMode(); + + scm = cluster.getStorageContainerManager(); + containerManager = scm.getContainerManager(); + } + + private void cleanUp() { + if (cluster != null) { + cluster.shutdown(); + } + } + + @Test + public void testAllocateContainerWithStorageTierAllTiersAvailable() throws Exception { + createCluster(Arrays.asList( + Arrays.asList(StorageType.DISK, StorageType.SSD, StorageType.ARCHIVE), + Arrays.asList(StorageType.DISK, StorageType.SSD, StorageType.ARCHIVE), + Arrays.asList(StorageType.DISK, StorageType.SSD, StorageType.ARCHIVE))); + try { + PipelineManager pipelineManager = scm.getPipelineManager(); + assertTrue(containerManager.getContainers().isEmpty()); + ContainerInfo container; + List<Pipeline> pipelines; + + // All datanodes have DISK, SSD, ARCHIVE volumes: allocate on any tier. + container = containerManager.allocateContainer(ratisThree, "admin", StorageTier.DISK); + assertContainer(container.containerID(), containerManager, ratisThree, StorageTier.DISK, 1); + + container = containerManager.allocateContainer(ratisThree, "admin", StorageTier.SSD); + assertContainer(container.containerID(), containerManager, ratisThree, StorageTier.SSD, 2); + + container = containerManager.allocateContainer(ratisThree, "admin", StorageTier.ARCHIVE); + assertContainer(container.containerID(), containerManager, ratisThree, StorageTier.ARCHIVE, 3); + + container = containerManager.allocateContainer(standaloneOne, "admin", StorageTier.DISK); + assertContainer(container.containerID(), containerManager, standaloneOne, StorageTier.DISK, 4); + + container = containerManager.allocateContainer(standaloneOne, "admin", StorageTier.SSD); + assertContainer(container.containerID(), containerManager, standaloneOne, StorageTier.SSD, 5); + + container = containerManager.allocateContainer(standaloneOne, "admin", StorageTier.ARCHIVE); + assertContainer(container.containerID(), containerManager, standaloneOne, StorageTier.ARCHIVE, 6); + + pipelines = pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.DISK); + assertPipeline(pipelines, 1, StorageTier.DISK); + + pipelines = pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.SSD); + assertPipeline(pipelines, 1, StorageTier.SSD); + + pipelines = pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.ARCHIVE); + assertPipeline(pipelines, 1, StorageTier.ARCHIVE); + + } finally { + cleanUp(); + } + } + + @Test + public void testAllocateContainerWithStorageTierPartialTierAvailability() throws Exception { + createCluster(Arrays.asList( + Arrays.asList(StorageType.DISK, StorageType.DISK, StorageType.ARCHIVE), + Arrays.asList(StorageType.DISK, StorageType.SSD, StorageType.ARCHIVE), + Arrays.asList(StorageType.DISK, StorageType.SSD, StorageType.SSD))); + try { + PipelineManager pipelineManager = scm.getPipelineManager(); + assertTrue(containerManager.getContainers().isEmpty()); + ContainerInfo container; + List<Pipeline> pipelines; + + // Every datanode has DISK; SSD and ARCHIVE are not uniformly available + // across all three, so RATIS_THREE only works on DISK. + assertThrows(IOException.class, () -> containerManager.allocateContainer( + ratisThree, "admin", StorageTier.SSD)); + + container = containerManager.allocateContainer(ratisThree, "admin", StorageTier.DISK); + assertContainer(container.containerID(), containerManager, ratisThree, StorageTier.DISK, 1); + + assertThrows(IOException.class, () -> containerManager.allocateContainer( + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), "admin", StorageTier.ARCHIVE)); + + container = containerManager.allocateContainer(standaloneOne, "admin", StorageTier.DISK); + assertContainer(container.containerID(), containerManager, standaloneOne, StorageTier.DISK, 2); + + container = containerManager.allocateContainer(standaloneOne, "admin", StorageTier.ARCHIVE); + assertContainer(container.containerID(), containerManager, standaloneOne, StorageTier.ARCHIVE, 3); + + pipelines = pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.DISK); + assertPipeline(pipelines, 1, StorageTier.DISK); + + assertTrue(pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.SSD).isEmpty()); + + assertTrue(pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.ARCHIVE).isEmpty()); + } finally { + cleanUp(); + } + } + + @Test + public void testAllocateContainerWithStorageTierOnlyDisk() throws Exception { + createCluster(Arrays.asList( + Arrays.asList(StorageType.DISK, StorageType.DISK, StorageType.DISK), + Arrays.asList(StorageType.DISK, StorageType.DISK, StorageType.DISK), + Arrays.asList(StorageType.DISK, StorageType.DISK, StorageType.DISK))); + try { + PipelineManager pipelineManager = scm.getPipelineManager(); + assertTrue(containerManager.getContainers().isEmpty()); + ContainerInfo container; + List<Pipeline> pipelines; + + // Only DISK is available. + assertThrows(IOException.class, () -> containerManager.allocateContainer( + ratisThree, "admin", StorageTier.SSD)); + + container = containerManager.allocateContainer(ratisThree, "admin", StorageTier.DISK); + assertContainer(container.containerID(), containerManager, ratisThree, StorageTier.DISK, 1); + + assertThrows(IOException.class, () -> containerManager.allocateContainer( + ratisThree, "admin", StorageTier.ARCHIVE)); + + assertThrows(IOException.class, () -> containerManager.allocateContainer( + standaloneOne, "admin", StorageTier.SSD)); + + container = containerManager.allocateContainer(standaloneOne, "admin", StorageTier.DISK); + assertContainer(container.containerID(), containerManager, standaloneOne, StorageTier.DISK, 2); + + assertThrows(IOException.class, () -> containerManager.allocateContainer( + standaloneOne, "admin", StorageTier.ARCHIVE)); + + pipelines = pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.DISK); + assertPipeline(pipelines, 1, StorageTier.DISK); + + assertTrue(pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.SSD).isEmpty()); + + assertTrue(pipelineManager.getPipelines( + ratisThree, PipelineState.OPEN, emptyList(), emptyList(), StorageTier.ARCHIVE).isEmpty()); + + } finally { + cleanUp(); + } + } + + private void assertContainer(ContainerID containerID, ContainerManager manager, + ReplicationConfig replicationConfig, StorageTier expectedStorageTier, int expectedTotalCount) + throws IOException { + ContainerInfo containerInfo = manager.getContainer(containerID); + Pipeline pipeline = scm.getPipelineManager().getPipeline(containerInfo.getPipelineID()); + + assertNotNull(containerInfo); + assertEquals(expectedTotalCount, manager.getContainers().size()); + assertEquals(replicationConfig.getRequiredNodes(), pipeline.getNodes().size()); + assertEquals(expectedStorageTier, pipeline.getSupportedStorageTier()); + assertEquals(expectedStorageTier, containerInfo.getStorageTier()); + } + + private void assertPipeline(List<Pipeline> pipelines, int expectedContainerCount, + StorageTier expectedStorageTier) { + assertFalse(pipelines.isEmpty()); + List<ContainerInfo> containerInfos = new ArrayList<>(); + for (Pipeline pipeline : pipelines) { + ContainerInfo containerInfo = containerManager.getMatchingContainer(0, "admin", + pipeline, new HashSet<>(), expectedStorageTier); + if (containerInfo != null) { + containerInfos.add(containerInfo); + } + } + assertEquals(expectedContainerCount, containerInfos.size()); + for (ContainerInfo containerInfo : containerInfos) { + assertEquals(expectedStorageTier, containerInfo.getStorageTier()); + } + } +} diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManagerIntegration.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManagerIntegration.java index f448a3a1336..39002405d7d 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManagerIntegration.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManagerIntegration.java @@ -25,6 +25,7 @@ import static org.junit.jupiter.api.Assertions.assertNull; import java.io.IOException; +import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Set; @@ -36,6 +37,7 @@ import java.util.concurrent.TimeoutException; import org.apache.commons.lang3.RandomUtils; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -102,7 +104,7 @@ public void testAllocateContainer() throws IOException { SCMTestUtils.getReplicationFactor(conf), OzoneConsts.OZONE); ContainerInfo info = containerManager .getMatchingContainer(OzoneConsts.GB * 3, OzoneConsts.OZONE, - container1.getPipeline()); + container1.getPipeline(), Collections.emptySet(), StorageTier.getDefaultTier()); assertNotEquals(container1.getContainerInfo().getContainerID(), info.getContainerID()); assertEquals(OzoneConsts.OZONE, info.getOwner()); @@ -130,13 +132,13 @@ public void testAllocateContainerWithDifferentOwner() throws IOException { SCMTestUtils.getReplicationFactor(conf), OzoneConsts.OZONE); ContainerInfo info = containerManager .getMatchingContainer(OzoneConsts.GB * 3, OzoneConsts.OZONE, - container1.getPipeline()); + container1.getPipeline(), Collections.emptySet(), StorageTier.getDefaultTier()); assertNotNull(info); String newContainerOwner = "OZONE_NEW"; ContainerInfo info2 = containerManager .getMatchingContainer(OzoneConsts.GB * 3, newContainerOwner, - container1.getPipeline()); + container1.getPipeline(), Collections.emptySet(), StorageTier.getDefaultTier()); assertNotNull(info2); assertNotEquals(info.containerID(), info2.containerID()); @@ -207,7 +209,7 @@ public void testGetMatchingContainer() throws IOException { for (int i = 1; i < numContainerPerOwnerInPipeline; i++) { ContainerInfo info = containerManager .getMatchingContainer(OzoneConsts.GB * 3, OzoneConsts.OZONE, - container1.getPipeline()); + container1.getPipeline(), Collections.emptySet(), StorageTier.getDefaultTier()); assertThat(info.getContainerID()).isGreaterThan(cid); cid = info.getContainerID(); } @@ -216,7 +218,7 @@ public void testGetMatchingContainer() throws IOException { // next container should be the same as first container ContainerInfo info = containerManager .getMatchingContainer(OzoneConsts.GB * 3, OzoneConsts.OZONE, - container1.getPipeline()); + container1.getPipeline(), Collections.emptySet(), StorageTier.getDefaultTier()); assertEquals(container1.getContainerInfo().getContainerID(), info.getContainerID()); } @@ -238,7 +240,7 @@ public void testGetMatchingContainerMultipleThreads() CompletableFuture.supplyAsync(() -> { ContainerInfo info = containerManager .getMatchingContainer(OzoneConsts.GB * 3, OzoneConsts.OZONE, - container1.getPipeline()); + container1.getPipeline(), Collections.emptySet(), StorageTier.getDefaultTier()); container2MatchedCount .compute(info.getContainerID(), (k, v) -> v == null ? 1L : v + 1); return null; diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java index cef98beaf0e..04319beb605 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java @@ -25,6 +25,7 @@ import java.util.List; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ContainerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; @@ -82,7 +83,8 @@ public void testAddContainerToClosedPipeline() throws Exception { // verify that SCMStateMachine is still functioning after the rejected // transaction. - assertNotNull(containerManager.allocateContainer(replication, "test")); + assertNotNull(containerManager.allocateContainer(replication, "test", + StorageTier.getDefaultTier())); } @Test diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/metrics/TestSCMContainerManagerMetrics.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/metrics/TestSCMContainerManagerMetrics.java index 1391e584766..8d3ea41e3a2 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/metrics/TestSCMContainerManagerMetrics.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/metrics/TestSCMContainerManagerMetrics.java @@ -27,6 +27,7 @@ import org.apache.commons.lang3.RandomUtils; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.ContainerInfo; @@ -78,7 +79,8 @@ public void testContainerOpsMetrics() throws Exception { ContainerInfo containerInfo = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - HddsProtos.ReplicationFactor.ONE), OzoneConsts.OZONE); + HddsProtos.ReplicationFactor.ONE), OzoneConsts.OZONE, + StorageTier.getDefaultTier()); metrics = getMetrics(SCMContainerManagerMetrics.class.getSimpleName()); assertEquals(getLongCounter("NumSuccessfulCreateContainers", @@ -86,7 +88,8 @@ public void testContainerOpsMetrics() throws Exception { assertThrows(IOException.class, () -> containerManager.allocateContainer( - new ECReplicationConfig(8, 5), OzoneConsts.OZONE)); + new ECReplicationConfig(8, 5), OzoneConsts.OZONE, + StorageTier.getDefaultTier())); // allocateContainer should fail, so it should have the old metric value. metrics = getMetrics(SCMContainerManagerMetrics.class.getSimpleName()); assertEquals(getLongCounter("NumSuccessfulCreateContainers", diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerIntegration.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerIntegration.java index 65ce765ca60..0024a903d39 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerIntegration.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerIntegration.java @@ -52,6 +52,7 @@ import java.util.UUID; import java.util.stream.Collectors; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -370,7 +371,9 @@ public void testOneDeadMaintenanceNodeAndOneLiveMaintenanceNodeAndOneDecommissio */ @Test public void testEmptyQuasiClosedContainerDeletion() throws Exception { - ContainerInfo containerInfo = containerManager.allocateContainer(RATIS_REPLICATION_CONFIG, "TestOwner"); + ContainerInfo containerInfo = containerManager.allocateContainer( + RATIS_REPLICATION_CONFIG, "TestOwner", + StorageTier.getDefaultTier()); ContainerID cid = containerInfo.containerID(); containerManager.updateContainerState(cid, HddsProtos.LifeCycleEvent.FINALIZE); containerManager.updateContainerState(cid, HddsProtos.LifeCycleEvent.QUASI_CLOSE); @@ -437,7 +440,9 @@ public void testEmptyQuasiClosedContainerDeletion() throws Exception { */ @Test public void testEmptyQuasiClosedContainerDeletionWithMixedReplicaStates() throws Exception { - ContainerInfo containerInfo = containerManager.allocateContainer(RATIS_REPLICATION_CONFIG, "TestOwner"); + ContainerInfo containerInfo = containerManager.allocateContainer( + RATIS_REPLICATION_CONFIG, "TestOwner", + StorageTier.getDefaultTier()); ContainerID cid = containerInfo.containerID(); containerManager.updateContainerState(cid, HddsProtos.LifeCycleEvent.FINALIZE); containerManager.updateContainerState(cid, HddsProtos.LifeCycleEvent.QUASI_CLOSE); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestNode2PipelineMap.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestNode2PipelineMap.java index f9f120f11ef..f6f4031ceec 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestNode2PipelineMap.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestNode2PipelineMap.java @@ -24,6 +24,7 @@ import java.util.List; import java.util.Set; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; @@ -53,7 +54,7 @@ public void init() throws Exception { pipelineManager = scm.getPipelineManager(); ContainerInfo containerInfo = containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "testOwner"); + ReplicationFactor.THREE), "testOwner", StorageTier.getDefaultTier()); ratisContainer = new ContainerWithPipeline(containerInfo, pipelineManager.getPipeline(containerInfo.getPipelineID())); pipelineManager = scm.getPipelineManager(); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineClose.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineClose.java index c45808504ef..b22ad283d9e 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineClose.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineClose.java @@ -36,6 +36,7 @@ import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -106,7 +107,7 @@ public void init() throws Exception { void createContainer() throws IOException { ContainerInfo containerInfo = containerManager .allocateContainer(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "testOwner"); + ReplicationFactor.THREE), "testOwner", StorageTier.getDefaultTier()); ratisContainer = new ContainerWithPipeline(containerInfo, pipelineManager.getPipeline(containerInfo.getPipelineID())); // At this stage, there should be 2 pipeline one with 1 open container each. @@ -219,7 +220,7 @@ void testPipelineCloseWithLogFailure() ContainerInfo containerInfo = containerManager .allocateContainer(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "testOwner"); + ReplicationFactor.THREE), "testOwner", StorageTier.getDefaultTier()); ContainerWithPipeline containerWithPipeline = new ContainerWithPipeline(containerInfo, pipelineManager.getPipeline(containerInfo.getPipelineID())); @@ -257,7 +258,7 @@ void testPipelineCloseWithLogFailure() void testPipelineCloseTriggersSkippedWhenAlreadyInProgress() throws Exception { ContainerInfo allocateContainer = containerManager .allocateContainer(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "newTestOwner"); + ReplicationFactor.THREE), "newTestOwner", StorageTier.getDefaultTier()); ContainerWithPipeline containerWithPipeline = new ContainerWithPipeline(allocateContainer, pipelineManager.getPipeline(allocateContainer.getPipelineID())); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSCMRestart.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSCMRestart.java index 5375206bab7..b19d87f2e0c 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSCMRestart.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSCMRestart.java @@ -26,6 +26,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; import org.apache.hadoop.hdds.scm.ScmConfigKeys; @@ -69,12 +70,14 @@ public static void init() throws Exception { ratisPipeline1 = pipelineManager.getPipeline( containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "Owner1").getPipelineID()); + ReplicationFactor.THREE), "Owner1", + StorageTier.getDefaultTier()).getPipelineID()); pipelineManager.openPipeline(ratisPipeline1.getId()); ratisPipeline2 = pipelineManager.getPipeline( containerManager.allocateContainer( RatisReplicationConfig.getInstance( - ReplicationFactor.ONE), "Owner2").getPipelineID()); + ReplicationFactor.ONE), "Owner2", + StorageTier.getDefaultTier()).getPipelineID()); pipelineManager.openPipeline(ratisPipeline2.getId()); // At this stage, there should be 2 pipeline one with 1 open container // each. Try restarting the SCM and then discover that pipeline are in @@ -113,7 +116,7 @@ public void testPipelineWithScmRestart() // as was before restart ContainerInfo containerInfo = newContainerManager .allocateContainer(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE), "Owner1"); + ReplicationFactor.THREE), "Owner1", StorageTier.getDefaultTier()); assertEquals(ratisPipeline1.getId(), containerInfo.getPipelineID()); } } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java index 878b57b6faa..e6e3453f3e6 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java @@ -112,6 +112,7 @@ import org.apache.hadoop.security.token.Token; import org.apache.hadoop.security.token.TokenIdentifier; import org.apache.ozone.test.GenericTestUtils; +import org.apache.ozone.test.tag.Unhealthy; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; @@ -444,6 +445,7 @@ public void testListBlock() throws Exception { } @Test + @Unhealthy("Need EC Pipeline to set supported StorageTier value") public void testCreateRecoveryContainer() throws Exception { try (XceiverClientManager xceiverClientManager = new XceiverClientManager(config)) { @@ -452,7 +454,7 @@ public void testCreateRecoveryContainer() throws Exception { scm.getPipelineManager().createPipeline(replicationConfig, StorageTier.getDefaultTier()); scm.getPipelineManager().activatePipeline(newPipeline.getId()); final ContainerInfo container = scm.getContainerManager() - .allocateContainer(replicationConfig, "test"); + .allocateContainer(replicationConfig, "test", StorageTier.getDefaultTier()); Token<ContainerTokenIdentifier> cToken = containerTokenGenerator .generateToken(ANY_USER, container.containerID()); scm.getContainerManager().getContainerStateManager() @@ -533,6 +535,7 @@ public void testCreateRecoveryContainer() throws Exception { } @Test + @Unhealthy("Need EC Pipeline to set supported StorageTier value") public void testCreateRecoveryContainerAfterDNRestart() throws Exception { try (XceiverClientManager xceiverClientManager = new XceiverClientManager(config)) { @@ -541,7 +544,7 @@ public void testCreateRecoveryContainerAfterDNRestart() throws Exception { scm.getPipelineManager().createPipeline(replicationConfig, StorageTier.getDefaultTier()); scm.getPipelineManager().activatePipeline(newPipeline.getId()); final ContainerInfo container = scm.getContainerManager() - .allocateContainer(replicationConfig, "test"); + .allocateContainer(replicationConfig, "test", StorageTier.getDefaultTier()); Token<ContainerTokenIdentifier> cToken = containerTokenGenerator .generateToken(ANY_USER, container.containerID()); scm.getContainerManager().getContainerStateManager() diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java index 19417f35bbe..4f80da69371 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java @@ -230,7 +230,7 @@ private void testPostUpgradePipelineCreation() assertEquals(0, scmPipelineManager.getNumberOfContainers(ratisPipeline1.getId())); PipelineID pid = scmContainerManager.allocateContainer(RATIS_THREE, - "Owner1").getPipelineID(); + "Owner1", StorageTier.getDefaultTier()).getPipelineID(); assertEquals(1, scmPipelineManager.getNumberOfContainers(pid)); assertEquals(pid, ratisPipeline1.getId()); } @@ -252,7 +252,7 @@ private void waitForPipelineCreated() throws Exception { private void createTestContainers() throws IOException, TimeoutException { XceiverClientManager xceiverClientManager = new XceiverClientManager(conf); ContainerInfo ci1 = scmContainerManager.allocateContainer( - RATIS_THREE, "Owner1"); + RATIS_THREE, "Owner1", StorageTier.getDefaultTier()); Pipeline ratisPipeline1 = scmPipelineManager.getPipeline(ci1.getPipelineID()); scmPipelineManager.openPipeline(ratisPipeline1.getId()); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestMiniOzoneCluster.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestMiniOzoneCluster.java index a227242d302..9930b7d375e 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestMiniOzoneCluster.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestMiniOzoneCluster.java @@ -23,13 +23,16 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import java.io.File; import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.HashSet; import java.util.List; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; import org.apache.hadoop.hdds.conf.OzoneConfiguration; @@ -279,4 +282,16 @@ public void testMultipleDataDirs() throws Exception { storageVolume.getVolumeUsage().getReservedInBytes())); } + @Test + public void testDatanodeStorageTypeCountMustMatchVolumeCount() { + List<List<StorageType>> storageTypes = Collections.singletonList( + Collections.singletonList(StorageType.DISK)); + + assertThrows(IllegalArgumentException.class, + () -> UniformDatanodesFactory.newBuilder() + .setNumDataVolumes(2) + .setDatanodeStorageType(storageTypes) + .build()); + } + } diff --git a/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/MiniOzoneCluster.java b/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/MiniOzoneCluster.java index 8765c2aaaae..3d770bc155e 100644 --- a/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/MiniOzoneCluster.java +++ b/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/MiniOzoneCluster.java @@ -20,10 +20,12 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.UUID; import java.util.concurrent.TimeoutException; import org.apache.hadoop.fs.CommonConfigurationKeysPublic; +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; @@ -268,6 +270,8 @@ abstract class Builder { protected String[] racks; protected String[] hosts; private final List<Service> services = new ArrayList<>(); + protected int numDataVolumes = 1; + protected List<List<StorageType>> datanodeStorageType = Collections.emptyList(); protected Builder(OzoneConfiguration conf) { this.conf = conf; @@ -395,6 +399,18 @@ public Builder setRacks(String[] racks) { return this; } + /** + * Sets the number of data volumes per datanode. Rebuilds the default + * {@link UniformDatanodesFactory} to honor the new count. If a custom + * {@link DatanodeFactory} was set via {@link #setDatanodeFactory(DatanodeFactory)}, + * this method has no effect on it. + */ + public Builder setNumDataVolumes(int val) { + this.numDataVolumes = val; + rebuildDefaultDatanodeFactory(); + return this; + } + /** * Sets the hostname for each datanode. When used together with * {@link #setRacks}, the hostnames are used as keys in the @@ -410,6 +426,25 @@ public Builder setHosts(String[] hosts) { return this; } + /** + * Per-datanode storage type list. Outer list size must equal number of datanodes; + * each inner list, when non-empty, must equal {@link #numDataVolumes}. When set, + * each data dir is advertised with the requested {@link StorageType}. + */ + public Builder setDatanodeStorageType(List<List<StorageType>> datanodeStorageType) { + this.datanodeStorageType = datanodeStorageType == null + ? Collections.emptyList() : datanodeStorageType; + rebuildDefaultDatanodeFactory(); + return this; + } + + private void rebuildDefaultDatanodeFactory() { + this.dnFactory = UniformDatanodesFactory.newBuilder() + .setNumDataVolumes(numDataVolumes) + .setDatanodeStorageType(datanodeStorageType) + .build(); + } + protected void validateDatanodeConfiguration() { if (racks != null && racks.length != numOfDatanodes) { throw new IllegalArgumentException("Number of racks must match the number of datanodes"); diff --git a/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/UniformDatanodesFactory.java b/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/UniformDatanodesFactory.java index b4a22c67dff..0a4249525cf 100644 --- a/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/UniformDatanodesFactory.java +++ b/hadoop-ozone/mini-cluster/src/main/java/org/apache/hadoop/ozone/UniformDatanodesFactory.java @@ -42,10 +42,12 @@ import java.nio.file.Path; import java.nio.file.Paths; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Objects; import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.DatanodeVersion; import org.apache.hadoop.hdds.conf.ConfigurationTarget; import org.apache.hadoop.hdds.conf.OzoneConfiguration; @@ -64,6 +66,7 @@ public class UniformDatanodesFactory implements MiniOzoneCluster.DatanodeFactory private final Integer layoutVersion; private final DatanodeVersion initialVersion; private final DatanodeVersion currentVersion; + private final List<List<StorageType>> datanodeStorageType; protected UniformDatanodesFactory(Builder builder) { numDataVolumes = builder.numDataVolumes; @@ -71,6 +74,7 @@ protected UniformDatanodesFactory(Builder builder) { reservedSpace = builder.reservedSpace; currentVersion = builder.currentVersion; initialVersion = builder.initialVersion != null ? builder.initialVersion : builder.currentVersion; + datanodeStorageType = builder.datanodeStorageType; } @Override @@ -86,12 +90,31 @@ public OzoneConfiguration apply(OzoneConfiguration conf) throws IOException { Files.createDirectories(metaDir); dnConf.set(OZONE_METADATA_DIRS, metaDir.toString()); + // Look up this datanode's per-volume StorageType list, if configured. `i` is 1-based. + List<StorageType> volumeStorageTypes = Collections.emptyList(); + if (datanodeStorageType != null && !datanodeStorageType.isEmpty()) { + if (i - 1 >= datanodeStorageType.size()) { + throw new IOException("Datanode index " + (i - 1) + + " has no entry in datanodeStorageType list (size=" + + datanodeStorageType.size() + ")."); + } + volumeStorageTypes = datanodeStorageType.get(i - 1); + if (!volumeStorageTypes.isEmpty() && volumeStorageTypes.size() != numDataVolumes) { + throw new IOException("Datanode " + (i - 1) + " storageType list size " + + volumeStorageTypes.size() + " must equal numDataVolumes " + numDataVolumes + "."); + } + } + List<String> dataDirs = new ArrayList<>(); List<String> reservedSpaceList = new ArrayList<>(); for (int j = 0; j < numDataVolumes; j++) { Path dir = baseDir.resolve("data-" + j); Files.createDirectories(dir); - dataDirs.add(dir.toString()); + if (!volumeStorageTypes.isEmpty()) { + dataDirs.add("[" + volumeStorageTypes.get(j) + "]" + dir); + } else { + dataDirs.add(dir.toString()); + } if (reservedSpace != null) { reservedSpaceList.add(dir + ":" + reservedSpace); } @@ -149,6 +172,7 @@ public static class Builder { private Integer layoutVersion; private DatanodeVersion initialVersion; private DatanodeVersion currentVersion; + private List<List<StorageType>> datanodeStorageType = Collections.emptyList(); /** * Sets the number of data volumes per datanode. @@ -186,7 +210,30 @@ public Builder setCurrentVersion(DatanodeVersion version) { return this; } + /** + * Per-datanode storage type list. Outer list is indexed by datanode; inner list + * is indexed by volume within a datanode. Each inner list, when non-empty, must + * have size == numDataVolumes. When set, each data dir is prefixed with + * {@code "[StorageType]"} so the DN advertises the requested type. + */ + public Builder setDatanodeStorageType(List<List<StorageType>> datanodeStorageType) { + this.datanodeStorageType = datanodeStorageType == null + ? Collections.emptyList() : datanodeStorageType; + return this; + } + public UniformDatanodesFactory build() { + for (int i = 0; i < datanodeStorageType.size(); i++) { + List<StorageType> storageTypes = Objects.requireNonNull( + datanodeStorageType.get(i), + "Datanode storageType list cannot be null"); + if (!storageTypes.isEmpty() + && storageTypes.size() != numDataVolumes) { + throw new IllegalArgumentException("Datanode " + i + + " storageType list size " + storageTypes.size() + + " must equal numDataVolumes " + numDataVolumes + "."); + } + } return new UniformDatanodesFactory(this); } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
