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 ca009c55287 Revert "HDDS-15146. Remove
getDatanodeInfo(DatanodeDetails) from NodeManager (#10601)"
ca009c55287 is described below
commit ca009c552873747ad62211bbdc5f0d0af80cc42c
Author: Doroszlai, Attila <[email protected]>
AuthorDate: Fri Jun 26 20:58:16 2026 +0200
Revert "HDDS-15146. Remove getDatanodeInfo(DatanodeDetails) from
NodeManager (#10601)"
This reverts commit c9736f308ef14bfb50a7beaaa1ec058f4e94c602.
Reason for revert: compile error
---
.../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 | 28 +++++++-
.../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, 206 insertions(+), 127 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 2326cd894e6..fe3438b1ea7 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,12 +136,11 @@ public void onMessage(final ContainerReportFromDatanode
reportFromDatanode,
final DatanodeDetails dnFromReport =
reportFromDatanode.getDatanodeDetails();
- final DatanodeInfo datanodeInfo =
getNodeManager().getNode(dnFromReport.getID());
- if (datanodeInfo == null) {
+ final DatanodeDetails datanodeDetails =
getNodeManager().getNode(dnFromReport.getID());
+ if (datanodeDetails == null) {
getLogger().warn("Datanode not found: {}", dnFromReport);
return;
}
- final DatanodeDetails datanodeDetails = datanodeInfo;
final ContainerReportsProto containerReport =
reportFromDatanode.getReport();
try {
@@ -154,6 +153,7 @@ 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,7 +178,9 @@ 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
- getNodeManager().removePendingAllocationForDatanode(datanodeInfo,
cid);
+ if (datanodeInfo != null) {
+
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 49af3cae740..51eb0e4e41f 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<DatanodeInfo> getNodes(
+ List<DatanodeDetails> getNodes(
NodeOperationalState opState, NodeState health);
/**
@@ -134,8 +134,10 @@ List<DatanodeInfo> getNodes(
int getNodeCount(
NodeOperationalState opState, NodeState health);
- /** @return a shadow copied list of all datanodes, sorted by {@link
DatanodeID}. */
- List<DatanodeInfo> getAllNodes();
+ /**
+ * @return all datanodes known to SCM.
+ */
+ List<? extends DatanodeDetails> getAllNodes();
/** @return the number of datanodes. */
default int getAllNodeCount() {
@@ -173,6 +175,15 @@ 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
@@ -393,7 +404,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 DatanodeInfo getNode(@Nullable DatanodeID id);
+ @Nullable DatanodeDetails 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 9539379fd84..57ee8ef9adb 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,8 +537,12 @@ public int getVolumeFailuresNodeCount() {
return getVolumeFailuresNodes().size();
}
- /** @return a shadow copied list of all datanodes, sorted by {@link
DatanodeID}. */
- List<DatanodeInfo> getAllNodes() {
+ /**
+ * Returns all the nodes which have registered to NodeStateManager.
+ *
+ * @return all the managed nodes
+ */
+ public 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 f628481b583..394793bb9a3 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,6 +26,7 @@
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;
@@ -269,9 +270,11 @@ public List<DatanodeDetails> getNodes(NodeStatus
nodeStatus) {
* @return List of Datanodes that are known to SCM in the requested states.
*/
@Override
- public List<DatanodeInfo> getNodes(
+ public List<DatanodeDetails> getNodes(
NodeOperationalState opState, NodeState health) {
- return nodeStateManager.getNodes(opState, health);
+ return nodeStateManager.getNodes(opState, health)
+ .stream()
+ .map(node -> (DatanodeDetails)node).collect(Collectors.toList());
}
@Override
@@ -1012,7 +1015,8 @@ public Map<DatanodeDetails, SCMNodeStat> getNodeStats() {
@Override
public List<DatanodeUsageInfo> getMostOrLeastUsedDatanodes(
boolean mostUsed) {
- final List<DatanodeInfo> healthyNodes = getNodes(IN_SERVICE,
NodeState.HEALTHY);
+ List<DatanodeDetails> healthyNodes =
+ getNodes(IN_SERVICE, NodeState.HEALTHY);
List<DatanodeUsageInfo> datanodeUsageInfoList =
new ArrayList<>(healthyNodes.size());
@@ -1059,6 +1063,24 @@ 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);
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 67c487e3ba1..8aea57b23ab 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 TreeMap<>();
+ private final Map<DatanodeID, DatanodeEntry> nodeMap = new HashMap<>();
private final ReadWriteLock lock = new ReentrantReadWriteLock();
@@ -166,7 +166,11 @@ public int getNodeCount() {
}
}
- /** @return a shadow copied list of all datanodes, sorted by {@link
DatanodeID}. */
+ /**
+ * Returns the list of all the nodes as DatanodeInfo objects.
+ *
+ * @return list of all the node ids
+ */
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 8003d8d52e9..fc8229c677d 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,8 +638,7 @@ 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) {
- // Refactored to use getNode instead of getDatanodeInfo
- final DatanodeInfo info = nodeManager.getNode(dn.getID());
+ final DatanodeInfo info = nodeManager.getDatanodeInfo(dn);
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 6883cee0127..010e0ca8b6c 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,6 +46,7 @@
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;
@@ -690,9 +691,9 @@ public List<HddsProtos.Node> queryNode(
}
try {
List<HddsProtos.Node> result = new ArrayList<>();
- for (DatanodeDetails node : scm.getScmNodeManager().getNodes(opState,
state)) {
+ for (DatanodeDetails node : queryNode(opState, state)) {
NodeStatus ns = scm.getScmNodeManager().getNodeStatus(node);
- DatanodeInfo datanodeInfo = node instanceof DatanodeInfo ?
(DatanodeInfo) node : null;
+ DatanodeInfo datanodeInfo =
scm.getScmNodeManager().getDatanodeInfo(node);
HddsProtos.Node.Builder nodeBuilder = HddsProtos.Node.newBuilder()
.setNodeID(node.toProto(clientVersion))
.addNodeStates(ns.getHealth())
@@ -716,22 +717,34 @@ public List<HddsProtos.Node> queryNode(
}
@Override
- public HddsProtos.Node queryNode(UUID uuid) {
+ public HddsProtos.Node queryNode(UUID uuid)
+ throws IOException {
final Map<String, String> auditMap = Maps.newHashMap();
auditMap.put("uuid", String.valueOf(uuid));
HddsProtos.Node result = null;
- 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();
+ 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;
}
AUDIT.logReadSuccess(buildAuditMessageForSuccess(
SCMAction.QUERY_NODE, auditMap));
@@ -1580,6 +1593,26 @@ 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;
@@ -1592,6 +1625,24 @@ 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 9df7c4f69d6..9290e8d83be 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,7 +33,6 @@
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;
@@ -85,15 +84,11 @@ 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 = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> result = dummyPlacementPolicy.getResultSet(3, list);
Set<DatanodeDetails> resultSet = new HashSet<>(result);
assertNotEquals(1, resultSet.size());
@@ -141,7 +136,7 @@ public void
testReplicasToFixMisreplicationWithOneMisreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 5)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -162,7 +157,7 @@ public void
testReplicasToFixMisreplicationWithTwoMisreplication() {
3, ImmutableList.of(3, 8),
4, ImmutableList.of(4, 9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 5)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -183,7 +178,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 8),
4, ImmutableList.of(4, 9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 5)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -205,7 +200,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 4, 8),
4, ImmutableList.of(9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 4)
.map(list::get).collect(Collectors.toList());
//Creating Replicas without replica Index
@@ -228,7 +223,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 4, 8),
4, ImmutableList.of(9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
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
@@ -251,7 +246,7 @@ public void
testReplicasToFixMisreplicationWithThreeMisreplication() {
3, ImmutableList.of(3, 4, 8),
4, ImmutableList.of(9))), 5);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
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
@@ -266,7 +261,7 @@ public void
testReplicasToFixMisreplicationMaxReplicaPerRack() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 2);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 2, 4, 6, 8)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -282,7 +277,7 @@ public void
testReplicasToFixMisreplicationMaxReplicaPerRack() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 2);
List<Node> racks = dummyPlacementPolicy.racks;
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 2, 4, 6, 8)
.map(list::get).collect(Collectors.toList());
List<ContainerReplica> replicas =
@@ -301,7 +296,7 @@ public void
testReplicasToFixMisreplicationMaxReplicaPerRack() {
public void testReplicasWithoutMisreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
List<DatanodeDetails> replicaDns = Stream.of(0, 1, 2, 3, 4)
.map(list::get).collect(Collectors.toList());
Map<ContainerReplica, Boolean> replicas =
@@ -318,7 +313,7 @@ public void testReplicasWithoutMisreplication() {
public void testReplicasToRemoveWithOneOverreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
ContainerID.valueOf(1), CLOSED, 0, 0, 0, list.subList(1,
6)));
@@ -339,7 +334,7 @@ public void testReplicasToRemoveWithOneOverreplication() {
public void testReplicasToRemoveWithTwoOverreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
@@ -360,7 +355,7 @@ public void testReplicasToRemoveWithTwoOverreplication() {
public void testReplicasToRemoveWith2CountPerUniqueReplica() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 3);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
@@ -386,7 +381,7 @@ public void
testReplicasToRemoveWith2CountPerUniqueReplica() {
public void testReplicasToRemoveWithoutReplicaIndex() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 3);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
Set<ContainerReplica> replicas = Sets.newHashSet(HddsTestUtils.getReplicas(
ContainerID.valueOf(1), CLOSED, 0, list.subList(0, 5)));
@@ -406,7 +401,7 @@ public void testReplicasToRemoveWithoutReplicaIndex() {
public void testReplicasToRemoveWithOverreplicationWithinSameRack() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 3);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
@@ -445,7 +440,7 @@ public void
testReplicasToRemoveWithOverreplicationWithinSameRack() {
public void testReplicasToRemoveWithNoOverreplication() {
DummyPlacementPolicy dummyPlacementPolicy =
new DummyPlacementPolicy(nodeManager, conf, 5);
- List<DatanodeDetails> list = getAllNodes(nodeManager);
+ List<DatanodeDetails> list = nodeManager.getAllNodes();
Set<ContainerReplica> replicas = Sets.newHashSet(
HddsTestUtils.getReplicasWithReplicaIndex(
ContainerID.valueOf(1), CLOSED, 0, 0, 0, list.subList(1,
6)));
@@ -622,9 +617,8 @@ private static class DummyPlacementPolicy extends
SCMCommonPlacementPolicy {
when(node.getNetworkFullPath()).thenReturn(String.valueOf(i));
return node;
}).collect(Collectors.toList());
- final List<DatanodeDetails> datanodeDetails = getAllNodes(nodeManager);
- rackMap = datanodeRackMap
- .entrySet().stream()
+ final List<? extends DatanodeDetails> datanodeDetails =
nodeManager.getAllNodes();
+ 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 46bccad1119..3da501228cb 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,7 +33,6 @@
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;
@@ -106,7 +105,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 = new
TreeMap<>();
+ private final Map<DatanodeDetails, SCMNodeStat> nodeMetricMap;
private final SCMNodeStat aggregateStat;
private final Map<DatanodeID, List<SCMCommand<?>>> commandMap;
private Node2PipelineMap node2PipelineMap;
@@ -122,6 +121,7 @@ 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 getDatanodeDetails(status.getOperationalState(),
status.getHealth());
+ return getNodes(status.getOperationalState(), status.getHealth());
}
/**
@@ -261,16 +261,7 @@ public List<DatanodeDetails> getNodes(NodeStatus status) {
* @return List of Datanodes that are Heartbeating SCM.
*/
@Override
- 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(
+ public List<DatanodeDetails> getNodes(
HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodestate)
{
if (nodestate == HEALTHY) {
// mock storage reports for SCMCommonPlacementPolicy.hasEnoughSpace()
@@ -331,7 +322,7 @@ public int getNodeCount(NodeStatus status) {
@Override
public int getNodeCount(
HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodestate)
{
- List<DatanodeDetails> nodes = getDatanodeDetails(opState, nodestate);
+ List<DatanodeDetails> nodes = getNodes(opState, nodestate);
if (nodes != null) {
return nodes.size();
}
@@ -344,9 +335,9 @@ public int getNodeCount(
* @return List of DatanodeDetails known to SCM.
*/
@Override
- public List<DatanodeInfo> getAllNodes() {
+ public List<DatanodeDetails> getAllNodes() {
// mock storage reports for TestDiskBalancer
- List<DatanodeInfo> healthyNodesWithInfo = new ArrayList<>();
+ List<DatanodeDetails> healthyNodesWithInfo = new ArrayList<>();
for (Map.Entry<DatanodeDetails, SCMNodeStat> entry:
nodeMetricMap.entrySet()) {
NodeStatus nodeStatus = NodeStatus.inServiceHealthy();
@@ -408,7 +399,7 @@ public Map<DatanodeDetails, SCMNodeStat> getNodeStats() {
public List<DatanodeUsageInfo> getMostOrLeastUsedDatanodes(
boolean mostUsed) {
List<DatanodeDetails> datanodeDetailsList =
- getDatanodeDetails(NodeOperationalState.IN_SERVICE, HEALTHY);
+ getNodes(NodeOperationalState.IN_SERVICE, HEALTHY);
if (datanodeDetailsList == null) {
return new ArrayList<>();
}
@@ -437,11 +428,9 @@ 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;
}
@@ -927,9 +916,9 @@ public List<SCMCommand<?>> getCommandQueue(DatanodeID dnID)
{
}
@Override
- public DatanodeInfo getNode(DatanodeID id) {
+ public DatanodeDetails getNode(DatanodeID id) {
Node node = clusterMap.getNode(NetConstants.DEFAULT_RACK + "/" + id);
- return node == null ? null : getDatanodeInfo((DatanodeDetails)node);
+ return node == null ? null : (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 3ad8d0c1ad2..5b10f9d2d75 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,6 +17,7 @@
package org.apache.hadoop.hdds.scm.container;
+import jakarta.annotation.Nullable;
import java.io.IOException;
import java.util.Collections;
import java.util.HashSet;
@@ -194,7 +195,7 @@ public List<DatanodeDetails> getNodes(NodeStatus
nodeStatus) {
}
@Override
- public List<DatanodeInfo> getNodes(
+ public List<DatanodeDetails> getNodes(
NodeOperationalState opState, HddsProtos.NodeState health) {
return null;
}
@@ -211,8 +212,8 @@ public int getNodeCount(NodeOperationalState opState,
}
@Override
- public List<DatanodeInfo> getAllNodes() {
- return Collections.emptyList();
+ public List<DatanodeDetails> getAllNodes() {
+ return null;
}
@Override
@@ -244,6 +245,12 @@ 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;
@@ -364,7 +371,7 @@ public List<SCMCommand<?>> getCommandQueue(DatanodeID dnID)
{
}
@Override
- public DatanodeInfo getNode(DatanodeID id) {
+ public DatanodeDetails 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 24d005e9815..b07f0da7851 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,6 +17,7 @@
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;
@@ -100,7 +101,7 @@ public class TestContainerReportHandler {
@BeforeEach
void setup() throws IOException {
final OzoneConfiguration conf = SCMTestUtils.getConf(testDir);
- nodeManager = new MockNodeManager(true, 20);
+ nodeManager = new MockNodeManager(true, 10);
containerManager = mock(ContainerManager.class);
dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get());
SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true);
@@ -633,11 +634,12 @@ private List<DatanodeDetails> setupECContainerForTesting(
container.getReplicationType());
final int numDatanodes =
container.getReplicationConfig().getRequiredNodes();
- // Get the required number of pre-registered, healthy datanodes from
NodeManager
- List<DatanodeDetails> dns =
nodeManager.getNodes(NodeStatus.inServiceHealthy())
- .stream()
- .limit(numDatanodes)
- .collect(Collectors.toList());
+ // 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);
+ }
// 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 dd5344f4f11..8d0938586bd 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,6 +17,7 @@
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;
@@ -311,9 +312,8 @@ public void
testDeletedContainerWithLowerBcsidStaleReplicaRatis()
names = {"DELETING", "DELETED"})
public void
testECContainerWithStaleClosedReplicaShouldForceDelete(HddsProtos.LifeCycleState
state)
throws IOException {
- //Get the first node from our list
- final DatanodeDetails datanode = nodeManager.getNodes(
- NodeStatus.inServiceHealthy()).get(0);
+ final DatanodeDetails datanode = randomDatanodeDetails();
+ nodeManager.register(datanode, null, null);
// 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 e4ca940cf56..49fab0afe22 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,6 +39,7 @@
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;
@@ -53,7 +54,6 @@
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<DatanodeInfo> nodes = nodeManager.getAllNodes();
+ List<DatanodeDetails> 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 fad2cf053f2..ffca82e231b 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,7 +49,6 @@
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;
@@ -299,16 +298,13 @@ 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 fullDnDetails =
MockDatanodeDetails.randomDatanodeDetails();
- DatanodeInfo fullDn = new DatanodeInfo(fullDnDetails,
NodeStatus.inServiceHealthy(), null, 1);
+ DatanodeDetails fullDn = MockDatanodeDetails.randomDatanodeDetails();
sizeScheduledMap.put(fullDn.getID(), new SizeAndTime(10 * oneGb,
clock.millis()));
// spaceAvailableDn should not be excluded as it has sufficient space
- DatanodeDetails spaceAvailableDnDetails =
MockDatanodeDetails.randomDatanodeDetails();
- DatanodeInfo spaceAvailableDn = new DatanodeInfo(spaceAvailableDnDetails,
NodeStatus.inServiceHealthy(), null, 1);
+ DatanodeDetails spaceAvailableDn =
MockDatanodeDetails.randomDatanodeDetails();
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 expiredOpDnDetails =
MockDatanodeDetails.randomDatanodeDetails();
- DatanodeInfo expiredOpDn = new DatanodeInfo(expiredOpDnDetails,
NodeStatus.inServiceHealthy(), null, 1);
+ DatanodeDetails expiredOpDn = MockDatanodeDetails.randomDatanodeDetails();
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 ef6f187b040..229d9283c74 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.getNode(dn.getID()),
containerID)) {
+ if
(!nodeManager.checkSpaceAndRecordAllocation(nodeManager.getDatanodeInfo(dn),
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 2871af693d7..341bbedf42d 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<DatanodeInfo> decomNodes = nm.getNodes(
+ List<DatanodeDetails> 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<DatanodeInfo> decomNodes = nm.getNodes(
+ List<DatanodeDetails> 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<DatanodeInfo> maintenanceNodes = nm.getNodes(
+ List<DatanodeDetails> 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 e1464f3e879..ca0e231a207 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.getNode(dn.getID());
+ DatanodeInfo info = nm.getDatanodeInfo(dn);
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.getNode(dn.getID());
+ DatanodeInfo info = nm.getDatanodeInfo(dn);
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 232554209c5..cfdd4c7d37d 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
- final List<DatanodeInfo> inServiceDatanodes = nm.getNodes(IN_SERVICE,
HddsProtos.NodeState.HEALTHY);
+ List<DatanodeDetails> 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 b8bee566a54..6e0686e1359 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,6 +136,13 @@ 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.
@@ -143,7 +150,7 @@ private DiskBalancerProtocol
getDiskBalancerProxy(DatanodeDetails dn) throws IOE
private <T> List<T> queryAllInServiceDatanodes(
DiskBalancerQuery<T> query) throws IOException {
NodeManager nm = cluster.getStorageContainerManager().getScmNodeManager();
- final List<DatanodeInfo> inServiceDatanodes = nm.getNodes(IN_SERVICE,
HddsProtos.NodeState.HEALTHY);
+ List<DatanodeDetails> inServiceDatanodes = getInServiceDatanodes(nm);
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 ec40c3a6a1a..f8985a0a867 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,12 +21,10 @@
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;
@@ -111,14 +109,13 @@ protected void runTask() throws IOException,
NodeNotFoundException {
*/
private void syncOperationalStateOnDeadNodes()
throws IOException, NodeNotFoundException {
- final Set<DatanodeID> deadNodesOnRecon = nodeManager.getNodes(null,
DEAD).stream()
- .map(info -> info.getID())
- .collect(Collectors.toSet());
+ List<DatanodeDetails> deadNodesOnRecon = nodeManager.getNodes(null, DEAD);
if (!deadNodesOnRecon.isEmpty()) {
List<Node> scmNodes = scmClient.getNodes();
List<Node> filteredScmNodes = scmNodes.stream()
- .filter(n ->
deadNodesOnRecon.contains(DatanodeDetails.getFromProtoBuf(n.getNodeID()).getID()))
+ .filter(n -> deadNodesOnRecon.contains(
+ DatanodeDetails.getFromProtoBuf(n.getNodeID())))
.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 7d6a217c0f8..c82eabb92e7 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,9 +48,7 @@
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;
@@ -138,12 +136,11 @@ public void testProcessICRStateMismatch()
ReconContainerManager containerManager = getContainerManager();
containerManager.addNewContainer(containerWithPipeline);
- DatanodeInfo datanodeInfo = new DatanodeInfo(
- containerWithPipeline.getPipeline().getFirstNode(),
- NodeStatus.inServiceHealthy(), null, 1000);
+ DatanodeDetails datanodeDetails =
+ containerWithPipeline.getPipeline().getFirstNode();
NodeManager nodeManagerMock = mock(NodeManager.class);
when(nodeManagerMock.getNode(any(DatanodeID.class)))
- .thenReturn(datanodeInfo);
+ .thenReturn(datanodeDetails);
IncrementalContainerReportFromDatanode reportMock =
mock(IncrementalContainerReportFromDatanode.class);
when(reportMock.getDatanodeDetails())
@@ -151,7 +148,7 @@ public void testProcessICRStateMismatch()
IncrementalContainerReportProto containerReport =
getIncrementalContainerReportProto(containerID, state,
- datanodeInfo.getUuidString());
+ datanodeDetails.getUuidString());
when(reportMock.getReport()).thenReturn(containerReport);
ReconIncrementalContainerReportHandler reconIcr =
new ReconIncrementalContainerReportHandler(nodeManagerMock,
@@ -242,10 +239,8 @@ 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(datanodeInfo);
+ .thenReturn(datanodeDetails);
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 508d391b32e..9c93441457d 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,7 +44,6 @@
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;
@@ -230,7 +229,7 @@ public void testUpdateNodeOperationalStateFromScm() throws
Exception {
reconNodeManager.updateNodeOperationalStateFromScm(node, datanodeDetails);
assertEquals(DECOMMISSIONING, reconNodeManager
.getNode(datanodeDetails.getID()).getPersistedOpState());
- List<DatanodeInfo> nodes =
+ List<DatanodeDetails> 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]