This is an automated email from the ASF dual-hosted git repository.
ivandika3 pushed a commit to branch HDDS-11233
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/HDDS-11233 by this push:
new ed07a4bfb69 HDDS-15408. SCM ContainerInfo Support Create And Get With
StorageTier (#10810)
ed07a4bfb69 is described below
commit ed07a4bfb6975257196f40055add066621a42e43
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 | 49 ++--
.../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, 818 insertions(+), 189 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 acefb5f0b38..17e7d5a646e 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;
@@ -89,6 +91,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);
@@ -105,6 +108,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() {
@@ -130,6 +134,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;
@@ -484,6 +503,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 f38eceb52ad..a0873ceb2b9 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.TestClock;
@@ -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 56e2ef6408f..e66bce755d0 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 933bb4a0087..4ae3a49ba19 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 81b1e89089e..97ab2b56b30 100644
--- a/hadoop-hdds/interface-client/src/main/proto/hdds.proto
+++ b/hadoop-hdds/interface-client/src/main/proto/hdds.proto
@@ -275,6 +275,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 691a1965ab8..750419a2a4f 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,11 +19,11 @@
import jakarta.annotation.Nullable;
import java.io.IOException;
-import java.util.Collections;
import java.util.List;
import java.util.Map;
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;
@@ -137,7 +137,7 @@ List<ContainerInfo> getContainers(ContainerID startID,
* @throws IOException
*/
ContainerInfo allocateContainer(ReplicationConfig replicationConfig,
- String owner)
+ String owner, StorageTier storageTier)
throws IOException;
/**
@@ -197,24 +197,21 @@ void removeContainerReplica(ContainerID containerID,
ContainerReplica replica)
void updateDeleteTransactionId(Map<ContainerID, Long> deleteTransactionMap)
throws IOException;
- 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 f9d68e5428f..b1a2cda6206 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
@@ -26,6 +26,7 @@
import java.util.List;
import java.util.Map;
import java.util.NavigableSet;
+import java.util.Objects;
import java.util.Random;
import java.util.Set;
import java.util.concurrent.locks.Lock;
@@ -171,7 +172,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
@@ -183,10 +185,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();
@@ -195,8 +197,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());
}
@@ -205,20 +206,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();
@@ -228,9 +229,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);
}
@@ -238,7 +240,8 @@ private ContainerInfo createContainer(Pipeline pipeline,
String owner)
}
private ContainerInfo allocateContainer(final Pipeline pipeline,
- final String owner)
+ final String owner,
+ StorageTier storageTier)
throws IOException {
if (!pipelineManager.hasEnoughSpace(pipeline)) {
LOG.debug("Cannot allocate a new container because pipeline {} does not
have enough space.", pipeline);
@@ -249,6 +252,8 @@ private ContainerInfo allocateContainer(final Pipeline
pipeline,
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);
final ContainerInfoProto.Builder containerInfoBuilder = ContainerInfoProto
.newBuilder()
@@ -260,7 +265,8 @@ private ContainerInfo allocateContainer(final Pipeline
pipeline,
.setOwner(owner)
.setContainerID(containerID.getId())
.setDeleteTransactionId(0)
- .setReplicationType(pipeline.getType());
+ .setReplicationType(pipeline.getType())
+ .setStorageTier(storageTier.toProto());
if (pipeline.getReplicationConfig() instanceof ECReplicationConfig) {
containerInfoBuilder.setEcReplicationConfig(
@@ -367,24 +373,25 @@ public void updateDeleteTransactionId(
@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;
}
@@ -395,20 +402,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 fa282396a3e..e3511595493 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
@@ -22,6 +22,7 @@
import java.util.Map;
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;
@@ -204,9 +205,10 @@ void updateDeleteTransactionId(Map<ContainerID, Long>
deleteTransactionMap)
/**
*
*/
- 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 80a9725fac9..ebc8b620e67 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;
@@ -482,8 +484,9 @@ public void updateDeleteTransactionId(
}
@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;
}
@@ -501,7 +504,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
@@ -513,7 +517,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
@@ -523,14 +527,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 12cc3785d28..e81346afd17 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
@@ -22,6 +22,7 @@
import java.util.Map;
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
@@ -267,56 +269,57 @@ public Message invokeLocal(String methodName, Object[] p)
throws Exception {
returnValue = getImpl().getContainerReplicas(arg14);
break;
- case "getMatchingContainer":
+ 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(arg15, arg16, arg17, arg18);
+ returnValue = getImpl().getMatchingContainerAndStorageTier(arg15, arg16,
arg17, arg18, arg19);
break;
case "reinitialize":
- final Table arg19 = p.length > 0 ? (Table) p[0] : null;
- getImpl().reinitialize(arg19);
+ final Table arg20 = p.length > 0 ? (Table) p[0] : null;
+ getImpl().reinitialize(arg20);
return Message.EMPTY;
case "removeContainer":
- final HddsProtos.ContainerID arg20 = p.length > 0 ?
(HddsProtos.ContainerID) p[0] : null;
- getImpl().removeContainer(arg20);
+ final HddsProtos.ContainerID arg21 = p.length > 0 ?
(HddsProtos.ContainerID) p[0] : null;
+ getImpl().removeContainer(arg21);
return Message.EMPTY;
case "removeContainerReplica":
- final ContainerReplica arg21 = p.length > 0 ? (ContainerReplica) p[0] :
null;
- getImpl().removeContainerReplica(arg21);
+ final ContainerReplica arg22 = p.length > 0 ? (ContainerReplica) p[0] :
null;
+ getImpl().removeContainerReplica(arg22);
return Message.EMPTY;
case "transitionDeletingOrDeletedToTargetState":
- final HddsProtos.ContainerID arg22 = p.length > 0 ?
(HddsProtos.ContainerID) p[0] : null;
- final LifeCycleState arg23 = p.length > 1 ? (LifeCycleState) p[1] : null;
- getImpl().transitionDeletingOrDeletedToTargetState(arg22, arg23);
+ final HddsProtos.ContainerID arg23 = p.length > 0 ?
(HddsProtos.ContainerID) p[0] : null;
+ final LifeCycleState arg24 = p.length > 1 ? (LifeCycleState) p[1] : null;
+ getImpl().transitionDeletingOrDeletedToTargetState(arg23, arg24);
return Message.EMPTY;
case "updateContainerInfo":
- final ContainerInfoProto arg24 = p.length > 0 ? (ContainerInfoProto)
p[0] : null;
- getImpl().updateContainerInfo(arg24);
+ final ContainerInfoProto arg25 = p.length > 0 ? (ContainerInfoProto)
p[0] : null;
+ getImpl().updateContainerInfo(arg25);
return Message.EMPTY;
case "updateContainerReplica":
- final ContainerReplica arg25 = p.length > 0 ? (ContainerReplica) p[0] :
null;
- getImpl().updateContainerReplica(arg25);
+ final ContainerReplica arg26 = p.length > 0 ? (ContainerReplica) p[0] :
null;
+ getImpl().updateContainerReplica(arg26);
return Message.EMPTY;
case "updateContainerStateWithSequenceId":
- final HddsProtos.ContainerID arg26 = p.length > 0 ?
(HddsProtos.ContainerID) p[0] : null;
- final LifeCycleEvent arg27 = p.length > 1 ? (LifeCycleEvent) p[1] : null;
- final Long arg28 = p.length > 2 ? (Long) p[2] : null;
- getImpl().updateContainerStateWithSequenceId(arg26, arg27, arg28);
+ final HddsProtos.ContainerID arg27 = p.length > 0 ?
(HddsProtos.ContainerID) p[0] : null;
+ final LifeCycleEvent arg28 = p.length > 1 ? (LifeCycleEvent) p[1] : null;
+ final Long arg29 = p.length > 2 ? (Long) p[2] : null;
+ getImpl().updateContainerStateWithSequenceId(arg27, arg28, arg29);
return Message.EMPTY;
case "updateDeleteTransactionId":
- final Map arg29 = p.length > 0 ? (Map) p[0] : null;
- getImpl().updateDeleteTransactionId(arg29);
+ final Map arg30 = p.length > 0 ? (Map) p[0] : null;
+ getImpl().updateDeleteTransactionId(arg30);
return Message.EMPTY;
default:
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 38c00ada778..9118b4d6290 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
@@ -250,15 +250,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 411679010b5..8107d3e6b3a 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
@@ -59,7 +59,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 9f3578512d5..9a6dcf0ad90 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 73bf92e9cd5..a02661eb066 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 85990a10761..74cc932231b 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);
final Pipeline pipeline = scm.getPipelineManager()
.getPipeline(container.getPipelineID());
ContainerWithPipeline cp = new ContainerWithPipeline(container,
pipeline);
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 69f7973ff1b..61bfd0701c1 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;
@@ -337,6 +341,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
@@ -458,6 +471,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 a0f5a14725f..233dc077208 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 003431416e2..9e93d8db9c1 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;
@@ -96,7 +97,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);
@@ -127,7 +128,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());
@@ -148,14 +149,15 @@ public void
testGetMatchingContainerReturnsNullWhenNotEnoughSpaceInDatanodes() t
// MockPipelineManager#hasEnoughSpace 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);
}
@@ -178,14 +180,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);
}
@@ -193,7 +196,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,
@@ -215,8 +218,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());
@@ -264,14 +269,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, () ->
@@ -285,7 +292,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());
}
@@ -343,7 +350,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()));
@@ -354,7 +361,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()
@@ -374,7 +381,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 ea9f8ae2f0a..3093996fe36 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.apache.hadoop.hdds.scm.HddsTestUtils.getECContainer;
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.Mockito.any;
@@ -36,7 +37,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;
@@ -213,6 +216,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"})
@@ -455,19 +481,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 e2a57cb924e..3bb8b9c1de4 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
@@ -231,7 +231,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 e740f9719e9..e20b457fab4 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 b8677ce3e16..95853b21fc7 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 327bcd6bed0..53a53944e11 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
@@ -955,7 +955,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 07e367b172f..fa1333863d1 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 2ac103e8483..9279c688eff 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 9ad018c0ed6..10ec5973634 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 bde82b2ab95..b9d178b910f 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;
@@ -108,7 +109,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 7d67358e785..acc197b6611 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;
@@ -103,7 +105,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());
@@ -131,13 +133,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());
@@ -209,7 +211,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();
}
@@ -218,7 +220,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());
}
@@ -240,7 +242,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 71d8e39338a..7cda4e8becf 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 14f4c849180..22571e7bcd6 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 d5441276c8b..20e896cfc11 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
@@ -25,6 +25,7 @@
import java.util.Set;
import java.util.concurrent.TimeoutException;
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;
@@ -55,7 +56,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 71666a48c23..856492a047c 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;
@@ -107,7 +108,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.
@@ -221,7 +222,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()));
@@ -259,7 +260,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 31bd6036e74..8c351fa0419 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
@@ -113,6 +113,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;
@@ -445,6 +446,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)) {
@@ -453,7 +455,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()
@@ -534,6 +536,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)) {
@@ -542,7 +545,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 e77b59c8679..4afaa8a7399 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 dbeeda5cdbf..84b10e33fdf 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
@@ -19,9 +19,11 @@
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeoutException;
+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;
@@ -261,6 +263,8 @@ abstract class Builder {
protected SecretKeyClient secretKeyClient;
protected DatanodeFactory dnFactory =
UniformDatanodesFactory.newBuilder().build();
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;
@@ -366,6 +370,37 @@ public Builder setDatanodeFactory(DatanodeFactory factory)
{
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;
+ }
+
+ /**
+ * 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();
+ }
+
public Builder addService(Service service) {
services.add(service);
return this;
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 7e0a1003115..7460a34c557 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
@@ -41,10 +41,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;
@@ -63,6 +65,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;
@@ -70,6 +73,7 @@ protected UniformDatanodesFactory(Builder builder) {
reservedSpace = builder.reservedSpace;
currentVersion = builder.currentVersion;
initialVersion = builder.initialVersion != null ? builder.initialVersion :
builder.currentVersion;
+ datanodeStorageType = builder.datanodeStorageType;
}
@Override
@@ -85,12 +89,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);
}
@@ -147,6 +170,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.
@@ -184,7 +208,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]