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..35512b2e4f3c 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; @@ -40,6 +42,7 @@ import org.apache.hadoop.hdds.scm.pipeline.leader.choose.algorithms.LeaderChoosePolicy; import org.apache.hadoop.hdds.scm.pipeline.leader.choose.algorithms.LeaderChoosePolicyFactory; import org.apache.hadoop.hdds.server.events.EventPublisher; +import org.apache.hadoop.ozone.OzoneConfigKeys; import org.apache.hadoop.ozone.protocol.commands.ClosePipelineCommand; import org.apache.hadoop.ozone.protocol.commands.CommandForDatanode; import org.apache.hadoop.ozone.protocol.commands.CreatePipelineCommand; @@ -64,6 +67,7 @@ public class RatisPipelineProvider private final SCMContext scmContext; private final long containerSizeBytes; private final long minRatisVolumeSizeBytes; + private final boolean isRatisStreamingEnabled; @VisibleForTesting public RatisPipelineProvider(NodeManager nodeManager, @@ -90,6 +94,9 @@ public RatisPipelineProvider(NodeManager nodeManager, ScmConfigKeys.OZONE_DATANODE_RATIS_VOLUME_FREE_SPACE_MIN, ScmConfigKeys.OZONE_DATANODE_RATIS_VOLUME_FREE_SPACE_MIN_DEFAULT, StorageUnit.BYTES); + this.isRatisStreamingEnabled = conf.getBoolean( + OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, + OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED_DEFAULT); try { leaderChoosePolicy = LeaderChoosePolicyFactory .getPolicy(conf, nodeManager, stateManager); @@ -163,17 +170,8 @@ 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); - } - } - dns = placementPolicy.chooseDatanodes(excludedNodes, - favoredNodes, factor.getNumber(), minRatisVolumeSizeBytes, - containerSizeBytes); + dns = chooseThreeFactorDatanodes(excludedNodes, favoredNodes, factor.getNumber()); + break; default: throw new IllegalStateException("Unknown factor: " + factor.name()); @@ -224,18 +222,45 @@ public Pipeline createForRead( .build(); } - private List filterPipelineEngagement() { + private List chooseThreeFactorDatanodes( + List excludedNodes, List favoredNodes, int requiredNode) + throws IOException { final NodeManager nodeManager = getNodeManager(); final PipelineStateManager stateManager = getPipelineStateManager(); final List healthyNodes = nodeManager.getNodes(NodeStatus.inServiceHealthy()); - final List excluded = new ArrayList<>(); + excludedNodes = excludedNodes.isEmpty() ? null : excludedNodes; + List additionalExcludedNodes = null; for (DatanodeDetails d : healthyNodes) { final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d); if (count >= nodeManager.pipelineLimit(d)) { - excluded.add(d); + if (excludedNodes == null) { + excludedNodes = new ArrayList<>(); + } + excludedNodes.add(d); + } else if (isRatisStreamingEnabled && !d.hasPort(RATIS_DATASTREAM)) { + if (additionalExcludedNodes == null) { + additionalExcludedNodes = new ArrayList<>(); + } + additionalExcludedNodes.add(d); + } + } + // If the cluster does not support Ratis streaming, or if all nodes support it, no fallback will occur. + if (additionalExcludedNodes != null) { + if (excludedNodes != null) { + additionalExcludedNodes.addAll(excludedNodes); + } + try { + return placementPolicy.chooseDatanodes(additionalExcludedNodes, + favoredNodes, requiredNode, minRatisVolumeSizeBytes, + containerSizeBytes); + } catch (SCMException scmException) { + LOG.debug("Failed to allocate datanodes with strict exclusion (non-streaming nodes excluded)." + + " Falling back.", scmException); } } - return excluded; + return placementPolicy.chooseDatanodes(excludedNodes, + favoredNodes, requiredNode, minRatisVolumeSizeBytes, + containerSizeBytes); } /** 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..5017a65aca42 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,10 +59,12 @@ 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; import org.apache.hadoop.ozone.ClientVersion; +import org.apache.hadoop.ozone.OzoneConfigKeys; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assumptions; import org.junit.jupiter.api.Test; @@ -93,10 +97,21 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf) init(maxPipelinePerNode, conf, testDir); } + public void initWithNodes(int maxPipelinePerNode, OzoneConfiguration conf, List nodes, int count) + throws Exception { + 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,79 @@ 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 { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setBoolean( + OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true); + 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, conf, 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 { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setBoolean( + OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_ENABLED, true); + 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, conf, 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 {