This is an automated email from the ASF dual-hosted git repository. devmadhuu pushed a commit to branch HDDS-11233 in repository https://gitbox.apache.org/repos/asf/ozone.git
commit da403ba156bca34a83a4e04a1f10841f4e35c7fc Author: XiChen <[email protected]> AuthorDate: Mon Jul 13 09:21:07 2026 +0800 HDDS-15409. SCM Pipeline Add SupportedStorageTier Field (#10573) --- .../org/apache/hadoop/hdds/client/StorageTier.java | 7 + .../apache/hadoop/hdds/client/StorageTierUtil.java | 43 +++--- .../apache/hadoop/hdds/scm/pipeline/Pipeline.java | 44 ++++++ .../hadoop/hdds/client/StorageTierUtilTest.java | 151 +++++++++++++++++++++ .../common/impl/StorageLocationReport.java | 25 ---- .../interface-client/src/main/proto/hdds.proto | 1 + .../org/apache/hadoop/hdds/scm/node/NodeUtils.java | 75 ++++++++++ .../hdds/scm/pipeline/ECPipelineProvider.java | 20 ++- .../hadoop/hdds/scm/pipeline/PipelineFactory.java | 5 +- .../hadoop/hdds/scm/pipeline/PipelineManager.java | 6 +- .../hdds/scm/pipeline/PipelineManagerImpl.java | 8 +- .../hdds/scm/pipeline/PipelinePlacementPolicy.java | 26 +++- .../hadoop/hdds/scm/pipeline/PipelineProvider.java | 6 +- .../hdds/scm/pipeline/RatisPipelineProvider.java | 30 +++- .../hdds/scm/pipeline/SimplePipelineProvider.java | 29 +++- .../hadoop/hdds/scm/container/MockNodeManager.java | 16 ++- .../hadoop/hdds/scm/node/TestSCMNodeManager.java | 2 +- .../hdds/scm/pipeline/MockPipelineManager.java | 7 +- .../scm/pipeline/MockRatisPipelineProvider.java | 6 +- .../hdds/scm/pipeline/TestECPipelineProvider.java | 15 ++ .../hdds/scm/pipeline/TestPipelineManagerImpl.java | 78 ++++++++++- .../scm/pipeline/TestPipelinePlacementPolicy.java | 20 ++- .../scm/pipeline/TestRatisPipelineProvider.java | 138 ++++++++++++++++--- .../scm/pipeline/TestSimplePipelineProvider.java | 69 +++++++++- .../hadoop/hdds/protocol/StorageTierTest.java | 3 +- .../hadoop/hdds/scm/pipeline/TestSCMRestart.java | 4 + .../apache/hadoop/ozone/om/TestKeyManagerImpl.java | 3 +- .../ozone/recon/scm/ReconPipelineFactory.java | 2 +- 28 files changed, 718 insertions(+), 121 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 1e5c40dfe64..24b783f1b14 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 @@ -104,6 +104,13 @@ public boolean isUniform() { return isUniform; } + public StorageType getUniformStorageType() { + if (!isUniform()) { + throw new IllegalArgumentException("Uniform storage type is not supported"); + } + return storageTypes.get(0); + } + /** * Maps a StorageTier to its corresponding StorageType based on replication type. * diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTierUtil.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTierUtil.java index d68c71a31da..f7ff2afe2b7 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTierUtil.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTierUtil.java @@ -17,7 +17,9 @@ package org.apache.hadoop.hdds.client; +import java.util.ArrayList; import java.util.List; +import java.util.Set; import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.scm.exceptions.SCMException; @@ -45,22 +47,31 @@ public static void validateNotEmpty(StorageTier storageTier) throws SCMException } } - /** - * Returns the StorageType for uniform StorageTier. - * - * @param storageTier the StorageTier to get StorageType from - * @return The uniform StorageTier corresponding StorageType - * @throws SCMException if the StorageTier is non-uniform or the EMPTY StorageTier - */ - public static StorageType getStorageTypeForUniformStorageTier(StorageTier storageTier, ReplicationConfig config) - throws SCMException { - validateNotEmpty(storageTier); - List<StorageType> storageTypes = storageTier.getStorageTypes(config.getRequiredNodes()); - if (storageTier.isUniform()) { - return storageTypes.get(0); - } else { - throw new SCMException("Unsupported non-uniform storage tier " + storageTier, - SCMException.ResultCodes.UNSUPPORTED_NON_UNIFORM_STORAGE_TIER); + public static List<StorageTier> findSupportedStorageTiers( + List<Set<StorageType>> dnStorageTypes) { + List<StorageTier> supportedStorageTiers = new ArrayList<>(); + if (dnStorageTypes.isEmpty()) { + return supportedStorageTiers; + } + // We only support uniform storage tiers currently + for (StorageTier storageTier : StorageTier.values()) { + if (storageTier.equals(StorageTier.EMPTY)) { + continue; + } + if (!storageTier.isUniform()) { + throw new UnsupportedOperationException(storageTier + " is not a uniform storage tier"); + } + boolean supportedTier = true; + for (Set<StorageType> dnStorageType : dnStorageTypes) { + if (!dnStorageType.contains(storageTier.getUniformStorageType())) { + supportedTier = false; + break; + } + } + if (supportedTier) { + supportedStorageTiers.add(storageTier); + } } + return supportedStorageTiers; } } diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java index 1379a9c52aa..a0fbed13ad2 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java @@ -33,6 +33,8 @@ import java.util.Objects; import java.util.Optional; import java.util.Set; +import java.util.concurrent.locks.ReadWriteLock; +import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.function.Function; import java.util.stream.Collectors; import org.apache.commons.lang3.StringUtils; @@ -42,6 +44,7 @@ import org.apache.hadoop.hdds.client.ReplicatedReplicationConfig; 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.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -87,6 +90,10 @@ public final class Pipeline { // suggested leader id with high priority private final DatanodeID suggestedLeaderId; + private final StorageTier supportedStorageTier; + + private final ReadWriteLock lock; + private final Instant stateEnterTime; /** @@ -114,6 +121,8 @@ private Pipeline(Builder b) { replicaIndexes = b.replicaIndexes; creationTimestamp = b.creationTimestamp != null ? b.creationTimestamp : Instant.now(); stateEnterTime = Instant.now(); + supportedStorageTier = b.supportedStorageTier; + lock = new ReentrantReadWriteLock(); } public static Codec<Pipeline> getCodec() { @@ -170,6 +179,20 @@ public DatanodeID getSuggestedLeaderId() { return suggestedLeaderId; } + /** + * Return the Pipeline supported StorageTier. + * + * @return Supported StorageTier + */ + public StorageTier getSupportedStorageTier() { + lock.readLock().lock(); + try { + return supportedStorageTier; + } finally { + lock.readLock().unlock(); + } + } + /** * Set the creation timestamp. Only for protobuf now. */ @@ -399,6 +422,12 @@ public HddsProtos.Pipeline getProtobufMessage(int clientVersion, Set<DatanodeDet builder.setSuggestedLeaderDatanodeID(suggestedLeaderId.toProto()); } + // To simplify Pipeline management, we currently only support Pipelines with a single StorageTier. + // However, to facilitate future expansion, StorageTier is still defined using `repeat` in the proto file. + if (supportedStorageTier != null) { + builder.addSupportedStorageTier(supportedStorageTier.toProto()); + } + // To save the message size on wire, only transfer the node order based on // network topology if (!nodesInOrder.isEmpty()) { @@ -488,6 +517,12 @@ public static Builder toBuilder(HddsProtos.Pipeline pipeline) { HddsProtos.UUID uuid = pipeline.getSuggestedLeaderID(); suggestedLeaderId = DatanodeID.of(uuid); } + StorageTier supportedStorageTier = null; + // To simplify Pipeline management, we currently only support Pipelines with a single StorageTier. + // However, to facilitate future expansion, StorageTier is still defined using `repeat` in the proto file. + for (HddsProtos.StorageTierProto storageTierProto : pipeline.getSupportedStorageTierList()) { + supportedStorageTier = StorageTier.fromProto(storageTierProto); + } final ReplicationConfig config = ReplicationConfig .fromProto(pipeline.getType(), pipeline.getFactor(), @@ -501,6 +536,7 @@ public static Builder toBuilder(HddsProtos.Pipeline pipeline) { .setLeaderId(leaderId) .setSuggestedLeaderId(suggestedLeaderId) .setNodeOrder(pipeline.getMemberOrdersList()) + .setSupportedStorageTier(supportedStorageTier) .setCreateTimestamp(pipeline.getCreationTimeStamp()); } @@ -549,6 +585,7 @@ public String toString() { b.append(']') .append(", ReplicationConfig: ").append(replicationConfig) .append(", State:").append(getPipelineState()) + .append(", SupportStorageTier:").append(getSupportedStorageTier()) .append(", leaderId:").append(leaderId != null ? leaderId.toString() : "") .append(", CreationTimestamp").append(getCreationTimestamp().atZone(ZoneId.systemDefault())) .append('}'); @@ -573,6 +610,7 @@ public static final class Builder { private Instant creationTimestamp = null; private DatanodeID suggestedLeaderId = null; private Map<DatanodeDetails, Integer> replicaIndexes = ImmutableMap.of(); + private StorageTier supportedStorageTier; private Builder() { } @@ -595,6 +633,7 @@ private Builder(Pipeline pipeline) { } replicaIndexes = b.build(); } + this.supportedStorageTier = pipeline.getSupportedStorageTier(); } public Builder setId(DatanodeID datanodeID) { @@ -677,6 +716,11 @@ public Builder setReplicaIndexes(Map<DatanodeDetails, Integer> indexes) { return this; } + public Builder setSupportedStorageTier(StorageTier supportedStorageTier) { + this.supportedStorageTier = supportedStorageTier; + return this; + } + public Pipeline build() { Objects.requireNonNull(id, "id == null"); Objects.requireNonNull(replicationConfig, "replicationConfig == null"); diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/client/StorageTierUtilTest.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/client/StorageTierUtilTest.java new file mode 100644 index 00000000000..f443d72f390 --- /dev/null +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/client/StorageTierUtilTest.java @@ -0,0 +1,151 @@ +/* + * 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.client; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import org.apache.hadoop.fs.StorageType; +import org.junit.jupiter.api.Test; + +/** + * Provides {@link StorageTierUtil} factory methods for testing. + */ +class StorageTierUtilTest { + @Test + void testFindSupportedStorageTiers() { + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD), + setOf(StorageType.SSD), + setOf(StorageType.SSD)), + 1, setOf(StorageTier.SSD)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD)), + 1, setOf(StorageTier.SSD)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK)), + 2, setOf(StorageTier.SSD, StorageTier.DISK)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD, StorageType.DISK, StorageType.ARCHIVE), + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK)), + 2, setOf(StorageTier.SSD, StorageTier.DISK)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD, StorageType.DISK, StorageType.ARCHIVE), + setOf(StorageType.SSD, StorageType.DISK, StorageType.ARCHIVE), + setOf(StorageType.SSD, StorageType.DISK, StorageType.ARCHIVE)), + 3, setOf(StorageTier.SSD, StorageTier.DISK, StorageTier.ARCHIVE)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD, StorageType.DISK)), + 2, setOf(StorageTier.SSD, StorageTier.DISK)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD)), + 1, setOf(StorageTier.SSD)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.DISK)), + 1, setOf(StorageTier.DISK)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.ARCHIVE)), + 1, setOf(StorageTier.ARCHIVE)); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK), + setOf(StorageType.SSD, StorageType.DISK)), + 2, setOf(StorageTier.SSD, StorageTier.DISK)); + } + + @Test + void testFindInvalidStorageTiers() { + assertStorageTiers(createDnStorageTypes(), 0, new HashSet<>()); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.DISK), + setOf(StorageType.SSD), + setOf(StorageType.SSD)), + 0, new HashSet<>()); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.DISK), + setOf(StorageType.DISK), + setOf(StorageType.SSD)), + 0, new HashSet<>()); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.ARCHIVE), + setOf(StorageType.DISK), + setOf(StorageType.SSD)), + 0, new HashSet<>()); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.RAM_DISK), + setOf(StorageType.SSD), + setOf(StorageType.SSD)), + 0, new HashSet<>()); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.RAM_DISK), + setOf(StorageType.RAM_DISK), + setOf(StorageType.RAM_DISK)), + 0, new HashSet<>()); + + assertStorageTiers(createDnStorageTypes( + setOf(StorageType.PROVIDED), + setOf(StorageType.PROVIDED), + setOf(StorageType.PROVIDED)), + 0, new HashSet<>()); + } + + private void assertStorageTiers(List<Set<StorageType>> dnStorageTypes, + int expectedSize, Set<StorageTier> expectedTiers) { + List<StorageTier> supportedTiers = + StorageTierUtil.findSupportedStorageTiers(dnStorageTypes); + HashSet<StorageTier> tiers = new HashSet<>(supportedTiers); + assertEquals(expectedSize, tiers.size()); + assertEquals(expectedTiers, tiers); + } + + private List<Set<StorageType>> createDnStorageTypes(Set<StorageType>... storageSets) { + return Arrays.asList(storageSets); + } + + private Set<StorageType> setOf(StorageType... types) { + return new HashSet<>(Arrays.asList(types)); + } + + private Set<StorageTier> setOf(StorageTier... tiers) { + return new HashSet<>(Arrays.asList(tiers)); + } +} diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/StorageLocationReport.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/StorageLocationReport.java index 75c2f973222..5bfcc78e9e6 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/StorageLocationReport.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/StorageLocationReport.java @@ -132,31 +132,6 @@ public long getFsAvailable() { return fsAvailable; } - public static StorageType getStorageType(StorageTypeProto proto) throws - IllegalArgumentException { - StorageType storageType; - switch (proto) { - case SSD: - storageType = StorageType.SSD; - break; - case DISK: - storageType = StorageType.DISK; - break; - case ARCHIVE: - storageType = StorageType.ARCHIVE; - break; - case PROVIDED: - storageType = StorageType.PROVIDED; - break; - case RAM_DISK: - storageType = StorageType.RAM_DISK; - break; - default: - throw new IllegalArgumentException("Illegal Storage Type specified: " + proto); - } - return storageType; - } - /** * Returns the StorageReportProto protoBuf message for the Storage Location * report. diff --git a/hadoop-hdds/interface-client/src/main/proto/hdds.proto b/hadoop-hdds/interface-client/src/main/proto/hdds.proto index 67eab258084..f14c888c644 100644 --- a/hadoop-hdds/interface-client/src/main/proto/hdds.proto +++ b/hadoop-hdds/interface-client/src/main/proto/hdds.proto @@ -143,6 +143,7 @@ message Pipeline { optional DatanodeIDProto leaderDatanodeID = 101; optional DatanodeIDProto suggestedLeaderDatanodeID = 102; + repeated StorageTierProto supportedStorageTier = 103; } message KeyValue { diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeUtils.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeUtils.java new file mode 100644 index 00000000000..f42a6e0f5a7 --- /dev/null +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeUtils.java @@ -0,0 +1,75 @@ +/* + * 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.node; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTier; +import org.apache.hadoop.hdds.client.StorageTierUtil; +import org.apache.hadoop.hdds.client.StorageTypeUtils; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.StorageReportProto; + +/** + * Util class for Node operations. + */ +public final class NodeUtils { + + private NodeUtils() { + } + + public static List<StorageTier> getDatanodesStorageTypes( + List<DatanodeDetails> dns, NodeManager nodeManager) { + List<Set<StorageType>> dnStorageTypes = new ArrayList<>(); + for (DatanodeDetails dn : dns) { + DatanodeInfo datanodeInfo = nodeManager.getDatanodeInfo(dn); + if (datanodeInfo == null) { + throw new IllegalStateException("Cannot get Datanode : " + dn.getUuidString() + " Info"); + } + dnStorageTypes.add(getDatanodeStorageTypes(datanodeInfo)); + } + return StorageTierUtil.findSupportedStorageTiers(dnStorageTypes); + } + + public static Set<StorageType> getDatanodeStorageTypes( + DatanodeInfo datanodeInfo) { + List<StorageReportProto> storageReportProtos = + datanodeInfo.getStorageReports(); + if (storageReportProtos == null || storageReportProtos.isEmpty()) { + return Collections.emptySet(); + } + + Set<StorageType> uniqueStorageTypes = new HashSet<>(); + for (StorageReportProto storageReportProto : storageReportProtos) { + uniqueStorageTypes.add(StorageTypeUtils.getFromProtobuf(( + storageReportProto.getStorageType()))); + } + return uniqueStorageTypes; + } + + public static StorageTypeProto getStorageTypeFromStorageReportProto( + StorageReportProto storageReportProto) { + return storageReportProto.hasStorageType() + ? storageReportProto.getStorageType() : StorageTypeProto.DISK; + } +} diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/ECPipelineProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/ECPipelineProvider.java index dc9618a1cb4..94380948200 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/ECPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/ECPipelineProvider.java @@ -28,15 +28,16 @@ import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.StorageTier; -import org.apache.hadoop.hdds.client.StorageTierUtil; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.scm.PlacementPolicy; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.exceptions.SCMException; import org.apache.hadoop.hdds.scm.node.NodeManager; import org.apache.hadoop.hdds.scm.node.NodeStatus; +import org.apache.hadoop.hdds.scm.node.NodeUtils; import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -90,17 +91,23 @@ protected Pipeline create(ECReplicationConfig replicationConfig, List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, StorageTier storageTier) throws IOException { - StorageType storageType = StorageTierUtil.getStorageTypeForUniformStorageTier(storageTier, replicationConfig); + StorageType storageType = storageTier.getUniformStorageType(); List<DatanodeDetails> dns = placementPolicy .chooseDatanodes(excludedNodes, favoredNodes, replicationConfig.getRequiredNodes(), 0, this.containerSizeBytes, storageType); - return create(replicationConfig, dns); + return create(replicationConfig, dns, storageTier); } @Override protected Pipeline create(ECReplicationConfig replicationConfig, - List<DatanodeDetails> nodes) { - + List<DatanodeDetails> nodes, StorageTier storageTier) throws IOException { + List<StorageTier> storageTiers = NodeUtils.getDatanodesStorageTypes(nodes, getNodeManager()); + if (!storageTiers.contains(storageTier)) { + throw new SCMException(String.format("Cannot create pipeline for " + + "StorageTier %s replicationConfig: %s", + storageTier, replicationConfig), + SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE); + } Map<DatanodeDetails, Integer> dnIndexes = new HashMap<>(); int ecIndex = 1; for (DatanodeDetails dn : nodes) { @@ -111,6 +118,7 @@ protected Pipeline create(ECReplicationConfig replicationConfig, return newPipelineBuilder(replicationConfig, nodes) .setId(PipelineID.randomId()) .setReplicaIndexes(dnIndexes) + .setSupportedStorageTier(storageTier) .build(); } @@ -140,9 +148,11 @@ public Pipeline createForRead( // Use insecureRandomId for throwaway read pipeline IDs to avoid // contention on the shared SecureRandom instance. + // Read Pipelines do not require storage tiers, so no supported tier is set. return newPipelineBuilder(replicationConfig, dns) .setId(PipelineID.insecureRandomId()) .setReplicaIndexes(map) + .setSupportedStorageTier(null) .build(); } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java index 9da5bac2f28..2268309012b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java @@ -108,10 +108,9 @@ private void checkPipeline(Pipeline pipeline) throws IOException { } public Pipeline create(ReplicationConfig replicationConfig, - List<DatanodeDetails> nodes - ) { + List<DatanodeDetails> nodes, StorageTier storageTier) throws IOException { return providers.get(replicationConfig.getReplicationType()) - .create(replicationConfig, nodes); + .create(replicationConfig, nodes, storageTier); } public Pipeline createForRead(ReplicationConfig replicationConfig, diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java index e0d9de1f10b..77c9fec4056 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java @@ -24,6 +24,7 @@ import java.util.NavigableSet; import java.util.Set; 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.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.ContainerReplica; @@ -53,8 +54,9 @@ Pipeline buildECPipeline(ReplicationConfig replicationConfig, Pipeline createPipeline( ReplicationConfig replicationConfig, - List<DatanodeDetails> nodes - ); + List<DatanodeDetails> nodes, + StorageTier storageTier + ) throws IOException; Pipeline createPipelineForRead( ReplicationConfig replicationConfig, Set<ContainerReplica> replicas); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java index c472e65b8b3..4dcfaad37a5 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java @@ -40,6 +40,7 @@ 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.ConfigurationSource; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -314,11 +315,12 @@ private boolean factorOne(ReplicationConfig replicationConfig) { @Override public Pipeline createPipeline( ReplicationConfig replicationConfig, - List<DatanodeDetails> nodes - ) { + List<DatanodeDetails> nodes, + StorageTier storageTier + ) throws IOException { // This will mostly be used to create dummy pipeline for SimplePipelines. // We don't update the metrics for SimplePipelines. - return pipelineFactory.create(replicationConfig, nodes); + return pipelineFactory.create(replicationConfig, nodes, storageTier); } @Override 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 e2f880121fc..8dbdafea23b 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 @@ -85,7 +85,8 @@ public PipelinePlacementPolicy(final NodeManager nodeManager, static int currentRatisThreePipelineCount( NodeManager nodeManager, PipelineStateManager stateManager, - DatanodeDetails datanodeDetails) { + DatanodeDetails datanodeDetails, + StorageType storageType) { // Safe to cast collection's size to int return (int) nodeManager.getPipelines(datanodeDetails).stream() .map(id -> { @@ -97,15 +98,15 @@ static int currentRatisThreePipelineCount( return null; } }) - .filter(PipelinePlacementPolicy::isNonClosedRatisThreePipeline) + .filter(p -> isNonClosedRatisThreePipeline(p, storageType)) .count(); } /** Filter the given datanodes within its pipeline limit. */ - List<DatanodeDetails> filterPipelineLimit(Iterable<DatanodeDetails> datanodes) { + List<DatanodeDetails> filterPipelineLimit(Iterable<DatanodeDetails> datanodes, StorageType storageType) { final SortedList<DatanodeDetails> sorted = new SortedList<>(DatanodeDetails.class); for (DatanodeDetails d : datanodes) { - final int count = currentRatisThreePipelineCount(nodeManager, stateManager, d); + final int count = currentRatisThreePipelineCount(nodeManager, stateManager, d, storageType); if (count < nodeManager.pipelineLimit(d)) { sorted.add(d, count); } @@ -113,10 +114,21 @@ List<DatanodeDetails> filterPipelineLimit(Iterable<DatanodeDetails> datanodes) { return sorted; } - private static boolean isNonClosedRatisThreePipeline(Pipeline p) { + 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; + } + } return p != null && p.getReplicationConfig() .equals(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)) - && !p.isClosed(); + && !p.isClosed() && matchedTier; } @Override @@ -182,7 +194,7 @@ List<DatanodeDetails> filterViableNodes( // filter nodes that meet the size and pipeline engagement criteria. // Pipeline placement doesn't take node space left into account. // Sort the DNs by pipeline load. - final List<DatanodeDetails> healthyList = filterPipelineLimit(healthyNodes); + final List<DatanodeDetails> healthyList = filterPipelineLimit(healthyNodes, storageType); if (healthyList.size() < nodesRequired) { if (LOG.isDebugEnabled()) { diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java index 983da957ded..a66c293459f 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java @@ -74,8 +74,8 @@ protected abstract Pipeline create(REPLICATION_CONFIG replicationConfig, protected abstract Pipeline create( REPLICATION_CONFIG replicationConfig, - List<DatanodeDetails> nodes - ); + List<DatanodeDetails> nodes, + StorageTier storageTier) throws IOException; protected abstract Pipeline createForRead( REPLICATION_CONFIG replicationConfig, @@ -89,7 +89,7 @@ List<DatanodeDetails> pickNodesNotUsed(REPLICATION_CONFIG replicationConfig, long dataSizeRequired, StorageTier storageTier) throws SCMException { StorageTierUtil.validateNotEmpty(storageTier); - StorageType storageType = StorageTierUtil.getStorageTypeForUniformStorageTier(storageTier, replicationConfig); + StorageType storageType = storageTier.getUniformStorageType(); int nodesRequired = replicationConfig.getRequiredNodes(); List<DatanodeDetails> healthyDNs = pickAllNodesNotUsed(replicationConfig); List<DatanodeDetails> healthyDNsWithSpace = healthyDNs.stream() diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java index e9cccd72b5d..8d558b5eee5 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java @@ -41,6 +41,7 @@ import org.apache.hadoop.hdds.scm.ha.SCMContext; import org.apache.hadoop.hdds.scm.node.NodeManager; import org.apache.hadoop.hdds.scm.node.NodeStatus; +import org.apache.hadoop.hdds.scm.node.NodeUtils; import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState; import org.apache.hadoop.hdds.scm.pipeline.leader.choose.algorithms.LeaderChoosePolicy; import org.apache.hadoop.hdds.scm.pipeline.leader.choose.algorithms.LeaderChoosePolicyFactory; @@ -160,8 +161,8 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, : String.format(": %d", pipelineNumberLimit); throw new SCMException( - String.format("Cannot create pipeline as it would exceed the limit %s replicationConfig: %s", - limitInfo, replicationConfig), + String.format("Cannot create pipeline for StorageTier %s as it would exceed the limit %s " + + "replicationConfig: %s", storageTier, limitInfo, replicationConfig), SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE ); } @@ -185,9 +186,15 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, DatanodeDetails suggestedLeader = leaderChoosePolicy.chooseLeader(dns); + List<StorageTier> storageTiers = NodeUtils.getDatanodesStorageTypes(dns, getNodeManager()); + if (!storageTiers.contains(storageTier)) { + throw new SCMException(String.format("Cannot create pipeline for StorageTier %s replicationConfig: %s", + storageTier, replicationConfig), SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE); + } Pipeline pipeline = newPipelineBuilder(RatisReplicationConfig.getInstance(factor), dns) .setId(PipelineID.randomId()) .setSuggestedLeaderId(suggestedLeader != null ? suggestedLeader.getID() : null) + .setSupportedStorageTier(storageTier) .build(); // Send command to datanodes to create pipeline @@ -200,8 +207,8 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, createCommand.setTerm(scmContext.getTermOfLeader()); dns.forEach(node -> { - LOG.info("Sending CreatePipelineCommand for pipeline:{} to datanode:{}", - pipeline.getId(), node); + LOG.info("Sending CreatePipelineCommand for pipeline:{} to datanode:{} storageTier:{}", + pipeline.getId(), node, storageTier); eventPublisher.fireEvent(SCMEvents.DATANODE_COMMAND, new CommandForDatanode<>(node, createCommand)); }); @@ -211,9 +218,17 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, @Override public Pipeline create(RatisReplicationConfig replicationConfig, - List<DatanodeDetails> nodes) { + List<DatanodeDetails> nodes, StorageTier storageTier) throws IOException { + List<StorageTier> storageTiers = NodeUtils.getDatanodesStorageTypes(nodes, getNodeManager()); + if (!storageTiers.contains(storageTier)) { + throw new SCMException(String.format("Cannot create pipeline for " + + "StorageTier %s replicationConfig: %s", + storageTier, replicationConfig), + SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE); + } return newPipelineBuilder(replicationConfig, nodes) .setId(PipelineID.randomId()) + .setSupportedStorageTier(storageTier) .build(); } @@ -223,8 +238,10 @@ public Pipeline createForRead( Set<ContainerReplica> replicas) { // Use insecureRandomId for throwaway read pipeline IDs to avoid // contention on the shared SecureRandom instance. + // Read Pipelines do not require storage tiers, so no supported tier is set. return newPipelineBuilder(replicationConfig, ContainerReplica.toDatanodeDetailsList(replicas)) .setId(PipelineID.insecureRandomId()) + .setSupportedStorageTier(null) .build(); } @@ -238,7 +255,8 @@ private List<DatanodeDetails> chooseThreeFactorDatanodes( excludedNodes = excludedNodes.isEmpty() ? null : excludedNodes; List<DatanodeDetails> additionalExcludedNodes = null; for (DatanodeDetails d : healthyNodes) { - final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d); + final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount( + nodeManager, stateManager, d, storageType); if (count >= nodeManager.pipelineLimit(d)) { if (excludedNodes == null) { excludedNodes = new ArrayList<>(); 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 1a41b183ed3..9e4d98feb8f 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 @@ -30,8 +30,10 @@ import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.exceptions.SCMException; import org.apache.hadoop.hdds.scm.node.DatanodeInfo; import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.node.NodeUtils; import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState; /** @@ -57,7 +59,7 @@ public Pipeline create(StandaloneReplicationConfig replicationConfig, StorageTie public Pipeline create(StandaloneReplicationConfig replicationConfig, List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, StorageTier storageTier) throws IOException { - StorageType storageType = StorageTierUtil.getStorageTypeForUniformStorageTier(storageTier, replicationConfig); + StorageType storageType = storageTier.getUniformStorageType(); List<DatanodeDetails> dns = pickNodesNotUsed(replicationConfig); dns = dns.stream().filter(dn -> ((DatanodeInfo) dn).getStorageReports().stream() @@ -74,26 +76,43 @@ public Pipeline create(StandaloneReplicationConfig replicationConfig, } Collections.shuffle(dns); - return newPipelineBuilder(replicationConfig, dns.subList(0, replicationConfig.getReplicationFactor().getNumber())) + return newPipelineBuilder(replicationConfig, + dns.subList(0, replicationConfig.getReplicationFactor().getNumber())) .setId(PipelineID.randomId()) + .setSupportedStorageTier(storageTier) .build(); } - @Override - public Pipeline create(StandaloneReplicationConfig replicationConfig, - List<DatanodeDetails> nodes) { + private Pipeline createPipelineInternal(StandaloneReplicationConfig replicationConfig, + List<DatanodeDetails> nodes, StorageTier storageTier) { return newPipelineBuilder(replicationConfig, nodes) .setId(PipelineID.randomId()) + .setSupportedStorageTier(storageTier) .build(); } + @Override + public Pipeline create(StandaloneReplicationConfig replicationConfig, + List<DatanodeDetails> nodes, StorageTier storageTier) throws IOException { + List<StorageTier> storageTiers = NodeUtils.getDatanodesStorageTypes(nodes, getNodeManager()); + if (!storageTiers.contains(storageTier)) { + throw new SCMException(String.format("Cannot create pipeline for " + + "StorageTier %s replicationConfig: %s", + storageTier, replicationConfig), + SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE); + } + return createPipelineInternal(replicationConfig, nodes, storageTier); + } + @Override public Pipeline createForRead(StandaloneReplicationConfig replicationConfig, Set<ContainerReplica> replicas) { // Use insecureRandomId for throwaway read pipeline IDs to avoid // contention on the shared SecureRandom instance. + // Read Pipelines do not require storage tiers, so no supported tier is set. return newPipelineBuilder(replicationConfig, ContainerReplica.toDatanodeDetailsList(replicas)) .setId(PipelineID.insecureRandomId()) + .setSupportedStorageTier(null) .build(); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java index ff0fb103f26..7365a643754 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java @@ -39,6 +39,7 @@ import java.util.concurrent.ConcurrentMap; import java.util.stream.Collectors; import org.apache.hadoop.fs.StorageType; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; @@ -121,6 +122,8 @@ public class MockNodeManager implements NodeManager { private final OzoneConfiguration conf = new OzoneConfiguration(); private StorageType storageType = null; + private Map<DatanodeID, StorageType> nodeStorageTypeMap = new HashMap<>(); + { this.healthyNodes = new LinkedList<>(); this.staleNodes = new LinkedList<>(); @@ -298,10 +301,12 @@ private List<DatanodeDetails> getDatanodeDetails( long capacity = nodeMetricMap.get(dd).getCapacity().get(); long used = nodeMetricMap.get(dd).getScmUsed().get(); long remaining = nodeMetricMap.get(dd).getRemaining().get(); + StorageType nodeStorageType = nodeStorageTypeMap.get(dd.getID()) != null ? + nodeStorageTypeMap.get(dd.getID()) : storageType; StorageReportProto storage1 = HddsTestUtils.createStorageReport( di.getID(), "/data1-" + di.getID(), capacity, used, remaining, - storageType == null ? null : getStorageTypeProto(storageType)); + nodeStorageType == null ? null : getStorageTypeProto(nodeStorageType)); MetadataStorageReportProto metaStorage1 = HddsTestUtils.createMetadataStorageReport( "/metadata1-" + di.getID(), capacity, used, remaining, null); @@ -468,9 +473,12 @@ public DatanodeInfo getDatanodeInfo(DatanodeDetails dd) { long capacity = nodeMetricMap.get(dd).getCapacity().get(); long used = nodeMetricMap.get(dd).getScmUsed().get(); long remaining = nodeMetricMap.get(dd).getRemaining().get(); + StorageType nodeStorageType = nodeStorageTypeMap.get(dd.getID()) != null ? + nodeStorageTypeMap.get(dd.getID()) : storageType; StorageReportProto storage1 = HddsTestUtils.createStorageReport( di.getID(), "/data1-" + di.getUuidString(), - capacity, used, remaining, null); + capacity, used, remaining, + nodeStorageType == null ? null : StorageTypeUtils.getStorageTypeProto(nodeStorageType)); MetadataStorageReportProto metaStorage1 = HddsTestUtils.createMetadataStorageReport( "/metadata1-" + di.getUuidString(), capacity, used, @@ -1074,4 +1082,8 @@ public void setCurrentState(long currentState) { } } + + public void setStorageTypeForNode(DatanodeID uuid, StorageType st) { + this.nodeStorageTypeMap.put(uuid, st); + } } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java index 78814b53e81..f1f755c699b 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java @@ -560,7 +560,7 @@ private void assertPipelineCreationFailsWithExceedingLimit(int limit) { () -> scm.getPipelineManager().createPipeline(config), "3 nodes should not have been found for a pipeline."); assertThat(ex.getMessage()) - .contains("Cannot create pipeline as it would exceed the limit per datanode: " + limit); + .contains("Cannot create pipeline for StorageTier DISK as it would exceed the limit per datanode: " + limit); } private void assertPipelines(HddsProtos.ReplicationFactor factor, 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 ef6f187b040..12bb97bfb5a 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 @@ -30,6 +30,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; 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.MockDatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -83,7 +84,8 @@ public Pipeline createPipeline(ReplicationConfig replicationConfig, pipeline = createPipeline(replicationConfig, ImmutableList.of(MockDatanodeDetails.randomDatanodeDetails(), MockDatanodeDetails.randomDatanodeDetails(), - MockDatanodeDetails.randomDatanodeDetails())); + MockDatanodeDetails.randomDatanodeDetails()), + StorageTier.getDefaultTier()); } stateManager.addPipeline(pipeline.getProtobufMessage( @@ -116,11 +118,12 @@ public void addEcPipeline(Pipeline pipeline) @Override public Pipeline createPipeline(final ReplicationConfig replicationConfig, - final List<DatanodeDetails> nodes) { + final List<DatanodeDetails> nodes, StorageTier storageTier) { return Pipeline.newBuilder() .setId(PipelineID.randomId()) .setReplicationConfig(replicationConfig) .setNodes(nodes) + .setSupportedStorageTier(storageTier) .setState(Pipeline.PipelineState.OPEN) .build(); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockRatisPipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockRatisPipelineProvider.java index 4ad8ee4a302..fbca53d4818 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockRatisPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockRatisPipelineProvider.java @@ -79,6 +79,7 @@ public Pipeline create(RatisReplicationConfig replicationConfig, StorageTier sto .fromProtoTypeAndFactor(initialPipeline.getType(), replicationConfig.getReplicationFactor())) .setNodes(initialPipeline.getNodes()) + .setSupportedStorageTier(initialPipeline.getSupportedStorageTier()) .build(); return pipeline; } @@ -94,12 +95,15 @@ public static void markPipelineHealthy(Pipeline pipeline) @Override public Pipeline create(RatisReplicationConfig replicationConfig, - List<DatanodeDetails> nodes) { + List<DatanodeDetails> nodes, StorageTier storageTier) + throws IOException { return Pipeline.newBuilder() .setId(PipelineID.randomId()) .setState(Pipeline.PipelineState.OPEN) .setReplicationConfig(replicationConfig) .setNodes(nodes) + .setSupportedStorageTier(storageTier) .build(); } + } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestECPipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestECPipelineProvider.java index 92a28c99d46..0657282ca74 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestECPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestECPipelineProvider.java @@ -37,6 +37,7 @@ import com.google.common.collect.ImmutableSet; import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.HashSet; import java.util.Iterator; import java.util.List; @@ -50,10 +51,12 @@ import org.apache.hadoop.hdds.protocol.DatanodeID; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos; +import org.apache.hadoop.hdds.scm.HddsTestUtils; import org.apache.hadoop.hdds.scm.PlacementPolicy; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.node.DatanodeInfo; import org.apache.hadoop.hdds.scm.node.NodeManager; import org.apache.hadoop.hdds.scm.node.NodeStatus; import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException; @@ -95,6 +98,8 @@ public void setup() throws IOException, NodeNotFoundException { when(nodeManager.getNodeStatus(any())) .thenReturn(NodeStatus.inServiceHealthy()); + when(nodeManager.getDatanodeInfo(any())) + .thenAnswer(invocation -> createDatanodeInfo(invocation.getArgument(0))); } @Test @@ -223,4 +228,14 @@ private Set<ContainerReplica> createContainerReplicas(int number) { return replicas; } + private DatanodeInfo createDatanodeInfo(DatanodeDetails dn) { + DatanodeInfo datanodeInfo = new DatanodeInfo(dn, + NodeStatus.inServiceHealthy(), null, + HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT); + datanodeInfo.updateStorageReports(Collections.singletonList( + HddsTestUtils.createStorageReport(dn.getID(), + "/data-" + dn.getUuidString(), 100))); + return datanodeInfo; + } + } 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 e980b81318d..2d6353c92ec 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 @@ -62,9 +62,11 @@ import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +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.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.DatanodeID; @@ -161,6 +163,18 @@ public void cleanup() throws Exception { } } + private PipelineManagerImpl createPipelineManager(NodeManager manager, boolean isLeader) + throws IOException { + return PipelineManagerImpl.newPipelineManager(conf, + SCMHAManagerStub.getInstance(isLeader), + manager, + SCMDBDefinition.PIPELINES.getTable(dbStore), + new EventQueue(), + scmContext, + serviceManager, + testClock); + } + private PipelineManagerImpl createPipelineManager(boolean isLeader) throws IOException { return PipelineManagerImpl.newPipelineManager(conf, @@ -375,7 +389,7 @@ public void testPipelineReport() throws Exception { // get pipeline report from each dn in the pipeline PipelineReportHandler pipelineReportHandler = new PipelineReportHandler(scmSafeModeManager, pipelineManager, - SCMContext.emptyContext(), conf); + scmContext, conf); nodes.subList(0, 2).forEach(dn -> sendPipelineReport(dn, pipeline, pipelineReportHandler, false)); sendPipelineReport(nodes.get(nodes.size() - 1), pipeline, @@ -403,6 +417,60 @@ public void testPipelineReport() throws Exception { } } + @Test + public void testPipelineReportDoesNotUpdateSupportedStorageTier() throws Exception { + // Set Env + int nodeCount = 3; + MockNodeManager localNodeManager = new MockNodeManager(true, nodeCount, StorageType.DISK); + SCMContext localScmContext = spy(SCMContext.emptyContext()); + StorageContainerManager localScm = mock(StorageContainerManager.class); + when(localScm.getScmNodeManager()).thenReturn(localNodeManager); + when(localScmContext.getScm()).thenReturn(localScm); + + PipelineManagerImpl pipelineManager = createPipelineManager(localNodeManager, true); + SCMSafeModeManager scmSafeModeManager = new SCMSafeModeManager(conf, + localNodeManager, pipelineManager, mock(ContainerManager.class), + serviceManager, new EventQueue(), scmContext); + List<DatanodeDetails> nodes = localNodeManager.getNodes(NodeStatus.inServiceHealthy()); + assertEquals(nodeCount, nodes.size()); + Pipeline pipeline = pipelineManager.createPipeline(RatisReplicationConfig + .getInstance(ReplicationFactor.THREE)); + assertEquals(StorageTier.DISK, pipeline.getSupportedStorageTier()); + + assertFalse(pipelineManager.getPipeline(pipeline.getId()).isHealthy()); + PipelineReportHandler pipelineReportHandler = new PipelineReportHandler( + scmSafeModeManager, pipelineManager, localScmContext, conf); + nodes.subList(0, 2).forEach(dn -> sendPipelineReport(dn, pipeline, + pipelineReportHandler, false)); + sendPipelineReport(nodes.get(nodes.size() - 1), pipeline, + pipelineReportHandler, true); + + // All the Datanode Volume StorageType is DISK so the Pipeline StorageTier will be StorageTier.DISK + assertTrue(pipelineManager.getPipeline(pipeline.getId()).isHealthy()); + assertTrue(pipelineManager.getPipeline(pipeline.getId()).isOpen()); + assertEquals(StorageTier.DISK, pipeline.getSupportedStorageTier()); + + // Only the first Datanode updated its NodeReport to SSD, + // but the Pipeline supportedStorageTier keeps the StorageTier selected at creation. + localNodeManager.setStorageTypeForNode(nodes.get(0).getID(), StorageType.SSD); + sendPipelineReport(nodes.get(0), pipeline, pipelineReportHandler, false); + assertEquals(StorageTier.DISK, pipelineManager.getPipeline(pipeline.getId()).getSupportedStorageTier()); + + // The first and second Datanode updated its NodeReport to SSD + localNodeManager.setStorageTypeForNode(nodes.get(1).getID(), StorageType.SSD); + sendPipelineReport(nodes.get(1), pipeline, pipelineReportHandler, false); + assertEquals(StorageTier.DISK, pipelineManager.getPipeline(pipeline.getId()).getSupportedStorageTier()); + + // The first and second Datanode updated its NodeReport to SSD + localNodeManager.setStorageTypeForNode(nodes.get(2).getID(), StorageType.SSD); + sendPipelineReport(nodes.get(2), pipeline, pipelineReportHandler, true); + assertEquals(StorageTier.DISK, pipelineManager.getPipeline(pipeline.getId()).getSupportedStorageTier()); + + // close the pipeline and clean up + pipelineManager.closePipeline(pipeline.getId()); + pipelineManager.close(); + } + @Test public void testPipelineCreationFailedMetric() throws Exception { PipelineManagerImpl pipelineManager = createPipelineManager(true); @@ -476,7 +544,7 @@ public void testPipelineOpenOnlyWhenLeaderReported() throws Exception { serviceManager, new EventQueue(), scmContext); PipelineReportHandler pipelineReportHandler = new PipelineReportHandler(scmSafeModeManager, pipelineManager, - SCMContext.emptyContext(), conf); + scmContext, conf); // Report pipelines with leaders List<DatanodeDetails> nodes = pipeline.getNodes(); @@ -853,7 +921,7 @@ public void testWaitForAllocatedPipeline() throws IOException { = new HealthyPipelineChoosePolicy(); ContainerManager containerManager = mock(ContainerManager.class); - + WritableContainerProvider<ReplicationConfig> provider; String owner = "TEST"; Pipeline allocatedPipeline; @@ -874,7 +942,7 @@ public void testWaitForAllocatedPipeline() throws IOException { ContainerInfo container = HddsTestUtils. getContainer(HddsProtos.LifeCycleState.OPEN, allocatedPipeline.getId()); - + pipelineManager.addContainerToPipeline( allocatedPipeline.getId(), container.containerID()); doReturn(container).when(containerManager).getMatchingContainer(anyLong(), @@ -900,7 +968,7 @@ public void testWaitForAllocatedPipeline() throws IOException { return call.callRealMethod(); }).when(pipelineManagerSpy).waitOnePipelineReady(any(), anyLong()); - + ContainerInfo c = provider.getContainer(1, repConfig, owner, new ExcludeList()); assertEquals(c, container, "Expected container was returned"); 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 d91d608b588..6b395aab0ce 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 @@ -46,6 +46,7 @@ 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.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -114,11 +115,13 @@ public class TestPipelinePlacementPolicy { private List<DatanodeDetails> nodesWithOutRackAwareness = new ArrayList<>(); private List<DatanodeDetails> nodesWithRackAwareness = new ArrayList<>(); - private final StorageType storageType = StorageType.DEFAULT; + private final StorageTier storageTier = StorageTier.DISK; + private StorageType storageType = StorageType.DEFAULT; @BeforeEach public void init() throws Exception { cluster = initTopology(); + storageType = storageTier.getUniformStorageType(); // start with nodes with rack awareness. nodeManager = new MockNodeManager(cluster, getNodesWithRackAwareness(), false, PIPELINE_PLACEMENT_MAX_NODES_COUNT, storageType); @@ -291,6 +294,7 @@ public void testPickLowestLoadAnchor() throws IOException, TimeoutException { .setState(Pipeline.PipelineState.ALLOCATED) .setReplicationConfig(RatisReplicationConfig.getInstance( ReplicationFactor.THREE)) + .setSupportedStorageTier(storageTier) .setNodes(nodes) .build(); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( @@ -647,6 +651,7 @@ private void insertHeavyNodesIntoNodeManager( .setReplicationConfig(ReplicationConfig .fromProtoTypeAndFactor(RATIS, THREE)) .setNodes(dnList) + .setSupportedStorageTier(storageTier) .build(); pipelineProto = pipeline.getProtobufMessage( @@ -674,7 +679,7 @@ public void testCurrentRatisThreePipelineCount() pipelineCount = placementPolicy.currentRatisThreePipelineCount(nodeManager, - stateManager, healthyNodes.get(0)); + stateManager, healthyNodes.get(0), StorageType.DEFAULT); assertEquals(pipelineCount, 0); // Check datanode with one RATIS/ONE pipeline @@ -684,7 +689,7 @@ public void testCurrentRatisThreePipelineCount() pipelineCount = placementPolicy.currentRatisThreePipelineCount(nodeManager, - stateManager, healthyNodes.get(1)); + stateManager, healthyNodes.get(1), StorageType.DEFAULT); assertEquals(pipelineCount, 0); // Check datanode with one RATIS/THREE pipeline @@ -696,7 +701,7 @@ public void testCurrentRatisThreePipelineCount() pipelineCount = placementPolicy.currentRatisThreePipelineCount(nodeManager, - stateManager, healthyNodes.get(2)); + stateManager, healthyNodes.get(2), StorageType.DEFAULT); assertEquals(pipelineCount, 1); // Check datanode with one RATIS/ONE and one STANDALONE/ONE pipeline @@ -706,7 +711,7 @@ public void testCurrentRatisThreePipelineCount() pipelineCount = placementPolicy.currentRatisThreePipelineCount(nodeManager, - stateManager, healthyNodes.get(1)); + stateManager, healthyNodes.get(1), StorageType.DEFAULT); assertEquals(pipelineCount, 0); // Check datanode with one RATIS/ONE and one STANDALONE/ONE pipeline and @@ -720,7 +725,7 @@ public void testCurrentRatisThreePipelineCount() pipelineCount = placementPolicy.currentRatisThreePipelineCount(nodeManager, - stateManager, healthyNodes.get(1)); + stateManager, healthyNodes.get(1), StorageType.DEFAULT); assertEquals(pipelineCount, 2); } @@ -767,7 +772,7 @@ public void testPipelinePlacementPolicyDefaultLimitFiltersNodeAtLimit() createPipelineWithReplicationConfig(p2Dns, RATIS, THREE); assertEquals(2, PipelinePlacementPolicy.currentRatisThreePipelineCount( - localNodeManager, localStateManager, target)); + localNodeManager, localStateManager, target, storageType)); // 3) Verifies node is filtered out when choosing nodes for new pipeline int nodesRequired = HddsProtos.ReplicationFactor.THREE.getNumber(); @@ -790,6 +795,7 @@ private void createPipelineWithReplicationConfig(List<DatanodeDetails> dnList, .setReplicationConfig(ReplicationConfig .fromProtoTypeAndFactor(replicationType, replicationFactor)) .setNodes(dnList) + .setSupportedStorageTier(storageTier) .build(); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java index 1ec92b31ea5..95b8b713af3 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java @@ -18,10 +18,12 @@ package org.apache.hadoop.hdds.scm.pipeline; import static org.apache.commons.collections4.CollectionUtils.intersection; +import static org.apache.hadoop.hdds.client.StorageTypeUtils.getStorageTypeProto; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_DATANODE_RATIS_VOLUME_FREE_SPACE_MIN; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_PIPELINE_PLACEMENT_IMPL_KEY; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_RATIS_PIPELINE_LIMIT; import static org.apache.hadoop.hdds.scm.exceptions.SCMException.ResultCodes.CANNOT_CREATE_PIPELINE_FOR_EMPTY_TIER; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -29,6 +31,9 @@ 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 static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.spy; import java.io.File; import java.io.IOException; @@ -52,25 +57,31 @@ import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; -import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.MetadataStorageReportProto; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.StorageReportProto; +import org.apache.hadoop.hdds.scm.HddsTestUtils; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.ContainerReplica; import org.apache.hadoop.hdds.scm.container.MockNodeManager; import org.apache.hadoop.hdds.scm.container.placement.algorithms.SCMContainerPlacementRackScatter; import org.apache.hadoop.hdds.scm.exceptions.SCMException; +import org.apache.hadoop.hdds.scm.ha.SCMContext; import org.apache.hadoop.hdds.scm.ha.SCMHAManager; import org.apache.hadoop.hdds.scm.ha.SCMHAManagerStub; import org.apache.hadoop.hdds.scm.metadata.SCMDBDefinition; import org.apache.hadoop.hdds.scm.net.NetworkTopologyImpl; +import org.apache.hadoop.hdds.scm.node.DatanodeInfo; import org.apache.hadoop.hdds.scm.node.NodeStatus; +import org.apache.hadoop.hdds.server.events.EventQueue; import org.apache.hadoop.hdds.utils.db.DBStore; import org.apache.hadoop.hdds.utils.db.DBStoreBuilder; import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.container.upgrade.UpgradeUtils; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assumptions; import org.junit.jupiter.api.Test; @@ -94,6 +105,7 @@ public class TestRatisPipelineProvider { private File testDir; private DBStore dbStore; private int nodeCount = 10; + private List<DatanodeDetails> datanodeList; public void init(int maxPipelinePerNode, StorageTier storageTier) throws Exception { init(maxPipelinePerNode, new OzoneConfiguration(), storageTier); @@ -119,8 +131,8 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf, File dir, StorageTier storageTier) throws Exception { assertTrue(storageTier.isUniform(), "Only support uniform StorageTier"); assertFalse(storageTier.equals(StorageTier.EMPTY), "not support the EMPTY StorageTier"); - StorageType storageType = storageTier.getStorageTypes( - RatisReplicationConfig.getInstance(ReplicationFactor.ONE).getRequiredNodes()).get(0); + StorageType storageType = storageTier. + getStorageTypes(ReplicationFactor.ONE.getNumber()).get(0); conf.set(HddsConfigKeys.OZONE_METADATA_DIRS, dir.getAbsolutePath()); nodeManager = new MockNodeManager(true, nodeCount, storageType); initializeCommonState(maxPipelinePerNode, conf); @@ -128,6 +140,7 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf, File dir, private void initializeCommonState(int maxPipelinePerNode, OzoneConfiguration conf) throws Exception { dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get()); + datanodeList = nodeManager.getNodes(NodeStatus.inServiceHealthy()); nodeManager.setNumPipelinePerDatanode(maxPipelinePerNode); long containerSize = (long) conf.getStorageSize( ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE, @@ -157,11 +170,12 @@ void cleanup() throws Exception { private static void assertPipelineProperties( Pipeline pipeline, HddsProtos.ReplicationFactor expectedFactor, HddsProtos.ReplicationType expectedReplicationType, - Pipeline.PipelineState expectedState) { + Pipeline.PipelineState expectedState, StorageTier expectedStorageTier) { assertEquals(expectedState, pipeline.getPipelineState()); assertEquals(expectedReplicationType, pipeline.getType()); assertEquals(expectedFactor.getNumber(), pipeline.getReplicationConfig().getRequiredNodes()); assertEquals(expectedFactor.getNumber(), pipeline.getNodes().size()); + assertEquals(expectedStorageTier, pipeline.getSupportedStorageTier()); } private void createPipelineAndAssertions( @@ -170,7 +184,7 @@ private void createPipelineAndAssertions( Pipeline pipeline = provider.create(RatisReplicationConfig .getInstance(factor), storageTier); assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, - Pipeline.PipelineState.ALLOCATED); + Pipeline.PipelineState.ALLOCATED, storageTier); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( ClientVersion.CURRENT_VERSION); stateManager.addPipeline(pipelineProto); @@ -181,7 +195,7 @@ private void createPipelineAndAssertions( HddsProtos.Pipeline pipelineProto1 = pipeline1.getProtobufMessage( ClientVersion.CURRENT_VERSION); assertPipelineProperties(pipeline1, factor, REPLICATION_TYPE, - Pipeline.PipelineState.ALLOCATED); + Pipeline.PipelineState.ALLOCATED, storageTier); // New pipeline should not overlap with the previous created pipeline assertThat(intersection(pipeline.getNodes(), pipeline1.getNodes()).size()) .isLessThan(factor.getNumber()); @@ -206,10 +220,67 @@ public void testCreatePipelineWithFactorOne(StorageTier storageTier) throws Exce createPipelineAndAssertions(HddsProtos.ReplicationFactor.ONE, storageTier); } + @Test + public void testPipelineEngagementLimitIsStorageTierAware() throws Exception { + init(1, StorageTier.DISK); + List<DatanodeDetails> diskAndSsdNodes = datanodeList.stream() + .map(TestRatisPipelineProvider::datanodeInfoWithDiskAndSsd) + .collect(Collectors.toList()); + MockNodeManager nodeManagerSpy = spy(nodeManager); + doAnswer(invocation -> diskAndSsdNodes) + .when(nodeManagerSpy).getNodes(any(NodeStatus.class)); + doAnswer(invocation -> datanodeInfoWithDiskAndSsd(invocation.getArgument(0))) + .when(nodeManagerSpy).getDatanodeInfo(any(DatanodeDetails.class)); + RatisPipelineProvider localProvider = new RatisPipelineProvider( + nodeManagerSpy, stateManager, new OzoneConfiguration(), new EventQueue(), SCMContext.emptyContext()); + for (int i = 0; i < 2; i++) { + addPipeline(diskAndSsdNodes.subList(i * 3, i * 3 + 3), + Pipeline.PipelineState.OPEN, + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), + StorageTier.DISK); + } + + Pipeline pipeline = localProvider.create( + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.SSD); + + assertPipelineProperties(pipeline, ReplicationFactor.THREE, + REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED, StorageTier.SSD); + } + + @Test + public void testGlobalPipelineNumberLimitIsNotStorageTierAware() throws Exception { + OzoneConfiguration pipelineLimitConf = new OzoneConfiguration(); + pipelineLimitConf.setInt(OZONE_SCM_RATIS_PIPELINE_LIMIT, 1); + init(0, pipelineLimitConf, StorageTier.DISK); + List<DatanodeDetails> diskAndSsdNodes = datanodeList.stream() + .map(TestRatisPipelineProvider::datanodeInfoWithDiskAndSsd) + .collect(Collectors.toList()); + MockNodeManager nodeManagerSpy = spy(nodeManager); + doAnswer(invocation -> diskAndSsdNodes) + .when(nodeManagerSpy).getNodes(any(NodeStatus.class)); + doAnswer(invocation -> datanodeInfoWithDiskAndSsd(invocation.getArgument(0))) + .when(nodeManagerSpy).getDatanodeInfo(any(DatanodeDetails.class)); + RatisPipelineProvider localProvider = new RatisPipelineProvider( + nodeManagerSpy, stateManager, pipelineLimitConf, new EventQueue(), + SCMContext.emptyContext()); + for (int i = 0; i < 2; i++) { + addPipeline(diskAndSsdNodes.subList(i * 3, i * 3 + 3), + Pipeline.PipelineState.OPEN, + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), + StorageTier.DISK); + } + + SCMException exception = assertThrows(SCMException.class, () -> + localProvider.create( + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.SSD)); + + assertThat(exception.getMessage()).contains("would exceed the limit"); + } + private List<DatanodeDetails> createListOfNodes(int count) { List<DatanodeDetails> nodes = new ArrayList<>(); for (int i = 0; i < count; i++) { - nodes.add(MockDatanodeDetails.randomDatanodeDetails()); + nodes.add(datanodeList.get(i)); } return nodes; } @@ -222,7 +293,7 @@ public void testCreatePipelineWithFactor(StorageTier storageTier) throws Excepti Pipeline pipeline = provider.create(RatisReplicationConfig .getInstance(factor), storageTier); assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, - Pipeline.PipelineState.ALLOCATED); + Pipeline.PipelineState.ALLOCATED, storageTier); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( ClientVersion.CURRENT_VERSION); stateManager.addPipeline(pipelineProto); @@ -231,7 +302,7 @@ public void testCreatePipelineWithFactor(StorageTier storageTier) throws Excepti Pipeline pipeline1 = provider.create(RatisReplicationConfig .getInstance(factor), storageTier); assertPipelineProperties(pipeline1, factor, REPLICATION_TYPE, - Pipeline.PipelineState.ALLOCATED); + Pipeline.PipelineState.ALLOCATED, storageTier); HddsProtos.Pipeline pipelineProto1 = pipeline1.getProtobufMessage( ClientVersion.CURRENT_VERSION); stateManager.addPipeline(pipelineProto1); @@ -246,15 +317,17 @@ public void testCreatePipelineWithNodes() throws Exception { HddsProtos.ReplicationFactor factor = HddsProtos.ReplicationFactor.THREE; Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(factor), - createListOfNodes(factor.getNumber())); + createListOfNodes(factor.getNumber()), + StorageTier.getDefaultTier()); assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, - Pipeline.PipelineState.OPEN); + Pipeline.PipelineState.OPEN, StorageTier.getDefaultTier()); factor = HddsProtos.ReplicationFactor.ONE; pipeline = provider.create(RatisReplicationConfig.getInstance(factor), - createListOfNodes(factor.getNumber())); + createListOfNodes(factor.getNumber()), + StorageTier.getDefaultTier()); assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, - Pipeline.PipelineState.OPEN); + Pipeline.PipelineState.OPEN, StorageTier.getDefaultTier()); } @Test @@ -268,13 +341,12 @@ public void testCreateFactorTHREEPipelineWithSameDatanodes() Pipeline pipeline1 = provider.create( RatisReplicationConfig.getInstance(ReplicationFactor.THREE), - healthyNodes); + healthyNodes, StorageTier.getDefaultTier()); Pipeline pipeline2 = provider.create( RatisReplicationConfig.getInstance(ReplicationFactor.THREE), - healthyNodes); + healthyNodes, StorageTier.getDefaultTier()); Pipeline pipeline3 = provider.createForRead( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE), - replicas); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), replicas); assertEquals(pipeline1.getNodeSet(), pipeline2.getNodeSet()); assertEquals(pipeline2.getNodeSet(), pipeline3.getNodeSet()); @@ -374,21 +446,21 @@ public void testCreatePipelinesDnExclude(StorageTier storageTier) throws Excepti for (int i = 0; i < maxPipelinePerNode; i++) { // Saturate pipeline counts on all the 1st 3 DNs. addPipeline(dns, Pipeline.PipelineState.OPEN, - RatisReplicationConfig.getInstance(factor)); + RatisReplicationConfig.getInstance(factor), storageTier); } Set<DatanodeDetails> membersOfOpenPipelines = new HashSet<>(dns); // Use up next 3 DNs for a closed pipeline. dns = healthyNodes.subList(3, 6); addPipeline(dns, Pipeline.PipelineState.CLOSED, - RatisReplicationConfig.getInstance(factor)); + RatisReplicationConfig.getInstance(factor), storageTier); Set<DatanodeDetails> membersOfClosedPipelines = new HashSet<>(dns); // only 2 healthy DNs left that are not part of any pipeline Pipeline pipeline = provider.create( RatisReplicationConfig.getInstance(factor), storageTier); assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, - Pipeline.PipelineState.ALLOCATED); + Pipeline.PipelineState.ALLOCATED, storageTier); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( ClientVersion.CURRENT_VERSION); nodeManager.addPipeline(pipeline); @@ -559,7 +631,8 @@ public void testCreatePipelineThrowErrorWithDataNodeLimit(int limit, int pipelin // Validate exception message. String expectedError = String.format( - "Cannot create pipeline as it would exceed the limit per datanode: %d replicationConfig: RATIS/THREE", limit); + "Cannot create pipeline for StorageTier DISK as it would exceed the limit per" + + " datanode: %d replicationConfig: RATIS/THREE", limit); assertEquals(expectedError, exception.getMessage()); } @@ -580,13 +653,14 @@ public void testCreatePipelinesInEmptyTier() throws Exception { private void addPipeline( List<DatanodeDetails> dns, - Pipeline.PipelineState open, ReplicationConfig replicationConfig) + Pipeline.PipelineState open, ReplicationConfig replicationConfig, StorageTier storageTier) throws IOException, TimeoutException { Pipeline openPipeline = Pipeline.newBuilder() .setReplicationConfig(replicationConfig) .setNodes(dns) .setState(open) .setId(PipelineID.randomId()) + .setSupportedStorageTier(storageTier) .build(); HddsProtos.Pipeline pipelineProto = openPipeline.getProtobufMessage( ClientVersion.CURRENT_VERSION); @@ -622,4 +696,24 @@ static Stream<StorageTier> storageTiers() { StorageTier.ARCHIVE ); } + + private static DatanodeInfo datanodeInfoWithDiskAndSsd(DatanodeDetails dn) { + DatanodeInfo datanodeInfo = new DatanodeInfo( + dn, NodeStatus.inServiceHealthy(), UpgradeUtils.defaultLayoutVersionProto(), + HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT); + List<StorageReportProto> storageReports = new ArrayList<>(); + long capacity = 10L * 1024 * 1024 * 1024; + storageReports.add(HddsTestUtils.createStorageReport( + dn.getID(), "/disk-" + dn.getUuidString(), + capacity, 0L, capacity, getStorageTypeProto(StorageType.DISK))); + storageReports.add(HddsTestUtils.createStorageReport( + dn.getID(), "/ssd-" + dn.getUuidString(), + capacity, 0L, capacity, getStorageTypeProto(StorageType.SSD))); + MetadataStorageReportProto metadataReport = HddsTestUtils.createMetadataStorageReport( + "/metadata-" + dn.getUuidString(), capacity, 0L, capacity, null); + datanodeInfo.updateStorageReports(storageReports); + datanodeInfo.updateMetaDataStorageReports(Collections.singletonList(metadataReport)); + return datanodeInfo; + } + } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java index 42aa1943aa7..3326c09d5cd 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java @@ -22,30 +22,41 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import java.io.File; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.stream.Stream; import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.HddsConfigKeys; 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.DatanodeDetails; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.StorageReportProto; +import org.apache.hadoop.hdds.scm.HddsTestUtils; import org.apache.hadoop.hdds.scm.container.MockNodeManager; import org.apache.hadoop.hdds.scm.exceptions.SCMException; import org.apache.hadoop.hdds.scm.ha.SCMHAManager; import org.apache.hadoop.hdds.scm.ha.SCMHAManagerStub; import org.apache.hadoop.hdds.scm.metadata.SCMDBDefinition; +import org.apache.hadoop.hdds.scm.node.DatanodeInfo; import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.node.NodeStatus; import org.apache.hadoop.hdds.utils.db.DBStore; import org.apache.hadoop.hdds.utils.db.DBStoreBuilder; import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.container.common.SCMTestUtils; +import org.apache.hadoop.ozone.container.upgrade.UpgradeUtils; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -63,6 +74,7 @@ public class TestSimplePipelineProvider { @TempDir private File testDir; private DBStore dbStore; + private List<DatanodeDetails> datanodeList; @BeforeEach public void startup() throws Exception { @@ -78,6 +90,7 @@ public void init(StorageTier storageTier) throws Exception { StorageType storageType = storageTier.getStorageTypes( RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE).getRequiredNodes()).get(0); NodeManager nodeManager = new MockNodeManager(true, 10, storageType); + datanodeList = nodeManager.getNodes(NodeStatus.inServiceHealthy()); SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); stateManager = PipelineStateManagerImpl.newBuilder() .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) @@ -109,6 +122,7 @@ public void testCreatePipelineWithFactor(StorageTier storageTier) throws Excepti assertEquals(pipeline.getReplicationConfig().getRequiredNodes(), factor.getNumber()); assertEquals(pipeline.getPipelineState(), Pipeline.PipelineState.OPEN); assertEquals(pipeline.getNodes().size(), factor.getNumber()); + assertEquals(storageTier, pipeline.getSupportedStorageTier()); factor = HddsProtos.ReplicationFactor.ONE; Pipeline pipeline1 = @@ -122,12 +136,13 @@ public void testCreatePipelineWithFactor(StorageTier storageTier) throws Excepti .getReplicationFactor(), factor); assertEquals(pipeline1.getPipelineState(), Pipeline.PipelineState.OPEN); assertEquals(pipeline1.getNodes().size(), factor.getNumber()); + assertEquals(storageTier, pipeline1.getSupportedStorageTier()); } private List<DatanodeDetails> createListOfNodes(int nodeCount) { List<DatanodeDetails> nodes = new ArrayList<>(); for (int i = 0; i < nodeCount; i++) { - nodes.add(MockDatanodeDetails.randomDatanodeDetails()); + nodes.add(datanodeList.get(i)); } return nodes; } @@ -139,7 +154,8 @@ public void testCreatePipelineWithNodes() HddsProtos.ReplicationFactor factor = HddsProtos.ReplicationFactor.THREE; Pipeline pipeline = provider.create(StandaloneReplicationConfig.getInstance(factor), - createListOfNodes(factor.getNumber())); + createListOfNodes(factor.getNumber()), + StorageTier.getDefaultTier()); assertEquals(pipeline.getType(), HddsProtos.ReplicationType.STAND_ALONE); assertEquals( @@ -147,16 +163,63 @@ public void testCreatePipelineWithNodes() .getReplicationFactor(), factor); assertEquals(pipeline.getPipelineState(), Pipeline.PipelineState.OPEN); assertEquals(pipeline.getNodes().size(), factor.getNumber()); + assertEquals(StorageTier.getDefaultTier(), pipeline.getSupportedStorageTier()); factor = HddsProtos.ReplicationFactor.ONE; pipeline = provider.create(StandaloneReplicationConfig.getInstance(factor), - createListOfNodes(factor.getNumber())); + createListOfNodes(factor.getNumber()), + StorageTier.getDefaultTier()); assertEquals(pipeline.getType(), HddsProtos.ReplicationType.STAND_ALONE); assertEquals( ((StandaloneReplicationConfig) pipeline.getReplicationConfig()) .getReplicationFactor(), factor); assertEquals(pipeline.getPipelineState(), Pipeline.PipelineState.OPEN); assertEquals(pipeline.getNodes().size(), factor.getNumber()); + assertEquals(StorageTier.getDefaultTier(), pipeline.getSupportedStorageTier()); + } + + @Test + public void testCreatedPipelineOnlySupportsRequestedStorageTier() + throws Exception { + NodeManager nodeManager = mock(NodeManager.class); + PipelineStateManager pipelineStateManager = mock(PipelineStateManager.class); + List<DatanodeDetails> nodes = new ArrayList<>(); + for (int i = 0; i < HddsProtos.ReplicationFactor.THREE.getNumber(); i++) { + nodes.add(createDatanodeInfoWithStorageReports( + StorageTypeProto.DISK, StorageTypeProto.SSD)); + } + when(nodeManager.getNodes(NodeStatus.inServiceHealthy())) + .thenReturn(nodes); + when(nodeManager.getDatanodeInfo(any())) + .thenAnswer(invocation -> invocation.getArgument(0)); + when(pipelineStateManager.getPipelines( + any(ReplicationConfig.class))) + .thenReturn(Collections.emptyList()); + + PipelineProvider pipelineProvider = + new SimplePipelineProvider(nodeManager, pipelineStateManager); + Pipeline pipeline = pipelineProvider.create(StandaloneReplicationConfig.getInstance( + HddsProtos.ReplicationFactor.THREE), StorageTier.DISK); + + assertEquals(StorageTier.DISK, pipeline.getSupportedStorageTier()); + } + + private DatanodeInfo createDatanodeInfoWithStorageReports( + HddsProtos.StorageTypeProto... storageTypes) { + DatanodeInfo datanodeInfo = new DatanodeInfo( + MockDatanodeDetails.randomDatanodeDetails(), + NodeStatus.inServiceHealthy(), UpgradeUtils.defaultLayoutVersionProto(), + HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT); + List<StorageReportProto> storageReports = new ArrayList<>(); + for (HddsProtos.StorageTypeProto storageType : storageTypes) { + storageReports.add(StorageReportProto.newBuilder() + .setStorageUuid(datanodeInfo.getUuidString() + "-" + storageType) + .setStorageLocation("/data-" + storageType) + .setStorageType(storageType) + .build()); + } + datanodeInfo.updateStorageReports(storageReports); + return datanodeInfo; } @Test diff --git a/hadoop-ozone/common/src/test/java/org/apache/hadoop/hdds/protocol/StorageTierTest.java b/hadoop-ozone/common/src/test/java/org/apache/hadoop/hdds/protocol/StorageTierTest.java index 0630dd45cd3..af978dfef36 100644 --- a/hadoop-ozone/common/src/test/java/org/apache/hadoop/hdds/protocol/StorageTierTest.java +++ b/hadoop-ozone/common/src/test/java/org/apache/hadoop/hdds/protocol/StorageTierTest.java @@ -59,7 +59,8 @@ void testGetStorageTypesWithReplicationConfig() { } for (ReplicationConfig replicationConfig : Arrays.asList(ratisOne, ratisThree, standaloneOne, standaloneThree)) { - List<StorageType> storageTypes = tier.getStorageTypes(replicationConfig.getRequiredNodes()); + List<StorageType> storageTypes = tier.getStorageTypes( + replicationConfig.getRequiredNodes()); if (tier.equals(StorageTier.EMPTY)) { assertEquals(0, storageTypes.size()); } else { 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 431e8857490..5375206bab7 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 @@ -104,6 +104,10 @@ public void testPipelineWithScmRestart() assertNotSame(ratisPipeline2AfterRestart, ratisPipeline2); assertEquals(ratisPipeline1AfterRestart, ratisPipeline1); assertEquals(ratisPipeline2AfterRestart, ratisPipeline2); + assertEquals(ratisPipeline1AfterRestart.getSupportedStorageTier(), + ratisPipeline1.getSupportedStorageTier()); + assertEquals(ratisPipeline2AfterRestart.getSupportedStorageTier(), + ratisPipeline2.getSupportedStorageTier()); // Try creating a new container, it should be from the same pipeline // as was before restart diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerImpl.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerImpl.java index 5a526cb3c5a..0aa821d1535 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerImpl.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerImpl.java @@ -82,6 +82,7 @@ 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.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; @@ -821,7 +822,7 @@ private static void createKeyWithPipeline(OmKeyArgs keyArgs) throws IOException assumeFalse(nodeList.get(0).equals(nodeList.get(2))); // create a pipeline using 3 datanodes Pipeline pipeline = scm.getPipelineManager().createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE), nodeList); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), nodeList, StorageTier.getDefaultTier()); List<OmKeyLocationInfo> locationInfoList = new ArrayList<>(); List<OmKeyLocationInfo> locationList = keySession.getKeyInfo().getLatestVersionLocations().getLocationList(); diff --git a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineFactory.java b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineFactory.java index 55fee8f8653..f8063a3a3a2 100644 --- a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineFactory.java +++ b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineFactory.java @@ -64,7 +64,7 @@ public Pipeline create(ReplicationConfig config, @Override public Pipeline create(ReplicationConfig config, - List<DatanodeDetails> nodes) { + List<DatanodeDetails> nodes, StorageTier storageTier) { throw new UnsupportedOperationException( "Trying to create pipeline in Recon, which is prohibited!"); } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
