This is an automated email from the ASF dual-hosted git repository.
adoroszlai pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 66d397813cb HDDS-15146. Remove getDatanodeInfo(DatanodeDetails) from
NodeManager (#10601)
66d397813cb is described below
commit 66d397813cbb60a9e7e81e5808021f3438d65f4d
Author: Navink <[email protected]>
AuthorDate: Sat Jun 27 00:21:59 2026 +0530
HDDS-15146. Remove getDatanodeInfo(DatanodeDetails) from NodeManager
(#10601)
Co-authored-by: Tsz-Wo Nicholas Sze <[email protected]>
---
.../hdds/scm/container/ContainerReportHandler.java | 10 ++-
.../apache/hadoop/hdds/scm/node/NodeManager.java | 19 ++---
.../hadoop/hdds/scm/node/NodeStateManager.java | 8 +--
.../hadoop/hdds/scm/node/SCMNodeManager.java | 32 ++-------
.../hadoop/hdds/scm/node/states/NodeStateMap.java | 10 +--
.../hdds/scm/pipeline/PipelineManagerImpl.java | 3 +-
.../hdds/scm/server/SCMClientProtocolServer.java | 81 ++++------------------
.../hdds/scm/TestSCMCommonPlacementPolicy.java | 42 ++++++-----
.../hadoop/hdds/scm/container/MockNodeManager.java | 33 ++++++---
.../hdds/scm/container/SimpleMockNodeManager.java | 15 ++--
.../scm/container/TestContainerReportHandler.java | 14 ++--
.../scm/container/TestContainerStateManager.java | 6 +-
.../TestReconcileContainerEventHandler.java | 4 +-
.../replication/TestReplicationManagerUtil.java | 10 ++-
.../hdds/scm/pipeline/MockPipelineManager.java | 2 +-
.../scm/node/TestDecommissionAndMaintenance.java | 6 +-
.../dn/TestDatanodeMinFreeSpaceIntegration.java | 4 +-
.../hadoop/ozone/scm/node/TestDiskBalancer.java | 2 +-
...skBalancerDuringDecommissionAndMaintenance.java | 9 +--
.../hadoop/ozone/recon/scm/PipelineSyncTask.java | 9 ++-
...TestReconIncrementalContainerReportHandler.java | 15 ++--
.../ozone/recon/scm/TestReconNodeManager.java | 3 +-
22 files changed, 129 insertions(+), 208 deletions(-)
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReportHandler.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReportHandler.java
index fe3438b1ea7..2326cd894e6 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReportHandler.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReportHandler.java
@@ -136,11 +136,12 @@ public void onMessage(final ContainerReportFromDatanode
reportFromDatanode,
final DatanodeDetails dnFromReport =
reportFromDatanode.getDatanodeDetails();
- final DatanodeDetails datanodeDetails =
getNodeManager().getNode(dnFromReport.getID());
- if (datanodeDetails == null) {
+ final DatanodeInfo datanodeInfo =
getNodeManager().getNode(dnFromReport.getID());
+ if (datanodeInfo == null) {
getLogger().warn("Datanode not found: {}", dnFromReport);
return;
}
+ final DatanodeDetails datanodeDetails = datanodeInfo;
final ContainerReportsProto containerReport =
reportFromDatanode.getReport();
try {
@@ -153,7 +154,6 @@ public void onMessage(final ContainerReportFromDatanode
reportFromDatanode,
containerReport.getReportsList();
final Set<ContainerID> expectedContainersInDatanode =
getNodeManager().getContainers(datanodeDetails);
- DatanodeInfo datanodeInfo = datanodeDetails instanceof DatanodeInfo ?
(DatanodeInfo) datanodeDetails : null;
for (ContainerReplicaProto replica : replicas) {
ContainerID cid = ContainerID.valueOf(replica.getContainerID());
@@ -178,9 +178,7 @@ public void onMessage(final ContainerReportFromDatanode
reportFromDatanode,
getNodeManager().addContainer(datanodeDetails, cid);
// Remove from pending tracker when container is added to DN
// This container was just confirmed for the first time on this DN
- if (datanodeInfo != null) {
-
getNodeManager().removePendingAllocationForDatanode(datanodeInfo, cid);
- }
+ getNodeManager().removePendingAllocationForDatanode(datanodeInfo,
cid);
}
if (container == null || ContainerReportValidator
.validate(container, datanodeDetails, replica)) {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
index 51eb0e4e41f..49af3cae740 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
@@ -113,7 +113,7 @@ default void registerSendCommandNotify(SCMCommandProto.Type
type,
* @param health - The health of the node
* @return List of Datanodes that are Heartbeating SCM.
*/
- List<DatanodeDetails> getNodes(
+ List<DatanodeInfo> getNodes(
NodeOperationalState opState, NodeState health);
/**
@@ -134,10 +134,8 @@ List<DatanodeDetails> getNodes(
int getNodeCount(
NodeOperationalState opState, NodeState health);
- /**
- * @return all datanodes known to SCM.
- */
- List<? extends DatanodeDetails> getAllNodes();
+ /** @return a shadow copied list of all datanodes, sorted by {@link
DatanodeID}. */
+ List<DatanodeInfo> getAllNodes();
/** @return the number of datanodes. */
default int getAllNodeCount() {
@@ -175,15 +173,6 @@ default int getAllNodeCount() {
*/
DatanodeUsageInfo getUsageInfo(DatanodeDetails dn);
- /**
- * Get the datanode info of a specified datanode.
- *
- * @param dn the usage of which we want to get
- * @return DatanodeInfo of the specified datanode
- */
- @Nullable
- DatanodeInfo getDatanodeInfo(DatanodeDetails dn);
-
/**
* Atomically checks if the datanode has space for a new container and
records the allocation
* if space is available. This prevents race conditions where multiple
threads check space
@@ -404,7 +393,7 @@ Map<SCMCommandProto.Type, Integer>
getTotalDatanodeCommandCounts(
List<SCMCommand<?>> getCommandQueue(DatanodeID dnID);
/** @return the datanode of the given id if it exists; otherwise, return
null. */
- @Nullable DatanodeDetails getNode(@Nullable DatanodeID id);
+ @Nullable DatanodeInfo getNode(@Nullable DatanodeID id);
/**
* Given datanode address(Ipaddress or hostname), returns a list of
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
index 57ee8ef9adb..9539379fd84 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
@@ -537,12 +537,8 @@ public int getVolumeFailuresNodeCount() {
return getVolumeFailuresNodes().size();
}
- /**
- * Returns all the nodes which have registered to NodeStateManager.
- *
- * @return all the managed nodes
- */
- public List<DatanodeInfo> getAllNodes() {
+ /** @return a shadow copied list of all datanodes, sorted by {@link
DatanodeID}. */
+ List<DatanodeInfo> getAllNodes() {
return nodeStateMap.getAllDatanodeInfos();
}
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
index 394793bb9a3..cbcfc553791 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
@@ -26,7 +26,6 @@
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
import com.google.common.base.Strings;
-import jakarta.annotation.Nullable;
import java.io.IOException;
import java.math.RoundingMode;
import java.net.InetAddress;
@@ -270,11 +269,9 @@ public List<DatanodeDetails> getNodes(NodeStatus
nodeStatus) {
* @return List of Datanodes that are known to SCM in the requested states.
*/
@Override
- public List<DatanodeDetails> getNodes(
+ public List<DatanodeInfo> getNodes(
NodeOperationalState opState, NodeState health) {
- return nodeStateManager.getNodes(opState, health)
- .stream()
- .map(node -> (DatanodeDetails)node).collect(Collectors.toList());
+ return nodeStateManager.getNodes(opState, health);
}
@Override
@@ -1015,8 +1012,7 @@ public Map<DatanodeDetails, SCMNodeStat> getNodeStats() {
@Override
public List<DatanodeUsageInfo> getMostOrLeastUsedDatanodes(
boolean mostUsed) {
- List<DatanodeDetails> healthyNodes =
- getNodes(IN_SERVICE, NodeState.HEALTHY);
+ final List<DatanodeInfo> healthyNodes = getNodes(IN_SERVICE,
NodeState.HEALTHY);
List<DatanodeUsageInfo> datanodeUsageInfoList =
new ArrayList<>(healthyNodes.size());
@@ -1063,24 +1059,6 @@ public DatanodeUsageInfo getUsageInfo(DatanodeDetails
dn) {
return usageInfo;
}
- /**
- * Get the usage info of a specified datanode.
- *
- * @param dn the usage of which we want to get
- * @return DatanodeUsageInfo of the specified datanode
- */
- @Override
- @Nullable
- public DatanodeInfo getDatanodeInfo(DatanodeDetails dn) {
- try {
- return nodeStateManager.getNode(dn);
- } catch (NodeNotFoundException e) {
- LOG.warn("Cannot retrieve DatanodeInfo, datanode {} not found.",
- dn.getID());
- return null;
- }
- }
-
@Override
public boolean checkSpaceAndRecordAllocation(DatanodeInfo datanodeInfo,
ContainerID containerID) {
return pendingContainerTracker.checkSpaceAndRecordAllocation(datanodeInfo,
containerID);
@@ -1105,7 +1083,7 @@ public void
removePendingAllocationForDatanode(DatanodeInfo datanodeInfo, Contai
@Override
public void opAdded(ContainerReplicaOp op, ContainerID containerID) {
if (op.getOpType() == ContainerReplicaOp.PendingOpType.ADD) {
- DatanodeInfo dnInfo = getDatanodeInfo(op.getTarget());
+ DatanodeInfo dnInfo = getNode(op.getTarget().getID());
if (dnInfo != null) {
recordAllocationForDatanode(dnInfo, containerID);
}
@@ -1115,7 +1093,7 @@ public void opAdded(ContainerReplicaOp op, ContainerID
containerID) {
@Override
public void opCompleted(ContainerReplicaOp op, ContainerID containerID,
boolean timedOut) {
if (op.getOpType() == ContainerReplicaOp.PendingOpType.ADD && !timedOut) {
- DatanodeInfo dnInfo = getDatanodeInfo(op.getTarget());
+ DatanodeInfo dnInfo = getNode(op.getTarget().getID());
if (dnInfo != null) {
removePendingAllocationForDatanode(dnInfo, containerID);
}
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/states/NodeStateMap.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/states/NodeStateMap.java
index 8aea57b23ab..67c487e3ba1 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/states/NodeStateMap.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/states/NodeStateMap.java
@@ -17,10 +17,10 @@
package org.apache.hadoop.hdds.scm.node.states;
-import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.TreeMap;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.Function;
@@ -41,7 +41,7 @@
*/
public class NodeStateMap {
/** Map: {@link DatanodeID} -> ({@link DatanodeInfo}, {@link ContainerID}s).
*/
- private final Map<DatanodeID, DatanodeEntry> nodeMap = new HashMap<>();
+ private final Map<DatanodeID, DatanodeEntry> nodeMap = new TreeMap<>();
private final ReadWriteLock lock = new ReentrantReadWriteLock();
@@ -166,11 +166,7 @@ public int getNodeCount() {
}
}
- /**
- * Returns the list of all the nodes as DatanodeInfo objects.
- *
- * @return list of all the node ids
- */
+ /** @return a shadow copied list of all datanodes, sorted by {@link
DatanodeID}. */
public List<DatanodeInfo> getAllDatanodeInfos() {
lock.readLock().lock();
try {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
index fc8229c677d..8003d8d52e9 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java
@@ -638,7 +638,8 @@ public boolean checkSpaceAndRecordAllocation(Pipeline
pipeline, ContainerID cont
final Set<DatanodeDetails> datanodeDetails = pipeline.getNodeSet();
final List<DatanodeInfo> datanodeInfos = new
ArrayList<>(datanodeDetails.size());
for (DatanodeDetails dn : datanodeDetails) {
- final DatanodeInfo info = nodeManager.getDatanodeInfo(dn);
+ // Refactored to use getNode instead of getDatanodeInfo
+ final DatanodeInfo info = nodeManager.getNode(dn.getID());
if (info == null) {
LOG.warn("DatanodeInfo not found for {}", dn.getID());
return false;
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
index 010e0ca8b6c..6883cee0127 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
@@ -46,7 +46,6 @@
import java.util.Map;
import java.util.Optional;
import java.util.Set;
-import java.util.TreeSet;
import java.util.UUID;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -691,9 +690,9 @@ public List<HddsProtos.Node> queryNode(
}
try {
List<HddsProtos.Node> result = new ArrayList<>();
- for (DatanodeDetails node : queryNode(opState, state)) {
+ for (DatanodeDetails node : scm.getScmNodeManager().getNodes(opState,
state)) {
NodeStatus ns = scm.getScmNodeManager().getNodeStatus(node);
- DatanodeInfo datanodeInfo =
scm.getScmNodeManager().getDatanodeInfo(node);
+ DatanodeInfo datanodeInfo = node instanceof DatanodeInfo ?
(DatanodeInfo) node : null;
HddsProtos.Node.Builder nodeBuilder = HddsProtos.Node.newBuilder()
.setNodeID(node.toProto(clientVersion))
.addNodeStates(ns.getHealth())
@@ -717,34 +716,22 @@ public List<HddsProtos.Node> queryNode(
}
@Override
- public HddsProtos.Node queryNode(UUID uuid)
- throws IOException {
+ public HddsProtos.Node queryNode(UUID uuid) {
final Map<String, String> auditMap = Maps.newHashMap();
auditMap.put("uuid", String.valueOf(uuid));
HddsProtos.Node result = null;
- try {
- DatanodeDetails node =
scm.getScmNodeManager().getNode(DatanodeID.of(uuid));
- if (node != null) {
- NodeStatus ns = scm.getScmNodeManager().getNodeStatus(node);
- DatanodeInfo datanodeInfo =
scm.getScmNodeManager().getDatanodeInfo(node);
- HddsProtos.Node.Builder nodeBuilder = HddsProtos.Node.newBuilder()
- .setNodeID(node.getProtoBufMessage())
- .addNodeStates(ns.getHealth())
- .addNodeOperationalStates(ns.getOperationalState());
-
- if (datanodeInfo != null) {
-
nodeBuilder.setTotalVolumeCount(datanodeInfo.getStorageReports().size());
-
nodeBuilder.setHealthyVolumeCount(datanodeInfo.getHealthyVolumeCount());
- addFailedVolumes(nodeBuilder, datanodeInfo);
- }
- result = nodeBuilder.build();
- }
- } catch (NodeNotFoundException e) {
- IOException ex = new IOException(
- "An unexpected error occurred querying the NodeStatus", e);
- AUDIT.logReadFailure(buildAuditMessageForFailure(
- SCMAction.QUERY_NODE, auditMap, ex));
- throw ex;
+ DatanodeInfo datanodeInfo =
scm.getScmNodeManager().getNode(DatanodeID.of(uuid));
+ if (datanodeInfo != null) {
+ NodeStatus ns = datanodeInfo.getNodeStatus();
+ HddsProtos.Node.Builder nodeBuilder = HddsProtos.Node.newBuilder()
+ .setNodeID(datanodeInfo.getProtoBufMessage())
+ .addNodeStates(ns.getHealth())
+ .addNodeOperationalStates(ns.getOperationalState());
+
+ nodeBuilder.setTotalVolumeCount(datanodeInfo.getStorageReports().size());
+ nodeBuilder.setHealthyVolumeCount(datanodeInfo.getHealthyVolumeCount());
+ addFailedVolumes(nodeBuilder, datanodeInfo);
+ result = nodeBuilder.build();
}
AUDIT.logReadSuccess(buildAuditMessageForSuccess(
SCMAction.QUERY_NODE, auditMap));
@@ -1593,26 +1580,6 @@ public List<ContainerID> getListOfContainerIDs(
}
}
- /**
- * Queries a list of Node that match a set of statuses.
- *
- * <p>For example, if the nodeStatuses is HEALTHY and RAFT_MEMBER, then
- * this call will return all
- * healthy nodes which members in Raft pipeline.
- *
- * <p>Right now we don't support operations, so we assume it is an AND
- * operation between the
- * operators.
- *
- * @param opState - NodeOperational State
- * @param state - NodeState.
- * @return List of Datanodes.
- */
- public List<DatanodeDetails> queryNode(
- HddsProtos.NodeOperationalState opState, HddsProtos.NodeState state) {
- return new ArrayList<>(queryNodeState(opState, state));
- }
-
@VisibleForTesting
public StorageContainerManager getScm() {
return scm;
@@ -1625,24 +1592,6 @@ public boolean getSafeModeStatus() {
return scm.getScmContext().isInSafeMode();
}
- /**
- * Query the System for Nodes.
- *
- * @params opState - The node operational state
- * @param nodeState - NodeState that we are interested in matching.
- * @return Set of Datanodes that match the NodeState.
- */
- private Set<DatanodeDetails> queryNodeState(
- HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodeState)
{
- Set<DatanodeDetails> returnSet = new TreeSet<>();
- List<DatanodeDetails> tmp = scm.getScmNodeManager()
- .getNodes(opState, nodeState);
- if ((tmp != null) && (!tmp.isEmpty())) {
- returnSet.addAll(tmp);
- }
- return returnSet;
- }
-
@Override
public AuditMessage buildAuditMessageForSuccess(
AuditAction op, Map<String, String> auditMap) {
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
index 9290e8d83be..9df7c4f69d6 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
@@ -33,6 +33,7 @@
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Sets;
import java.io.File;
+import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
@@ -84,11 +85,15 @@ void setup(@TempDir File testDir) {
conf = SCMTestUtils.getConf(testDir);
}
+ static List<DatanodeDetails> getAllNodes(NodeManager nm) {
+ return new ArrayList<>(nm.getAllNodes());
+ }
+
@Test
public void testGetResultSet() throws SCMException {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> result = dummyPlacementPolicy.getResultSet(3, list);
Set<DatanodeDetails> resultSet = new HashSet<>(result);
assertNotEquals(1, resultSet.size());
@@ -136,7 +141,7 @@ public void
testReplicasToFixMisreplicationWithOneMisreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 5)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -157,7 +162,7 @@ public void
testReplicasToFixMisreplicationWithTwoMisreplication() {
3, ImmutableList.of(3, 8),
4, ImmutableList.of(4, 9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 5)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -178,7 +183,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 8),
4, ImmutableList.of(4, 9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 5)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -200,7 +205,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 4, 8),
4, ImmutableList.of(9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 4)
.map(list::get).collect(Collectors.toList());
//Creating Replicas without replica Index
@@ -223,7 +228,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 4, 8),
4, ImmutableList.of(9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 3, 4)
.map(list::get).collect(Collectors.toList());
//Creating Replicas without replica Index for replicas < number of racks
@@ -246,7 +251,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 4, 8),
4, ImmutableList.of(9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 4, 6)
.map(list::get).collect(Collectors.toList());
//Creating Replicas without replica Index for replicas >number of racks
@@ -261,7 +266,7 @@ public void
testReplicasToFixMisreplicationMaxReplicaPerRack() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 2);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 2, 4, 6, 8)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -277,7 +282,7 @@ public void
testReplicasToFixMisreplicationMaxReplicaPerRack() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 2);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 2, 4, 6, 8)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -296,7 +301,7 @@ public void
testReplicasToFixMisreplicationMaxReplicaPerRack() {
public void testReplicasWithoutMisreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 4)
.map(list::get).collect(Collectors.toList());
Map<ContainerReplica, Boolean> replicas =
@@ -313,7 +318,7 @@ public void testReplicasWithoutMisreplication() {
public void testReplicasToRemoveWithOneOverreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
ContainerID.valueOf(1), CLOSED, 0, 0, 0, list.subList(1,
6)));
@@ -334,7 +339,7 @@ public void testReplicasToRemoveWithOneOverreplication() {
public void testReplicasToRemoveWithTwoOverreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
@@ -355,7 +360,7 @@ public void testReplicasToRemoveWithTwoOverreplication() {
public void testReplicasToRemoveWith2CountPerUniqueReplica() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 3);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
@@ -381,7 +386,7 @@ public void
testReplicasToRemoveWith2CountPerUniqueReplica() {
public void testReplicasToRemoveWithoutReplicaIndex() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 3);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
Set<ContainerReplica> replicas = Sets.newHashSet(HddsTestUtils.getReplicas(
ContainerID.valueOf(1), CLOSED, 0, list.subList(0, 5)));
@@ -401,7 +406,7 @@ public void testReplicasToRemoveWithoutReplicaIndex() {
public void testReplicasToRemoveWithOverreplicationWithinSameRack() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 3);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
@@ -440,7 +445,7 @@ public void
testReplicasToRemoveWithOverreplicationWithinSameRack() {
public void testReplicasToRemoveWithNoOverreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = nodeManager.getAllNodes();
+ List<DatanodeDetails> list = getAllNodes(nodeManager);
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
ContainerID.valueOf(1), CLOSED, 0, 0, 0, list.subList(1,
6)));
@@ -617,8 +622,9 @@ private static class DummyPlacementPolicy extends
SCMCommonPlacementPolicy {
when(node.getNetworkFullPath()).thenReturn(String.valueOf(i));
return node;
}).collect(Collectors.toList());
- final List<? extends DatanodeDetails> datanodeDetails =
nodeManager.getAllNodes();
- rackMap = datanodeRackMap.entrySet().stream()
+ final List<DatanodeDetails> datanodeDetails = getAllNodes(nodeManager);
+ rackMap = datanodeRackMap
+ .entrySet().stream()
.collect(Collectors.toMap(
entry -> datanodeDetails.get(entry.getKey()),
entry -> racks.get(entry.getValue())));
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
index 3da501228cb..46bccad1119 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
@@ -33,6 +33,7 @@
import java.util.Map;
import java.util.Objects;
import java.util.Set;
+import java.util.TreeMap;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors;
@@ -105,7 +106,7 @@ public class MockNodeManager implements NodeManager {
private final List<DatanodeDetails> healthyNodes;
private final List<DatanodeDetails> staleNodes;
private final List<DatanodeDetails> deadNodes;
- private final Map<DatanodeDetails, SCMNodeStat> nodeMetricMap;
+ private final Map<DatanodeDetails, SCMNodeStat> nodeMetricMap = new
TreeMap<>();
private final SCMNodeStat aggregateStat;
private final Map<DatanodeID, List<SCMCommand<?>>> commandMap;
private Node2PipelineMap node2PipelineMap;
@@ -121,7 +122,6 @@ public class MockNodeManager implements NodeManager {
this.healthyNodes = new LinkedList<>();
this.staleNodes = new LinkedList<>();
this.deadNodes = new LinkedList<>();
- this.nodeMetricMap = new HashMap<>();
this.node2PipelineMap = new Node2PipelineMap();
this.node2ContainerMap = new NodeStateMap();
this.dnsToUuidMap = new ConcurrentHashMap<>();
@@ -250,7 +250,7 @@ private void populateNodeMetric(DatanodeDetails
datanodeDetails, int x) {
*/
@Override
public List<DatanodeDetails> getNodes(NodeStatus status) {
- return getNodes(status.getOperationalState(), status.getHealth());
+ return getDatanodeDetails(status.getOperationalState(),
status.getHealth());
}
/**
@@ -261,7 +261,16 @@ public List<DatanodeDetails> getNodes(NodeStatus status) {
* @return List of Datanodes that are Heartbeating SCM.
*/
@Override
- public List<DatanodeDetails> getNodes(
+ public List<DatanodeInfo> getNodes(
+ HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodestate)
{
+ final List<DatanodeDetails> details = getDatanodeDetails(opState,
nodestate);
+ if (details == null) {
+ return null;
+ }
+ return
details.stream().map(this::getDatanodeInfo).collect(Collectors.toList());
+ }
+
+ private List<DatanodeDetails> getDatanodeDetails(
HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodestate)
{
if (nodestate == HEALTHY) {
// mock storage reports for SCMCommonPlacementPolicy.hasEnoughSpace()
@@ -322,7 +331,7 @@ public int getNodeCount(NodeStatus status) {
@Override
public int getNodeCount(
HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodestate)
{
- List<DatanodeDetails> nodes = getNodes(opState, nodestate);
+ List<DatanodeDetails> nodes = getDatanodeDetails(opState, nodestate);
if (nodes != null) {
return nodes.size();
}
@@ -335,9 +344,9 @@ public int getNodeCount(
* @return List of DatanodeDetails known to SCM.
*/
@Override
- public List<DatanodeDetails> getAllNodes() {
+ public List<DatanodeInfo> getAllNodes() {
// mock storage reports for TestDiskBalancer
- List<DatanodeDetails> healthyNodesWithInfo = new ArrayList<>();
+ List<DatanodeInfo> healthyNodesWithInfo = new ArrayList<>();
for (Map.Entry<DatanodeDetails, SCMNodeStat> entry:
nodeMetricMap.entrySet()) {
NodeStatus nodeStatus = NodeStatus.inServiceHealthy();
@@ -399,7 +408,7 @@ public Map<DatanodeDetails, SCMNodeStat> getNodeStats() {
public List<DatanodeUsageInfo> getMostOrLeastUsedDatanodes(
boolean mostUsed) {
List<DatanodeDetails> datanodeDetailsList =
- getNodes(NodeOperationalState.IN_SERVICE, HEALTHY);
+ getDatanodeDetails(NodeOperationalState.IN_SERVICE, HEALTHY);
if (datanodeDetailsList == null) {
return new ArrayList<>();
}
@@ -428,9 +437,11 @@ public DatanodeUsageInfo getUsageInfo(DatanodeDetails
datanodeDetails) {
return new DatanodeUsageInfo(datanodeDetails, stat);
}
- @Override
@Nullable
public DatanodeInfo getDatanodeInfo(DatanodeDetails dd) {
+ if (dd instanceof DatanodeInfo) {
+ return (DatanodeInfo) dd;
+ }
if (nodeMetricMap.get(dd) == null) {
return null;
}
@@ -916,9 +927,9 @@ public List<SCMCommand<?>> getCommandQueue(DatanodeID dnID)
{
}
@Override
- public DatanodeDetails getNode(DatanodeID id) {
+ public DatanodeInfo getNode(DatanodeID id) {
Node node = clusterMap.getNode(NetConstants.DEFAULT_RACK + "/" + id);
- return node == null ? null : (DatanodeDetails)node;
+ return node == null ? null : getDatanodeInfo((DatanodeDetails)node);
}
@Override
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
index 5b10f9d2d75..3ad8d0c1ad2 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
@@ -17,7 +17,6 @@
package org.apache.hadoop.hdds.scm.container;
-import jakarta.annotation.Nullable;
import java.io.IOException;
import java.util.Collections;
import java.util.HashSet;
@@ -195,7 +194,7 @@ public List<DatanodeDetails> getNodes(NodeStatus
nodeStatus) {
}
@Override
- public List<DatanodeDetails> getNodes(
+ public List<DatanodeInfo> getNodes(
NodeOperationalState opState, HddsProtos.NodeState health) {
return null;
}
@@ -212,8 +211,8 @@ public int getNodeCount(NodeOperationalState opState,
}
@Override
- public List<DatanodeDetails> getAllNodes() {
- return null;
+ public List<DatanodeInfo> getAllNodes() {
+ return Collections.emptyList();
}
@Override
@@ -245,12 +244,6 @@ public DatanodeUsageInfo getUsageInfo(DatanodeDetails
datanodeDetails) {
return null;
}
- @Override
- @Nullable
- public DatanodeInfo getDatanodeInfo(DatanodeDetails dn) {
- return null;
- }
-
@Override
public boolean checkSpaceAndRecordAllocation(DatanodeInfo datanodeInfo,
ContainerID containerID) {
return true;
@@ -371,7 +364,7 @@ public List<SCMCommand<?>> getCommandQueue(DatanodeID dnID)
{
}
@Override
- public DatanodeDetails getNode(DatanodeID id) {
+ public DatanodeInfo getNode(DatanodeID id) {
return null;
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java
index b07f0da7851..24d005e9815 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java
@@ -17,7 +17,6 @@
package org.apache.hadoop.hdds.scm.container;
-import static
org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
import static org.apache.hadoop.hdds.scm.HddsTestUtils.getContainer;
import static org.apache.hadoop.hdds.scm.HddsTestUtils.getContainerReports;
import static org.apache.hadoop.hdds.scm.HddsTestUtils.getECContainer;
@@ -101,7 +100,7 @@ public class TestContainerReportHandler {
@BeforeEach
void setup() throws IOException {
final OzoneConfiguration conf = SCMTestUtils.getConf(testDir);
- nodeManager = new MockNodeManager(true, 10);
+ nodeManager = new MockNodeManager(true, 20);
containerManager = mock(ContainerManager.class);
dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get());
SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true);
@@ -634,12 +633,11 @@ private List<DatanodeDetails> setupECContainerForTesting(
container.getReplicationType());
final int numDatanodes =
container.getReplicationConfig().getRequiredNodes();
- // Register required number of datanodes with NodeManager
- List<DatanodeDetails> dns = new ArrayList<>(numDatanodes);
- for (int i = 0; i < numDatanodes; i++) {
- dns.add(randomDatanodeDetails());
- nodeManager.register(dns.get(i), null, null);
- }
+ // Get the required number of pre-registered, healthy datanodes from
NodeManager
+ List<DatanodeDetails> dns =
nodeManager.getNodes(NodeStatus.inServiceHealthy())
+ .stream()
+ .limit(numDatanodes)
+ .collect(Collectors.toList());
// Add this container to ContainerStateManager
containerStateManager.addContainer(container.getProtobuf());
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java
index 8d0938586bd..dd5344f4f11 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java
@@ -17,7 +17,6 @@
package org.apache.hadoop.hdds.scm.container;
-import static
org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
import static org.apache.hadoop.hdds.scm.HddsTestUtils.getContainer;
import static org.apache.hadoop.hdds.scm.HddsTestUtils.getECContainer;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
@@ -312,8 +311,9 @@ public void
testDeletedContainerWithLowerBcsidStaleReplicaRatis()
names = {"DELETING", "DELETED"})
public void
testECContainerWithStaleClosedReplicaShouldForceDelete(HddsProtos.LifeCycleState
state)
throws IOException {
- final DatanodeDetails datanode = randomDatanodeDetails();
- nodeManager.register(datanode, null, null);
+ //Get the first node from our list
+ final DatanodeDetails datanode = nodeManager.getNodes(
+ NodeStatus.inServiceHealthy()).get(0);
// Create an EC container
ECReplicationConfig repConfig = new ECReplicationConfig(3, 2);
final ContainerInfo ecContainer = getECContainer(
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
index 49fab0afe22..e4ca940cf56 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
@@ -39,7 +39,6 @@
import org.apache.hadoop.hdds.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
-import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.DatanodeID;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto.State;
@@ -54,6 +53,7 @@
import
org.apache.hadoop.hdds.scm.container.reconciliation.ReconciliationEligibilityHandler.EligibilityResult;
import
org.apache.hadoop.hdds.scm.container.reconciliation.ReconciliationEligibilityHandler.Result;
import org.apache.hadoop.hdds.scm.ha.SCMContext;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.server.events.EventPublisher;
import org.apache.hadoop.ozone.protocol.commands.CommandForDatanode;
import org.apache.hadoop.ozone.protocol.commands.SCMCommand;
@@ -300,7 +300,7 @@ private Set<ContainerReplica>
addReplicasToContainer(State... replicaStates) thr
// If no states are specified, replica list will be empty.
Set<ContainerReplica> replicas = new HashSet<>();
try (MockNodeManager nodeManager = new MockNodeManager(true,
replicaStates.length)) {
- List<DatanodeDetails> nodes = nodeManager.getAllNodes();
+ List<DatanodeInfo> nodes = nodeManager.getAllNodes();
for (int i = 0; i < replicaStates.length; i++) {
replicas.addAll(HddsTestUtils.getReplicas(CONTAINER_ID,
replicaStates[i], nodes.get(i)));
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerUtil.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerUtil.java
index ffca82e231b..fad2cf053f2 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerUtil.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerUtil.java
@@ -49,6 +49,7 @@
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
import org.apache.hadoop.hdds.scm.container.placement.metrics.SCMNodeMetric;
import
org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.node.NodeStatus;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
@@ -298,13 +299,16 @@ public void
testDatanodesWithInSufficientDiskSpaceAreExcluded() throws NodeNotFo
ConcurrentHashMap<DatanodeID, SizeAndTime> sizeScheduledMap = new
ConcurrentHashMap<>();
// fullDn has 10GB size scheduled, 30GB available and 20GB min free space,
so it should be excluded
- DatanodeDetails fullDn = MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeDetails fullDnDetails =
MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeInfo fullDn = new DatanodeInfo(fullDnDetails,
NodeStatus.inServiceHealthy(), null, 1);
sizeScheduledMap.put(fullDn.getID(), new SizeAndTime(10 * oneGb,
clock.millis()));
// spaceAvailableDn should not be excluded as it has sufficient space
- DatanodeDetails spaceAvailableDn =
MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeDetails spaceAvailableDnDetails =
MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeInfo spaceAvailableDn = new DatanodeInfo(spaceAvailableDnDetails,
NodeStatus.inServiceHealthy(), null, 1);
sizeScheduledMap.put(spaceAvailableDn.getID(), new SizeAndTime(10 * oneGb,
clock.millis()));
// expiredOpDn is the same as fullDn, however its op has expired - so it
should not be excluded
- DatanodeDetails expiredOpDn = MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeDetails expiredOpDnDetails =
MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeInfo expiredOpDn = new DatanodeInfo(expiredOpDnDetails,
NodeStatus.inServiceHealthy(), null, 1);
sizeScheduledMap.put(expiredOpDn.getID(), new SizeAndTime(10 * oneGb,
clock.millis() - rmConf.getEventTimeout() - 1));
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java
index 229d9283c74..ef6f187b040 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java
@@ -336,7 +336,7 @@ public boolean isPipelineCreationFrozen() {
@Override
public boolean checkSpaceAndRecordAllocation(Pipeline pipeline, ContainerID
containerID) {
for (DatanodeDetails dn : pipeline.getNodes()) {
- if
(!nodeManager.checkSpaceAndRecordAllocation(nodeManager.getDatanodeInfo(dn),
containerID)) {
+ if
(!nodeManager.checkSpaceAndRecordAllocation(nodeManager.getNode(dn.getID()),
containerID)) {
return false;
}
}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/node/TestDecommissionAndMaintenance.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/node/TestDecommissionAndMaintenance.java
index 341bbedf42d..2871af693d7 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/node/TestDecommissionAndMaintenance.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/node/TestDecommissionAndMaintenance.java
@@ -213,7 +213,7 @@ public void
testNodeWithOpenPipelineCanBeDecommissionedAndRecommissioned()
waitForDnToReachOpState(nm, toDecommission, DECOMMISSIONED);
// Ensure one node transitioned to DECOMMISSIONING
- List<DatanodeDetails> decomNodes = nm.getNodes(
+ List<DatanodeInfo> decomNodes = nm.getNodes(
DECOMMISSIONED,
HEALTHY);
assertEquals(1, decomNodes.size());
@@ -325,7 +325,7 @@ public void testInsufficientNodesCannotBeDecommissioned()
toDecommission.get(3).getIpAddress(),
toDecommission.get(4).getIpAddress()), false);
// Ensure no nodes transitioned to DECOMMISSIONING or DECOMMISSIONED
- List<DatanodeDetails> decomNodes = nm.getNodes(
+ List<DatanodeInfo> decomNodes = nm.getNodes(
DECOMMISSIONING,
HEALTHY);
assertEquals(0, decomNodes.size());
@@ -717,7 +717,7 @@ public void testInsufficientNodesCannotBePutInMaintenance()
getDNHostAndPort(toMaintenance.get(5))), 0, false);
// Ensure no nodes transitioned to MAINTENANCE
- List<DatanodeDetails> maintenanceNodes = nm.getNodes(
+ List<DatanodeInfo> maintenanceNodes = nm.getNodes(
ENTERING_MAINTENANCE,
HEALTHY);
assertEquals(0, maintenanceNodes.size());
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeMinFreeSpaceIntegration.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeMinFreeSpaceIntegration.java
index ca0e231a207..e1464f3e879 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeMinFreeSpaceIntegration.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeMinFreeSpaceIntegration.java
@@ -71,7 +71,7 @@ public void
storageReportsAtScmMatchSoftMinFreeSpaceFromConfig() throws Exceptio
() -> storageReportsMatchSoftMinFree(nm, dn, dnConf);
GenericTestUtils.waitFor(softSpareVisibleAtScm, 500, 120_000);
- DatanodeInfo info = nm.getDatanodeInfo(dn);
+ DatanodeInfo info = nm.getNode(dn.getID());
assertNotNull(info);
assertFalse(info.getStorageReports().isEmpty());
@@ -101,7 +101,7 @@ public void
storageReportsAtScmMatchSoftMinFreeSpaceFromConfig() throws Exceptio
*/
private static boolean storageReportsMatchSoftMinFree(
NodeManager nm, DatanodeDetails dn, DatanodeConfiguration dnConf) {
- DatanodeInfo info = nm.getDatanodeInfo(dn);
+ DatanodeInfo info = nm.getNode(dn.getID());
if (info == null || info.getStorageReports().isEmpty()) {
return false;
}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancer.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancer.java
index cfdd4c7d37d..232554209c5 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancer.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancer.java
@@ -219,7 +219,7 @@ public void testDatanodeDiskBalancerStatus() throws
IOException, InterruptedExce
}
// Query status from remaining IN_SERVICE DNs and verify they still show
RUNNING
- List<DatanodeDetails> inServiceDatanodes = nm.getNodes(IN_SERVICE,
HddsProtos.NodeState.HEALTHY);
+ final List<DatanodeInfo> inServiceDatanodes = nm.getNodes(IN_SERVICE,
HddsProtos.NodeState.HEALTHY);
statusProtoList.clear();
for (DatanodeDetails dn : inServiceDatanodes) {
try (DiskBalancerProtocol proxy = getDiskBalancerProxy(dn)) {
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancerDuringDecommissionAndMaintenance.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancerDuringDecommissionAndMaintenance.java
index 6e0686e1359..b8bee566a54 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancerDuringDecommissionAndMaintenance.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/scm/node/TestDiskBalancerDuringDecommissionAndMaintenance.java
@@ -136,13 +136,6 @@ private DiskBalancerProtocol
getDiskBalancerProxy(DatanodeDetails dn) throws IOE
return new DiskBalancerProtocolClientSideTranslatorPB(nodeAddr, user,
conf);
}
- /**
- * Helper method to get all IN_SERVICE datanodes.
- */
- private List<DatanodeDetails> getInServiceDatanodes(NodeManager nm) {
- return nm.getNodes(IN_SERVICE, HddsProtos.NodeState.HEALTHY);
- }
-
/**
* Helper method to query DiskBalancer info from all IN_SERVICE datanodes.
* Similar to --in-service-datanodes option in CLI.
@@ -150,7 +143,7 @@ private List<DatanodeDetails>
getInServiceDatanodes(NodeManager nm) {
private <T> List<T> queryAllInServiceDatanodes(
DiskBalancerQuery<T> query) throws IOException {
NodeManager nm = cluster.getStorageContainerManager().getScmNodeManager();
- List<DatanodeDetails> inServiceDatanodes = getInServiceDatanodes(nm);
+ final List<DatanodeInfo> inServiceDatanodes = nm.getNodes(IN_SERVICE,
HddsProtos.NodeState.HEALTHY);
List<T> results = new ArrayList<>();
for (DatanodeDetails dn : inServiceDatanodes) {
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/PipelineSyncTask.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/PipelineSyncTask.java
index f8985a0a867..ec40c3a6a1a 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/PipelineSyncTask.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/PipelineSyncTask.java
@@ -21,10 +21,12 @@
import java.io.IOException;
import java.util.List;
+import java.util.Set;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.stream.Collectors;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.DatanodeID;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.Node;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
@@ -109,13 +111,14 @@ protected void runTask() throws IOException,
NodeNotFoundException {
*/
private void syncOperationalStateOnDeadNodes()
throws IOException, NodeNotFoundException {
- List<DatanodeDetails> deadNodesOnRecon = nodeManager.getNodes(null, DEAD);
+ final Set<DatanodeID> deadNodesOnRecon = nodeManager.getNodes(null,
DEAD).stream()
+ .map(info -> info.getID())
+ .collect(Collectors.toSet());
if (!deadNodesOnRecon.isEmpty()) {
List<Node> scmNodes = scmClient.getNodes();
List<Node> filteredScmNodes = scmNodes.stream()
- .filter(n -> deadNodesOnRecon.contains(
- DatanodeDetails.getFromProtoBuf(n.getNodeID())))
+ .filter(n ->
deadNodesOnRecon.contains(DatanodeDetails.getFromProtoBuf(n.getNodeID()).getID()))
.collect(Collectors.toList());
for (Node deadNode : filteredScmNodes) {
diff --git
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconIncrementalContainerReportHandler.java
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconIncrementalContainerReportHandler.java
index c82eabb92e7..7d6a217c0f8 100644
---
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconIncrementalContainerReportHandler.java
+++
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconIncrementalContainerReportHandler.java
@@ -48,7 +48,9 @@
import org.apache.hadoop.hdds.scm.ha.SCMContext;
import org.apache.hadoop.hdds.scm.net.NetworkTopology;
import org.apache.hadoop.hdds.scm.net.NetworkTopologyImpl;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.NodeManager;
+import org.apache.hadoop.hdds.scm.node.NodeStatus;
import org.apache.hadoop.hdds.scm.node.SCMNodeManager;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
import
org.apache.hadoop.hdds.scm.server.SCMDatanodeHeartbeatDispatcher.IncrementalContainerReportFromDatanode;
@@ -136,11 +138,12 @@ public void testProcessICRStateMismatch()
ReconContainerManager containerManager = getContainerManager();
containerManager.addNewContainer(containerWithPipeline);
- DatanodeDetails datanodeDetails =
- containerWithPipeline.getPipeline().getFirstNode();
+ DatanodeInfo datanodeInfo = new DatanodeInfo(
+ containerWithPipeline.getPipeline().getFirstNode(),
+ NodeStatus.inServiceHealthy(), null, 1000);
NodeManager nodeManagerMock = mock(NodeManager.class);
when(nodeManagerMock.getNode(any(DatanodeID.class)))
- .thenReturn(datanodeDetails);
+ .thenReturn(datanodeInfo);
IncrementalContainerReportFromDatanode reportMock =
mock(IncrementalContainerReportFromDatanode.class);
when(reportMock.getDatanodeDetails())
@@ -148,7 +151,7 @@ public void testProcessICRStateMismatch()
IncrementalContainerReportProto containerReport =
getIncrementalContainerReportProto(containerID, state,
- datanodeDetails.getUuidString());
+ datanodeInfo.getUuidString());
when(reportMock.getReport()).thenReturn(containerReport);
ReconIncrementalContainerReportHandler reconIcr =
new ReconIncrementalContainerReportHandler(nodeManagerMock,
@@ -239,8 +242,10 @@ private LifeCycleState getContainerStateFromReplicaState(
private static NodeManager getNodeManagerMock(DatanodeDetails
datanodeDetails)
throws NodeNotFoundException {
NodeManager nodeManagerMock = mock(NodeManager.class);
+ DatanodeInfo datanodeInfo = new DatanodeInfo(
+ datanodeDetails, NodeStatus.inServiceHealthy(), null, 1000);
when(nodeManagerMock.getNode(any(DatanodeID.class)))
- .thenReturn(datanodeDetails);
+ .thenReturn(datanodeInfo);
return nodeManagerMock;
}
diff --git
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconNodeManager.java
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconNodeManager.java
index 9c93441457d..508d391b32e 100644
---
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconNodeManager.java
+++
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/scm/TestReconNodeManager.java
@@ -44,6 +44,7 @@
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
import org.apache.hadoop.hdds.scm.net.NetworkTopology;
import org.apache.hadoop.hdds.scm.net.NetworkTopologyImpl;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
import org.apache.hadoop.hdds.server.events.EventQueue;
import org.apache.hadoop.hdds.upgrade.HDDSLayoutVersionManager;
@@ -229,7 +230,7 @@ public void testUpdateNodeOperationalStateFromScm() throws
Exception {
reconNodeManager.updateNodeOperationalStateFromScm(node, datanodeDetails);
assertEquals(DECOMMISSIONING, reconNodeManager
.getNode(datanodeDetails.getID()).getPersistedOpState());
- List<DatanodeDetails> nodes =
+ List<DatanodeInfo> nodes =
reconNodeManager.getNodes(DECOMMISSIONING, null);
assertEquals(1, nodes.size());
assertEquals(datanodeDetails.getID(), nodes.get(0).getID());
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]