Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -163,17 +165,8 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig,
dns = pickNodesNotUsed(replicationConfig, minRatisVolumeSizeBytes, containerSizeBytes);
break;
case THREE:
List<DatanodeDetails> 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());
Expand Down Expand Up @@ -224,18 +217,39 @@ public Pipeline createForRead(
.build();
}

private List<DatanodeDetails> filterPipelineEngagement() {
private List<DatanodeDetails> chooseThreeFactorDatanodes(
List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, int requiredNode)
throws IOException {
final NodeManager nodeManager = getNodeManager();
final PipelineStateManager stateManager = getPipelineStateManager();
final List<DatanodeDetails> healthyNodes = nodeManager.getNodes(NodeStatus.inServiceHealthy());
final List<DatanodeDetails> excluded = new ArrayList<>();
excludedNodes = excludedNodes.isEmpty() ? new ArrayList<>() : excludedNodes;
List<DatanodeDetails> strictExcludedNodes = null;
for (DatanodeDetails d : healthyNodes) {
final int count = PipelinePlacementPolicy.currentRatisThreePipelineCount(nodeManager, stateManager, d);
if (count >= nodeManager.pipelineLimit(d)) {
excluded.add(d);
excludedNodes.add(d);
} else if (!d.hasPort(RATIS_DATASTREAM)) {
if (strictExcludedNodes == null) {
strictExcludedNodes = new ArrayList<>();
}
strictExcludedNodes.add(d);
}
}
return excluded;
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 placementPolicy.chooseDatanodes(excludedNodes,
favoredNodes, requiredNode, minRatisVolumeSizeBytes,
containerSizeBytes);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -93,10 +96,22 @@ public void init(int maxPipelinePerNode, OzoneConfiguration conf)
init(maxPipelinePerNode, conf, testDir);
}

public void initWithNodes(int maxPipelinePerNode, List<DatanodeDetails> 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,
Expand Down Expand Up @@ -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);
Comment on lines +264 to +269

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<DatanodeDetails> 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<DatanodeDetails> 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 {

Expand Down
Loading