From 02cba80d38ef96177483e34556429207b93e4fa3 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 09:45:14 +0800 Subject: [PATCH 01/11] HDDS-15969. Prioritize Ratis Streaming capable DataNodes during Pipeline creation with graceful fallback --- .../scm/pipeline/RatisPipelineProvider.java | 37 +++++--- .../pipeline/TestRatisPipelineProvider.java | 92 +++++++++++++++++++ 2 files changed, 116 insertions(+), 13 deletions(-) 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..f18a32c37e54 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,18 @@ 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(DatanodeDetails.Port.Name.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..97b3acb80e70 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 @@ -60,6 +60,7 @@ 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.hdds.scm.net.NetworkTopologyImpl; import org.apache.hadoop.ozone.ClientVersion; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assumptions; @@ -246,6 +247,97 @@ public void testCreateFactorTHREEPipelineWithSameDatanodes() assertEquals(pipeline2.getNodeSet(), pipeline3.getNodeSet()); } + private DatanodeDetails createDatanodeDetails(boolean supportRatisStreaming) { + DatanodeDetails.Builder dn = DatanodeDetails.newBuilder() + .setID(DatanodeID.randomID()) + .setHostName("localhost") + .setIpAddress("127.0.0.1") + .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 { + init(1); + 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)); + } + + // Initialize mock node manager with these nodes + nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 0); + nodeManager.setNumPipelinePerDatanode(1); + + // We must rebuild the provider with our custom nodeManager + SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT, 1); + stateManager = PipelineStateManagerImpl.newBuilder() + .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) + .setRatisServer(scmhaManager.getRatisServer()) + .setNodeManager(nodeManager) + .setSCMDBTransactionBuffer(scmhaManager.getDBTransactionBuffer()) + .build(); + provider = new MockRatisPipelineProvider(nodeManager, stateManager, conf); + + 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 { + init(1); + 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)); + + // Initialize mock node manager with these nodes + nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 0); + nodeManager.setNumPipelinePerDatanode(1); + + SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT, 1); + stateManager = PipelineStateManagerImpl.newBuilder() + .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) + .setRatisServer(scmhaManager.getRatisServer()) + .setNodeManager(nodeManager) + .setSCMDBTransactionBuffer(scmhaManager.getDBTransactionBuffer()) + .build(); + provider = new MockRatisPipelineProvider(nodeManager, stateManager, conf); + + 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 { From 64c7ff976e7e83fdbae2522ac32f612ecb6b0982 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 09:56:03 +0800 Subject: [PATCH 02/11] fix checkstyle --- .../hadoop/hdds/scm/pipeline/RatisPipelineProvider.java | 9 +++++---- .../hdds/scm/pipeline/TestRatisPipelineProvider.java | 2 +- 2 files changed, 6 insertions(+), 5 deletions(-) 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 f18a32c37e54..277e8f0c47a6 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 @@ -165,9 +165,9 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, dns = pickNodesNotUsed(replicationConfig, minRatisVolumeSizeBytes, containerSizeBytes); break; case THREE: - List excludeDueToEngagement = filterNodes(true); - List currentExcluded = new ArrayList<>(excludedNodes); - currentExcluded.addAll(excludeDueToEngagement); + List excludeDueToEngagement = filterNodes(true); + List currentExcluded = new ArrayList<>(excludedNodes); + currentExcluded.addAll(excludeDueToEngagement); try { dns = placementPolicy.chooseDatanodes(currentExcluded, favoredNodes, factor.getNumber(), minRatisVolumeSizeBytes, @@ -242,7 +242,8 @@ private List filterNodes(boolean filterRatisStreaming) { final List excluded = new ArrayList<>(); for (DatanodeDetails d : healthyNodes) { final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d); - if (count >= nodeManager.pipelineLimit(d) || (filterRatisStreaming && !d.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM))) { + if (count >= nodeManager.pipelineLimit(d) || + (filterRatisStreaming && !d.hasPort(DatanodeDetails.Port.Name.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 97b3acb80e70..e7ace0fdc908 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 @@ -57,10 +57,10 @@ 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.hdds.scm.net.NetworkTopologyImpl; import org.apache.hadoop.ozone.ClientVersion; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assumptions; From f80898661b7081b1fa90f29cb1852f4d94fdb26b Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 10:00:17 +0800 Subject: [PATCH 03/11] fix checkstyle --- .../apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 277e8f0c47a6..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 @@ -243,7 +243,7 @@ private List filterNodes(boolean filterRatisStreaming) { for (DatanodeDetails d : healthyNodes) { final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d); if (count >= nodeManager.pipelineLimit(d) || - (filterRatisStreaming && !d.hasPort(DatanodeDetails.Port.Name.RATIS_DATASTREAM))) { + (filterRatisStreaming && !d.hasPort(RATIS_DATASTREAM))) { excluded.add(d); } } From 38e12e562692c147e800224044f988802cc70d16 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 10:13:53 +0800 Subject: [PATCH 04/11] update TestRatisPipelineProvider --- .../hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 e7ace0fdc908..ba25063cedb6 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 @@ -278,7 +278,7 @@ public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception } // Initialize mock node manager with these nodes - nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 0); + MockNodeManager nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 3); nodeManager.setNumPipelinePerDatanode(1); // We must rebuild the provider with our custom nodeManager @@ -313,7 +313,7 @@ public void testCreatePipelineFallsBackWhenNotEnoughRatisStreamingNodes() throws nodes.add(createDatanodeDetails(false)); // Initialize mock node manager with these nodes - nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 0); + MockNodeManager nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 3); nodeManager.setNumPipelinePerDatanode(1); SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); From 0986f0c88b74abd9985c59f0642e640bae010e46 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 10:39:41 +0800 Subject: [PATCH 05/11] refactor TestRatisPipelineProvider --- .../pipeline/TestRatisPipelineProvider.java | 60 +++++++++---------- 1 file changed, 27 insertions(+), 33 deletions(-) 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 ba25063cedb6..5c3c0cc46d60 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 @@ -94,6 +94,31 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf) init(maxPipelinePerNode, conf, testDir); } + public void initWithNodes(int maxPipelinePerNode, List nodes, int nodeCount) + throws Exception { + + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(HddsConfigKeys.OZONE_METADATA_DIRS, testDir.getAbsolutePath()); + dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get()); + nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, nodeCount); + nodeManager.setNumPipelinePerDatanode(maxPipelinePerNode); + long containerSize = (long) conf.getStorageSize( + ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE, + ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT, StorageUnit.BYTES); + nodeManager.setPendingContainerMaxSize(containerSize); + SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); + conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT, + maxPipelinePerNode); + stateManager = PipelineStateManagerImpl.newBuilder() + .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) + .setRatisServer(scmhaManager.getRatisServer()) + .setNodeManager(nodeManager) + .setSCMDBTransactionBuffer(scmhaManager.getDBTransactionBuffer()) + .build(); + provider = new MockRatisPipelineProvider(nodeManager, + stateManager, 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()); @@ -266,7 +291,6 @@ private DatanodeDetails createDatanodeDetails(boolean supportRatisStreaming) { @Test public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception { - init(1); List nodes = new ArrayList<>(); // Add 3 nodes WITH RATIS_DATASTREAM for (int i = 0; i < 3; i++) { @@ -277,22 +301,7 @@ public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception nodes.add(createDatanodeDetails(false)); } - // Initialize mock node manager with these nodes - MockNodeManager nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 3); - nodeManager.setNumPipelinePerDatanode(1); - - // We must rebuild the provider with our custom nodeManager - SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); - OzoneConfiguration conf = new OzoneConfiguration(); - conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT, 1); - stateManager = PipelineStateManagerImpl.newBuilder() - .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) - .setRatisServer(scmhaManager.getRatisServer()) - .setNodeManager(nodeManager) - .setSCMDBTransactionBuffer(scmhaManager.getDBTransactionBuffer()) - .build(); - provider = new MockRatisPipelineProvider(nodeManager, stateManager, conf); - + initWithNodes(1, nodes, 3); Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); assertEquals(3, pipeline.getNodes().size()); for (DatanodeDetails dn : pipeline.getNodes()) { @@ -303,7 +312,6 @@ public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception @Test public void testCreatePipelineFallsBackWhenNotEnoughRatisStreamingNodes() throws Exception { - init(1); List nodes = new ArrayList<>(); // Add 2 nodes WITH RATIS_DATASTREAM for (int i = 0; i < 2; i++) { @@ -312,21 +320,7 @@ public void testCreatePipelineFallsBackWhenNotEnoughRatisStreamingNodes() throws // Add 1 node WITHOUT RATIS_DATASTREAM nodes.add(createDatanodeDetails(false)); - // Initialize mock node manager with these nodes - MockNodeManager nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, 3); - nodeManager.setNumPipelinePerDatanode(1); - - SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); - OzoneConfiguration conf = new OzoneConfiguration(); - conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT, 1); - stateManager = PipelineStateManagerImpl.newBuilder() - .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) - .setRatisServer(scmhaManager.getRatisServer()) - .setNodeManager(nodeManager) - .setSCMDBTransactionBuffer(scmhaManager.getDBTransactionBuffer()) - .build(); - provider = new MockRatisPipelineProvider(nodeManager, stateManager, conf); - + initWithNodes(1, nodes, 3); Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); assertEquals(3, pipeline.getNodes().size()); From 53c0e2a4319242ef71d3f1dc105ba5ca863c6ad1 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 10:44:13 +0800 Subject: [PATCH 06/11] refactor TestRatisPipelineProvider --- .../pipeline/TestRatisPipelineProvider.java | 25 +++++-------------- 1 file changed, 6 insertions(+), 19 deletions(-) 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 5c3c0cc46d60..201218e3e7f3 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 @@ -96,33 +96,20 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf) public void initWithNodes(int maxPipelinePerNode, List nodes, int nodeCount) throws Exception { - OzoneConfiguration conf = new OzoneConfiguration(); conf.set(HddsConfigKeys.OZONE_METADATA_DIRS, testDir.getAbsolutePath()); - dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get()); nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, nodeCount); - nodeManager.setNumPipelinePerDatanode(maxPipelinePerNode); - long containerSize = (long) conf.getStorageSize( - ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE, - ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT, StorageUnit.BYTES); - nodeManager.setPendingContainerMaxSize(containerSize); - SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true); - conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT, - maxPipelinePerNode); - stateManager = PipelineStateManagerImpl.newBuilder() - .setPipelineStore(SCMDBDefinition.PIPELINES.getTable(dbStore)) - .setRatisServer(scmhaManager.getRatisServer()) - .setNodeManager(nodeManager) - .setSCMDBTransactionBuffer(scmhaManager.getDBTransactionBuffer()) - .build(); - provider = new MockRatisPipelineProvider(nodeManager, - stateManager, conf); + 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, From 3b88b99b23c5c73c021c9f5be53519232653b5c7 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 10:57:24 +0800 Subject: [PATCH 07/11] update TestRatisPipelineProvider --- .../scm/pipeline/TestRatisPipelineProvider.java | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) 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 201218e3e7f3..811e3b4e42bb 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; @@ -260,10 +262,17 @@ public void testCreateFactorTHREEPipelineWithSameDatanodes() } 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") - .setIpAddress("127.0.0.1") + .setHostName("localhost" + "-" + ipAddress) + .setIpAddress(ipAddress) + .setNetworkLocation(null) .setPersistedOpState(HddsProtos.NodeOperationalState.IN_SERVICE) .setPersistedOpStateExpiry(0); From 67a8035b987b6c9f0707794917a6d4bfa49ae8bc Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 10:59:38 +0800 Subject: [PATCH 08/11] fix checkstyle --- .../hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 811e3b4e42bb..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 @@ -96,11 +96,11 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf) init(maxPipelinePerNode, conf, testDir); } - public void initWithNodes(int maxPipelinePerNode, List nodes, int nodeCount) + 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, nodeCount); + nodeManager = new MockNodeManager(new NetworkTopologyImpl(new OzoneConfiguration()), nodes, false, count); initializeCommonState(maxPipelinePerNode, conf); } From cecb9b04a8c1956dd6260736afa24c2229f3353a Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 28 Jul 2026 12:27:08 +0800 Subject: [PATCH 09/11] address comment --- .../scm/pipeline/RatisPipelineProvider.java | 54 ++++++++++--------- 1 file changed, 28 insertions(+), 26 deletions(-) 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 583ef5cac17b..1746cb3733a0 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 @@ -157,7 +157,7 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, ); } - List dns; + final List dns; final ReplicationFactor factor = replicationConfig.getReplicationFactor(); switch (factor) { @@ -165,21 +165,7 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, dns = pickNodesNotUsed(replicationConfig, minRatisVolumeSizeBytes, containerSizeBytes); break; case THREE: - 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 = chooseThreeFactorDatanodes(excludedNodes, favoredNodes, factor.getNumber()); break; default: @@ -231,23 +217,39 @@ public Pipeline createForRead( .build(); } - /** - * - * @return - */ - private List filterNodes(boolean filterRatisStreaming) { + 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() ? new ArrayList<>() : excludedNodes; + List strictExcludedNodes = null; for (DatanodeDetails d : healthyNodes) { final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d); - if (count >= nodeManager.pipelineLimit(d) || - (filterRatisStreaming && !d.hasPort(RATIS_DATASTREAM))) { - excluded.add(d); + if (count >= nodeManager.pipelineLimit(d)) { + excludedNodes.add(d); + } else if (!d.hasPort(RATIS_DATASTREAM)) { + if (strictExcludedNodes == null) { + strictExcludedNodes = new ArrayList<>(); + } + strictExcludedNodes.add(d); + } + } + if (strictExcludedNodes != null) { + strictExcludedNodes.addAll(excludedNodes); + try { + return placementPolicy.chooseDatanodes(strictExcludedNodes, + 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); } /** From 54ea91941966d48ead40f7aef97559b6899c29d9 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Tue, 28 Jul 2026 13:01:17 +0800 Subject: [PATCH 10/11] fix checkstyle --- .../apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 1746cb3733a0..a0d7d8d46e60 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 @@ -219,7 +219,7 @@ public Pipeline createForRead( private List chooseThreeFactorDatanodes( List excludedNodes, List favoredNodes, int requiredNode) - throws IOException{ + throws IOException { final NodeManager nodeManager = getNodeManager(); final PipelineStateManager stateManager = getPipelineStateManager(); final List healthyNodes = nodeManager.getNodes(NodeStatus.inServiceHealthy()); From 50dce0b8231690b6a9d30a9c7750da2f4200bbc2 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Wed, 29 Jul 2026 18:46:18 +0800 Subject: [PATCH 11/11] address comments --- .../scm/pipeline/RatisPipelineProvider.java | 29 +++++++++++++------ .../pipeline/TestRatisPipelineProvider.java | 14 ++++++--- 2 files changed, 30 insertions(+), 13 deletions(-) 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 a0d7d8d46e60..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 @@ -42,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; @@ -66,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, @@ -92,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); @@ -223,23 +228,29 @@ private List chooseThreeFactorDatanodes( final NodeManager nodeManager = getNodeManager(); final PipelineStateManager stateManager = getPipelineStateManager(); final List healthyNodes = nodeManager.getNodes(NodeStatus.inServiceHealthy()); - excludedNodes = excludedNodes.isEmpty() ? new ArrayList<>() : excludedNodes; - List strictExcludedNodes = null; + 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)) { + if (excludedNodes == null) { + excludedNodes = new ArrayList<>(); + } excludedNodes.add(d); - } else if (!d.hasPort(RATIS_DATASTREAM)) { - if (strictExcludedNodes == null) { - strictExcludedNodes = new ArrayList<>(); + } else if (isRatisStreamingEnabled && !d.hasPort(RATIS_DATASTREAM)) { + if (additionalExcludedNodes == null) { + additionalExcludedNodes = new ArrayList<>(); } - strictExcludedNodes.add(d); + additionalExcludedNodes.add(d); } } - if (strictExcludedNodes != null) { - strictExcludedNodes.addAll(excludedNodes); + // 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(strictExcludedNodes, + return placementPolicy.chooseDatanodes(additionalExcludedNodes, favoredNodes, requiredNode, minRatisVolumeSizeBytes, containerSizeBytes); } catch (SCMException scmException) { 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 ab2a1c3d77d4..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 @@ -64,6 +64,7 @@ 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; @@ -96,9 +97,8 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf) init(maxPipelinePerNode, conf, testDir); } - public void initWithNodes(int maxPipelinePerNode, List nodes, int count) + public void initWithNodes(int maxPipelinePerNode, OzoneConfiguration conf, 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); @@ -287,6 +287,9 @@ private DatanodeDetails createDatanodeDetails(boolean supportRatisStreaming) { @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++) { @@ -297,7 +300,7 @@ public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception nodes.add(createDatanodeDetails(false)); } - initWithNodes(1, nodes, 3); + initWithNodes(1, conf, nodes, 3); Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); assertEquals(3, pipeline.getNodes().size()); for (DatanodeDetails dn : pipeline.getNodes()) { @@ -308,6 +311,9 @@ public void testCreatePipelinePrioritizesRatisStreamingNodes() throws Exception @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++) { @@ -316,7 +322,7 @@ public void testCreatePipelineFallsBackWhenNotEnoughRatisStreamingNodes() throws // Add 1 node WITHOUT RATIS_DATASTREAM nodes.add(createDatanodeDetails(false)); - initWithNodes(1, nodes, 3); + initWithNodes(1, conf, nodes, 3); Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); assertEquals(3, pipeline.getNodes().size());