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 30eb83ab735f..583ef5cac17b 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 @@ -17,6 +17,8 @@ package org.apache.hadoop.hdds.scm.pipeline; +import static org.apache.hadoop.hdds.protocol.DatanodeDetails.Port.Name.RATIS_DATASTREAM; + import com.google.common.annotations.VisibleForTesting; import java.io.IOException; import java.util.ArrayList; @@ -155,7 +157,7 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, ); } - final List dns; + List dns; final ReplicationFactor factor = replicationConfig.getReplicationFactor(); switch (factor) { @@ -163,17 +165,22 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, dns = pickNodesNotUsed(replicationConfig, minRatisVolumeSizeBytes, containerSizeBytes); break; case THREE: - List excludeDueToEngagement = filterPipelineEngagement(); - if (!excludeDueToEngagement.isEmpty()) { - if (excludedNodes.isEmpty()) { - excludedNodes = excludeDueToEngagement; - } else { - excludedNodes.addAll(excludeDueToEngagement); - } + List excludeDueToEngagement = filterNodes(true); + List currentExcluded = new ArrayList<>(excludedNodes); + currentExcluded.addAll(excludeDueToEngagement); + try { + dns = placementPolicy.chooseDatanodes(currentExcluded, + favoredNodes, factor.getNumber(), minRatisVolumeSizeBytes, + containerSizeBytes); + } catch (SCMException scmException) { + excludeDueToEngagement = filterNodes(false); + currentExcluded = new ArrayList<>(excludedNodes); + currentExcluded.addAll(excludeDueToEngagement); + dns = placementPolicy.chooseDatanodes(currentExcluded, + favoredNodes, factor.getNumber(), minRatisVolumeSizeBytes, + containerSizeBytes); } - dns = placementPolicy.chooseDatanodes(excludedNodes, - favoredNodes, factor.getNumber(), minRatisVolumeSizeBytes, - containerSizeBytes); + break; default: throw new IllegalStateException("Unknown factor: " + factor.name()); @@ -224,14 +231,19 @@ public Pipeline createForRead( .build(); } - private List filterPipelineEngagement() { + /** + * + * @return + */ + private List filterNodes(boolean filterRatisStreaming) { final NodeManager nodeManager = getNodeManager(); final PipelineStateManager stateManager = getPipelineStateManager(); final List healthyNodes = nodeManager.getNodes(NodeStatus.inServiceHealthy()); final List excluded = new ArrayList<>(); for (DatanodeDetails d : healthyNodes) { final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d); - if (count >= nodeManager.pipelineLimit(d)) { + if (count >= nodeManager.pipelineLimit(d) || + (filterRatisStreaming && !d.hasPort(RATIS_DATASTREAM))) { excluded.add(d); } } 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 bf352f2051ef..ab2a1c3d77d4 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 @@ -34,7 +34,9 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; +import java.util.Random; import java.util.Set; +import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeoutException; import java.util.stream.Collectors; import org.apache.hadoop.hdds.HddsConfigKeys; @@ -57,6 +59,7 @@ 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.NodeStatus; import org.apache.hadoop.hdds.utils.db.DBStore; import org.apache.hadoop.hdds.utils.db.DBStoreBuilder; @@ -93,10 +96,22 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf) init(maxPipelinePerNode, conf, testDir); } + public void initWithNodes(int maxPipelinePerNode, List nodes, int count) + throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(HddsConfigKeys.OZONE_METADATA_DIRS, testDir.getAbsolutePath()); + nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, count); + initializeCommonState(maxPipelinePerNode, conf); + } + public void init(int maxPipelinePerNode, OzoneConfiguration conf, File dir) throws Exception { conf.set(HddsConfigKeys.OZONE_METADATA_DIRS, dir.getAbsolutePath()); - dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get()); nodeManager = new MockNodeManager(true, nodeCount); + initializeCommonState(maxPipelinePerNode, conf); + } + + private void initializeCommonState(int maxPipelinePerNode, OzoneConfiguration conf) throws Exception { + dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get()); nodeManager.setNumPipelinePerDatanode(maxPipelinePerNode); long containerSize = (long) conf.getStorageSize( ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE, @@ -246,6 +261,73 @@ public void testCreateFactorTHREEPipelineWithSameDatanodes() assertEquals(pipeline2.getNodeSet(), pipeline3.getNodeSet()); } + private DatanodeDetails createDatanodeDetails(boolean supportRatisStreaming) { + Random random = ThreadLocalRandom.current(); + String ipAddress = random.nextInt(256) + + "." + random.nextInt(256) + + "." + random.nextInt(256) + + "." + random.nextInt(256); + + DatanodeDetails.Builder dn = DatanodeDetails.newBuilder() + .setID(DatanodeID.randomID()) + .setHostName("localhost" + "-" + ipAddress) + .setIpAddress(ipAddress) + .setNetworkLocation(null) + .setPersistedOpState(HddsProtos.NodeOperationalState.IN_SERVICE) + .setPersistedOpStateExpiry(0); + + for (DatanodeDetails.Port.Name name : DatanodeDetails.Port.Name.values()) { + if (!supportRatisStreaming && name == DatanodeDetails.Port.Name.RATIS_DATASTREAM) { + continue; + } + dn.addPort(DatanodeDetails.newPort(name, 0)); + } + return dn.build(); + } + + @Test + public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception { + List nodes = new ArrayList<>(); + // Add 3 nodes WITH RATIS_DATASTREAM + for (int i = 0; i < 3; i++) { + nodes.add(createDatanodeDetails(true)); + } + // Add 3 nodes WITHOUT RATIS_DATASTREAM + for (int i = 0; i < 3; i++) { + nodes.add(createDatanodeDetails(false)); + } + + initWithNodes(1, nodes, 3); + Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); + assertEquals(3, pipeline.getNodes().size()); + for (DatanodeDetails dn : pipeline.getNodes()) { + assertTrue(dn.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM), + "Pipeline should only contain datanodes with RATIS_DATASTREAM when available"); + } + } + + @Test + public void testCreatePipelineFallsBackWhenNotEnoughRatisStreamingNodes() throws Exception { + List nodes = new ArrayList<>(); + // Add 2 nodes WITH RATIS_DATASTREAM + for (int i = 0; i < 2; i++) { + nodes.add(createDatanodeDetails(true)); + } + // Add 1 node WITHOUT RATIS_DATASTREAM + nodes.add(createDatanodeDetails(false)); + + initWithNodes(1, nodes, 3); + Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); + assertEquals(3, pipeline.getNodes().size()); + + long streamingNodeCount = pipeline.getNodes().stream() + .filter(dn -> dn.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM)) + .count(); + + assertEquals(2, streamingNodeCount, + "Pipeline should contain exactly 2 nodes with RATIS_DATASTREAM as fallback was required"); + } + @Test public void testCreatePipelinesDnExclude() throws Exception {