Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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 @@ -155,25 +157,30 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig,
);
}

final List<DatanodeDetails> dns;
List<DatanodeDetails> dns;
final ReplicationFactor factor =
replicationConfig.getReplicationFactor();
switch (factor) {
case ONE:
dns = pickNodesNotUsed(replicationConfig, minRatisVolumeSizeBytes, containerSizeBytes);
break;
case THREE:
List<DatanodeDetails> excludeDueToEngagement = filterPipelineEngagement();
if (!excludeDueToEngagement.isEmpty()) {
if (excludedNodes.isEmpty()) {
excludedNodes = excludeDueToEngagement;
} else {
excludedNodes.addAll(excludeDueToEngagement);
}
List<DatanodeDetails> excludeDueToEngagement = filterNodes(true);
List<DatanodeDetails> 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());
Expand Down Expand Up @@ -224,14 +231,19 @@ public Pipeline createForRead(
.build();
}

private List<DatanodeDetails> filterPipelineEngagement() {
/**
*
* @return
*/
private List<DatanodeDetails> filterNodes(boolean filterRatisStreaming) {
final NodeManager nodeManager = getNodeManager();
final PipelineStateManager stateManager = getPipelineStateManager();
final List<DatanodeDetails> healthyNodes = nodeManager.getNodes(NodeStatus.inServiceHealthy());
final List<DatanodeDetails> 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);
}
}
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);

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