From 02cba80d38ef96177483e34556429207b93e4fa3 Mon Sep 17 00:00:00 2001 From: chungen0126 Date: Sat, 25 Jul 2026 09:45:14 +0800 Subject: [PATCH 1/8] 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 2/8] 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 3/8] 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 4/8] 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 5/8] 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 6/8] 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 7/8] 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 8/8] 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); }