Skip to content

Commit 4004e4a

Browse files
committed
Storage tiering. Add creation of the typed pipelines
1 parent 94a2b64 commit 4004e4a

11 files changed

Lines changed: 256 additions & 18 deletions

File tree

hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/ScmConfigKeys.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -445,6 +445,11 @@ public final class ScmConfigKeys {
445445
public static final String OZONE_SCM_PIPELINE_SCRUB_INTERVAL_DEFAULT =
446446
"150s";
447447

448+
public static final String OZONE_SCM_PIPELINE_CREATION_STORAGE_TYPE_AWARE =
449+
"ozone.scm.pipeline.creation.storage-type-aware.enabled";
450+
public static final boolean
451+
OZONE_SCM_PIPELINE_CREATION_STORAGE_TYPE_AWARE_DEFAULT = false;
452+
448453
// Allow SCM to auto create factor ONE ratis pipeline.
449454
public static final String OZONE_SCM_PIPELINE_AUTO_CREATE_FACTOR_ONE =
450455
"ozone.scm.pipeline.creation.auto.factor.one";

hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/BackgroundPipelineCreator.java

Lines changed: 63 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,14 @@
2525
import static org.apache.hadoop.hdds.scm.ha.SCMService.Event.PRE_CHECK_COMPLETED;
2626
import static org.apache.hadoop.hdds.scm.ha.SCMService.Event.UNHEALTHY_TO_HEALTHY_NODE_HANDLER_TRIGGERED;
2727

28+
import com.google.common.annotations.VisibleForTesting;
2829
import com.google.common.util.concurrent.ThreadFactoryBuilder;
2930
import java.io.IOException;
3031
import java.time.Clock;
32+
import java.util.AbstractMap;
3133
import java.util.ArrayList;
3234
import java.util.List;
35+
import java.util.Map;
3336
import java.util.concurrent.TimeUnit;
3437
import java.util.concurrent.atomic.AtomicBoolean;
3538
import java.util.concurrent.locks.Lock;
@@ -40,6 +43,7 @@
4043
import org.apache.hadoop.hdds.client.ReplicationConfig;
4144
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
4245
import org.apache.hadoop.hdds.conf.ConfigurationSource;
46+
import org.apache.hadoop.hdds.protocol.StorageType;
4347
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
4448
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
4549
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
@@ -88,6 +92,7 @@ public class BackgroundPipelineCreator implements SCMService {
8892
private final AtomicBoolean running = new AtomicBoolean(false);
8993
private final long intervalInMillis;
9094
private final Clock clock;
95+
private final boolean storageTypeAwareCreation;
9196

9297
BackgroundPipelineCreator(PipelineManager pipelineManager,
9398
ConfigurationSource conf, SCMContext scmContext, Clock clock) {
@@ -110,6 +115,10 @@ public class BackgroundPipelineCreator implements SCMService {
110115
ScmConfigKeys.OZONE_SCM_PIPELINE_CREATION_INTERVAL_DEFAULT,
111116
TimeUnit.MILLISECONDS);
112117

118+
this.storageTypeAwareCreation = conf.getBoolean(
119+
ScmConfigKeys.OZONE_SCM_PIPELINE_CREATION_STORAGE_TYPE_AWARE,
120+
ScmConfigKeys.OZONE_SCM_PIPELINE_CREATION_STORAGE_TYPE_AWARE_DEFAULT);
121+
113122
threadName = scmContext.threadNamePrefix() + THREAD_NAME;
114123
}
115124

@@ -203,7 +212,8 @@ private boolean skipCreation(ReplicationConfig replicationConfig,
203212
return true;
204213
}
205214

206-
private void createPipelines() throws RuntimeException {
215+
@VisibleForTesting
216+
void createPipelines() throws RuntimeException {
207217
// TODO: #CLUTIL Different replication factor may need to be supported
208218
HddsProtos.ReplicationType type = HddsProtos.ReplicationType.valueOf(
209219
conf.get(OzoneConfigKeys.OZONE_REPLICATION_TYPE,
@@ -212,8 +222,7 @@ private void createPipelines() throws RuntimeException {
212222
ScmConfigKeys.OZONE_SCM_PIPELINE_AUTO_CREATE_FACTOR_ONE,
213223
ScmConfigKeys.OZONE_SCM_PIPELINE_AUTO_CREATE_FACTOR_ONE_DEFAULT);
214224

215-
List<ReplicationConfig> list =
216-
new ArrayList<>();
225+
List<ReplicationConfig> replicationConfigs = new ArrayList<>();
217226
for (HddsProtos.ReplicationFactor factor : HddsProtos.ReplicationFactor
218227
.values()) {
219228
if (factor == ReplicationFactor.ZERO) {
@@ -233,10 +242,20 @@ private void createPipelines() throws RuntimeException {
233242
// Skip this iteration for creating pipeline
234243
continue;
235244
}
236-
list.add(replicationConfig);
245+
replicationConfigs.add(replicationConfig);
246+
}
247+
248+
if (storageTypeAwareCreation) {
249+
createTypedPipelines(replicationConfigs);
250+
} else {
251+
createUntypedPipelines(replicationConfigs);
237252
}
238253

239-
LoopingIterator it = new LoopingIterator(list);
254+
LOG.debug("BackgroundPipelineCreator createPipelines finished.");
255+
}
256+
257+
private void createUntypedPipelines(List<ReplicationConfig> configs) {
258+
LoopingIterator it = new LoopingIterator(configs);
240259
while (it.hasNext()) {
241260
ReplicationConfig replicationConfig =
242261
(ReplicationConfig) it.next();
@@ -251,8 +270,46 @@ private void createPipelines() throws RuntimeException {
251270
it.remove();
252271
}
253272
}
273+
}
254274

255-
LOG.debug("BackgroundPipelineCreator createPipelines finished.");
275+
private void createTypedPipelines(List<ReplicationConfig> configs) {
276+
// Build (ReplicationConfig, StorageType) pairs: for each config,
277+
// one null entry (untyped) plus one per concrete StorageType.
278+
StorageType[] storageTypes = {
279+
StorageType.SSD, StorageType.DISK, StorageType.ARCHIVE
280+
};
281+
List<Map.Entry<ReplicationConfig, StorageType>> pairs = new ArrayList<>();
282+
for (ReplicationConfig config : configs) {
283+
pairs.add(new AbstractMap.SimpleEntry<>(config, null));
284+
for (StorageType st : storageTypes) {
285+
pairs.add(new AbstractMap.SimpleEntry<>(config, st));
286+
}
287+
}
288+
289+
LoopingIterator it = new LoopingIterator(pairs);
290+
while (it.hasNext()) {
291+
@SuppressWarnings("unchecked")
292+
Map.Entry<ReplicationConfig, StorageType> entry =
293+
(Map.Entry<ReplicationConfig, StorageType>) it.next();
294+
295+
try {
296+
Pipeline pipeline;
297+
if (entry.getValue() == null) {
298+
pipeline = pipelineManager.createPipeline(entry.getKey());
299+
} else {
300+
pipeline = pipelineManager.createPipeline(
301+
entry.getKey(), entry.getValue());
302+
}
303+
LOG.info("Created new pipeline {} with StorageType {}",
304+
pipeline, entry.getValue());
305+
} catch (IOException ioe) {
306+
it.remove();
307+
} catch (Throwable t) {
308+
LOG.error("Error while creating pipelines for StorageType "
309+
+ entry.getValue(), t);
310+
it.remove();
311+
}
312+
}
256313
}
257314

258315
@Override

hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import java.util.Set;
2626
import org.apache.hadoop.hdds.client.ReplicationConfig;
2727
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
28+
import org.apache.hadoop.hdds.protocol.StorageType;
2829
import org.apache.hadoop.hdds.scm.container.ContainerID;
2930
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
3031
import org.apache.hadoop.hdds.utils.db.CodecException;
@@ -39,6 +40,11 @@ public interface PipelineManager extends Closeable, PipelineManagerMXBean {
3940
Pipeline createPipeline(ReplicationConfig replicationConfig)
4041
throws IOException;
4142

43+
default Pipeline createPipeline(ReplicationConfig replicationConfig,
44+
StorageType storageType) throws IOException {
45+
return createPipeline(replicationConfig);
46+
}
47+
4248
Pipeline createPipeline(ReplicationConfig replicationConfig,
4349
List<DatanodeDetails> excludedNodes,
4450
List<DatanodeDetails> favoredNodes)

hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
import java.util.Map;
3131
import java.util.NavigableSet;
3232
import java.util.Set;
33+
import java.util.UUID;
3334
import java.util.concurrent.TimeUnit;
3435
import java.util.concurrent.atomic.AtomicBoolean;
3536
import java.util.concurrent.locks.ReentrantReadWriteLock;
@@ -41,6 +42,7 @@
4142
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
4243
import org.apache.hadoop.hdds.conf.ConfigurationSource;
4344
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
45+
import org.apache.hadoop.hdds.protocol.StorageType;
4446
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
4547
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
4648
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType;
@@ -56,6 +58,7 @@
5658
import org.apache.hadoop.hdds.scm.ha.SCMServiceManager;
5759
import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
5860
import org.apache.hadoop.hdds.scm.node.NodeManager;
61+
import org.apache.hadoop.hdds.scm.node.NodeStatus;
5962
import org.apache.hadoop.hdds.scm.server.upgrade.FinalizationManager;
6063
import org.apache.hadoop.hdds.server.events.EventPublisher;
6164
import org.apache.hadoop.hdds.utils.db.CodecException;
@@ -268,6 +271,33 @@ public Pipeline createPipeline(ReplicationConfig replicationConfig,
268271
}
269272
}
270273

274+
@Override
275+
public Pipeline createPipeline(ReplicationConfig replicationConfig,
276+
StorageType storageType) throws IOException {
277+
if (storageType == null) {
278+
return createPipeline(replicationConfig);
279+
}
280+
// Compute excluded nodes: all healthy nodes that do NOT have the
281+
// requested StorageType.
282+
List<DatanodeDetails> allHealthy =
283+
nodeManager.getNodes(NodeStatus.inServiceHealthy());
284+
Set<UUID> qualifiedNodeIds =
285+
PipelineStorageTypeFilter.getNodesWithStorageType(
286+
nodeManager, storageType);
287+
288+
if (qualifiedNodeIds.isEmpty()) {
289+
throw new IOException("No healthy nodes with StorageType "
290+
+ storageType + " available for pipeline creation");
291+
}
292+
293+
List<DatanodeDetails> excludedNodes = allHealthy.stream()
294+
.filter(dn -> !qualifiedNodeIds.contains(dn.getUuid()))
295+
.collect(Collectors.toList());
296+
297+
return createPipeline(replicationConfig, excludedNodes,
298+
Collections.emptyList());
299+
}
300+
271301
private void checkIfPipelineCreationIsAllowed(
272302
ReplicationConfig replicationConfig) throws IOException {
273303
if (!isPipelineCreationAllowed() && !factorOne(replicationConfig)) {
Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,103 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.hadoop.hdds.scm.pipeline;
19+
20+
import static org.mockito.Mockito.any;
21+
import static org.mockito.Mockito.atLeastOnce;
22+
import static org.mockito.Mockito.mock;
23+
import static org.mockito.Mockito.never;
24+
import static org.mockito.Mockito.verify;
25+
import static org.mockito.Mockito.when;
26+
27+
import java.io.IOException;
28+
import java.time.Instant;
29+
import java.time.ZoneOffset;
30+
import org.apache.hadoop.hdds.client.ReplicationConfig;
31+
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
32+
import org.apache.hadoop.hdds.protocol.StorageType;
33+
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
34+
import org.apache.hadoop.hdds.scm.ha.SCMContext;
35+
import org.apache.ozone.test.TestClock;
36+
import org.junit.jupiter.api.Test;
37+
38+
/**
39+
* Tests for storage-type-aware pipeline creation in
40+
* BackgroundPipelineCreator.
41+
*/
42+
public class TestBackgroundPipelineCreatorStorageType {
43+
44+
@Test
45+
public void testStorageTypeAwareDisabled() throws IOException {
46+
OzoneConfiguration conf = new OzoneConfiguration();
47+
conf.setBoolean(
48+
ScmConfigKeys.OZONE_SCM_PIPELINE_CREATION_STORAGE_TYPE_AWARE,
49+
false);
50+
51+
PipelineManager pipelineManager = mock(PipelineManager.class);
52+
when(pipelineManager.createPipeline(any(ReplicationConfig.class)))
53+
.thenThrow(new IOException("exhausted"));
54+
55+
SCMContext scmContext = SCMContext.emptyContext();
56+
57+
TestClock clock = new TestClock(Instant.now(), ZoneOffset.UTC);
58+
BackgroundPipelineCreator creator =
59+
new BackgroundPipelineCreator(pipelineManager, conf, scmContext,
60+
clock);
61+
62+
creator.createPipelines();
63+
64+
// Untyped createPipeline(ReplicationConfig) should have been called.
65+
verify(pipelineManager, atLeastOnce())
66+
.createPipeline(any(ReplicationConfig.class));
67+
// Typed createPipeline(ReplicationConfig, StorageType) should NOT
68+
// have been called.
69+
verify(pipelineManager, never())
70+
.createPipeline(any(ReplicationConfig.class),
71+
any(StorageType.class));
72+
}
73+
74+
@Test
75+
public void testStorageTypeAwareEnabled() throws IOException {
76+
OzoneConfiguration conf = new OzoneConfiguration();
77+
conf.setBoolean(
78+
ScmConfigKeys.OZONE_SCM_PIPELINE_CREATION_STORAGE_TYPE_AWARE,
79+
true);
80+
81+
PipelineManager pipelineManager = mock(PipelineManager.class);
82+
when(pipelineManager.createPipeline(any(ReplicationConfig.class)))
83+
.thenThrow(new IOException("exhausted"));
84+
when(pipelineManager.createPipeline(any(ReplicationConfig.class),
85+
any(StorageType.class)))
86+
.thenThrow(new IOException("exhausted"));
87+
88+
SCMContext scmContext = SCMContext.emptyContext();
89+
90+
TestClock clock = new TestClock(Instant.now(), ZoneOffset.UTC);
91+
BackgroundPipelineCreator creator =
92+
new BackgroundPipelineCreator(pipelineManager, conf, scmContext,
93+
clock);
94+
95+
creator.createPipelines();
96+
97+
// When storage-type-aware is enabled, the typed method should be called
98+
// for SSD, DISK, and ARCHIVE.
99+
verify(pipelineManager, atLeastOnce())
100+
.createPipeline(any(ReplicationConfig.class),
101+
any(StorageType.class));
102+
}
103+
}

hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@
7171
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
7272
import org.apache.hadoop.hdds.protocol.DatanodeID;
7373
import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
74+
import org.apache.hadoop.hdds.protocol.StorageType;
7475
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
7576
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
7677
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos;
@@ -920,6 +921,34 @@ public void testWaitForAllocatedPipeline() throws IOException {
920921
pipelineManager.close();
921922
}
922923

924+
@Test
925+
public void testCreatePipelineWithStorageType() throws Exception {
926+
PipelineManagerImpl pipelineManager = createPipelineManager(true);
927+
928+
// MockNodeManager creates storage reports with DISK type by default.
929+
// DISK-typed pipeline should succeed.
930+
Pipeline diskPipeline = pipelineManager.createPipeline(
931+
RatisReplicationConfig.getInstance(ReplicationFactor.THREE),
932+
StorageType.DISK);
933+
assertNotNull(diskPipeline);
934+
assertEquals(3, diskPipeline.getNodes().size());
935+
936+
// SSD-typed pipeline should fail since no nodes have SSD storage.
937+
assertThrows(IOException.class,
938+
() -> pipelineManager.createPipeline(
939+
RatisReplicationConfig.getInstance(ReplicationFactor.THREE),
940+
StorageType.SSD));
941+
942+
// null StorageType should fall through to untyped creation.
943+
Pipeline untypedPipeline = pipelineManager.createPipeline(
944+
RatisReplicationConfig.getInstance(ReplicationFactor.THREE),
945+
(StorageType) null);
946+
assertNotNull(untypedPipeline);
947+
assertEquals(3, untypedPipeline.getNodes().size());
948+
949+
pipelineManager.close();
950+
}
951+
923952
public void testCreatePipelineForRead() throws IOException {
924953
PipelineManager pipelineManager = createPipelineManager(true);
925954
List<DatanodeDetails> dns = nodeManager

hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestKeyManagerImpl.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,7 @@
8484
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
8585
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
8686
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
87+
import org.apache.hadoop.hdds.protocol.StorageType;
8788
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
8889
import org.apache.hadoop.hdds.scm.HddsTestUtils;
8990
import org.apache.hadoop.hdds.scm.HddsWhiteboxTestUtils;
@@ -227,7 +228,8 @@ public static void setUp() throws Exception {
227228
any(ReplicationConfig.class),
228229
anyString(),
229230
any(ExcludeList.class),
230-
anyString())).thenThrow(
231+
anyString(),
232+
any(StorageType.class))).thenThrow(
231233
new SCMException("SafeModePrecheck failed for allocateBlock",
232234
ResultCodes.SAFE_MODE_EXCEPTION));
233235
createVolume(VOLUME_NAME);

0 commit comments

Comments
 (0)