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 b1dc0abff667..ac6e92c1b44f 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 @@ -137,11 +137,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 { @@ -154,7 +155,6 @@ public void onMessage(final ContainerReportFromDatanode reportFromDatanode, containerReport.getReportsList(); final Set expectedContainersInDatanode = getNodeManager().getContainers(datanodeDetails); - DatanodeInfo datanodeInfo = datanodeDetails instanceof DatanodeInfo ? (DatanodeInfo) datanodeDetails : null; for (ContainerReplicaProto replica : replicas) { ContainerID cid = ContainerID.valueOf(replica.getContainerID()); @@ -179,9 +179,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 6a2d92b64a5e..13f5c959d1d0 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 getNodes( + List getNodes( NodeOperationalState opState, NodeState health); /** @@ -134,10 +134,8 @@ List getNodes( int getNodeCount( NodeOperationalState opState, NodeState health); - /** - * @return all datanodes known to SCM. - */ - List getAllNodes(); + /** @return a shadow copied list of all datanodes, sorted by {@link DatanodeID}. */ + List 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 @@ -387,7 +376,7 @@ Map getTotalDatanodeCommandCounts( List> 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 57ee8ef9adb6..9539379fd844 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 getAllNodes() { + /** @return a shadow copied list of all datanodes, sorted by {@link DatanodeID}. */ + List 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 926884c5b615..21fda8522f00 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 @@ -27,7 +27,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; @@ -268,11 +267,9 @@ public List getNodes(NodeStatus nodeStatus) { * @return List of Datanodes that are known to SCM in the requested states. */ @Override - public List getNodes( + public List getNodes( NodeOperationalState opState, NodeState health) { - return nodeStateManager.getNodes(opState, health) - .stream() - .map(node -> (DatanodeDetails)node).collect(Collectors.toList()); + return nodeStateManager.getNodes(opState, health); } @Override @@ -1013,8 +1010,7 @@ public Map getNodeStats() { @Override public List getMostOrLeastUsedDatanodes( boolean mostUsed) { - List healthyNodes = - getNodes(IN_SERVICE, NodeState.HEALTHY); + final List healthyNodes = getNodes(IN_SERVICE, NodeState.HEALTHY); List datanodeUsageInfoList = new ArrayList<>(healthyNodes.size()); @@ -1061,24 +1057,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); 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 8aea57b23ab0..67c487e3ba1c 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 nodeMap = new HashMap<>(); + private final Map 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 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 bdb286292d10..88feb47dbbdc 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 @@ -640,7 +640,8 @@ public boolean checkSpaceAndRecordAllocation(Pipeline pipeline, ContainerID cont final Set datanodeDetails = pipeline.getNodeSet(); final List 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 010e0ca8b6c4..6883cee0127d 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 queryNode( } try { List 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 queryNode( } @Override - public HddsProtos.Node queryNode(UUID uuid) - throws IOException { + public HddsProtos.Node queryNode(UUID uuid) { final Map 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 getListOfContainerIDs( } } - /** - * Queries a list of Node that match a set of statuses. - * - *

For example, if the nodeStatuses is HEALTHY and RAFT_MEMBER, then - * this call will return all - * healthy nodes which members in Raft pipeline. - * - *

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 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 queryNodeState( - HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodeState) { - Set returnSet = new TreeSet<>(); - List tmp = scm.getScmNodeManager() - .getNodes(opState, nodeState); - if ((tmp != null) && (!tmp.isEmpty())) { - returnSet.addAll(tmp); - } - return returnSet; - } - @Override public AuditMessage buildAuditMessageForSuccess( AuditAction op, Map 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 5ecd3fa1c74c..5296b1e8e240 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 @@ -34,6 +34,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; @@ -85,11 +86,15 @@ void setup(@TempDir File testDir) { conf = SCMTestUtils.getConf(testDir); } + static List getAllNodes(NodeManager nm) { + return new ArrayList<>(nm.getAllNodes()); + } + @Test public void testGetResultSet() throws SCMException { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 5); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List result = dummyPlacementPolicy.getResultSet(3, list); Set resultSet = new HashSet<>(result); assertNotEquals(1, resultSet.size()); @@ -137,7 +142,7 @@ public void testReplicasToFixMisreplicationWithOneMisreplication() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 5); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 1, 2, 3, 5) .map(list::get).collect(Collectors.toList()); List replicas = @@ -158,7 +163,7 @@ public void testReplicasToFixMisreplicationWithTwoMisreplication() { 3, ImmutableList.of(3, 8), 4, ImmutableList.of(4, 9))), 5); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 1, 2, 3, 5) .map(list::get).collect(Collectors.toList()); List replicas = @@ -179,7 +184,7 @@ public void testReplicasToFixMisreplicationWithThreeMisreplication() { 3, ImmutableList.of(3, 8), 4, ImmutableList.of(4, 9))), 5); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 1, 2, 3, 5) .map(list::get).collect(Collectors.toList()); List replicas = @@ -201,7 +206,7 @@ public void testReplicasToFixMisreplicationWithThreeMisreplication() { 3, ImmutableList.of(3, 4, 8), 4, ImmutableList.of(9))), 5); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 1, 2, 3, 4) .map(list::get).collect(Collectors.toList()); //Creating Replicas without replica Index @@ -224,7 +229,7 @@ public void testReplicasToFixMisreplicationWithThreeMisreplication() { 3, ImmutableList.of(3, 4, 8), 4, ImmutableList.of(9))), 5); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 1, 3, 4) .map(list::get).collect(Collectors.toList()); //Creating Replicas without replica Index for replicas < number of racks @@ -247,7 +252,7 @@ public void testReplicasToFixMisreplicationWithThreeMisreplication() { 3, ImmutableList.of(3, 4, 8), 4, ImmutableList.of(9))), 5); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List 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 @@ -262,7 +267,7 @@ public void testReplicasToFixMisreplicationMaxReplicaPerRack() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 2); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 2, 4, 6, 8) .map(list::get).collect(Collectors.toList()); List replicas = @@ -278,7 +283,7 @@ public void testReplicasToFixMisreplicationMaxReplicaPerRack() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 2); List racks = dummyPlacementPolicy.racks; - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 2, 4, 6, 8) .map(list::get).collect(Collectors.toList()); List replicas = @@ -297,7 +302,7 @@ public void testReplicasToFixMisreplicationMaxReplicaPerRack() { public void testReplicasWithoutMisreplication() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 5); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); List replicaDns = Stream.of(0, 1, 2, 3, 4) .map(list::get).collect(Collectors.toList()); Map replicas = @@ -314,7 +319,7 @@ public void testReplicasWithoutMisreplication() { public void testReplicasToRemoveWithOneOverreplication() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 5); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); Set replicas = Sets.newHashSet( HddsTestUtils.getReplicasWithReplicaIndex( ContainerID.valueOf(1), CLOSED, 0, 0, 0, list.subList(1, 6))); @@ -335,7 +340,7 @@ public void testReplicasToRemoveWithOneOverreplication() { public void testReplicasToRemoveWithTwoOverreplication() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 5); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); Set replicas = Sets.newHashSet( HddsTestUtils.getReplicasWithReplicaIndex( @@ -356,7 +361,7 @@ public void testReplicasToRemoveWithTwoOverreplication() { public void testReplicasToRemoveWith2CountPerUniqueReplica() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 3); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); Set replicas = Sets.newHashSet( HddsTestUtils.getReplicasWithReplicaIndex( @@ -382,7 +387,7 @@ public void testReplicasToRemoveWith2CountPerUniqueReplica() { public void testReplicasToRemoveWithoutReplicaIndex() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 3); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); Set replicas = Sets.newHashSet(HddsTestUtils.getReplicas( ContainerID.valueOf(1), CLOSED, 0, list.subList(0, 5))); @@ -402,7 +407,7 @@ public void testReplicasToRemoveWithoutReplicaIndex() { public void testReplicasToRemoveWithOverreplicationWithinSameRack() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 3); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); Set replicas = Sets.newHashSet( HddsTestUtils.getReplicasWithReplicaIndex( @@ -441,7 +446,7 @@ public void testReplicasToRemoveWithOverreplicationWithinSameRack() { public void testReplicasToRemoveWithNoOverreplication() { DummyPlacementPolicy dummyPlacementPolicy = new DummyPlacementPolicy(nodeManager, conf, 5); - List list = nodeManager.getAllNodes(); + List list = getAllNodes(nodeManager); Set replicas = Sets.newHashSet( HddsTestUtils.getReplicasWithReplicaIndex( ContainerID.valueOf(1), CLOSED, 0, 0, 0, list.subList(1, 6))); @@ -642,8 +647,9 @@ private static class DummyPlacementPolicy extends SCMCommonPlacementPolicy { when(node.getNetworkFullPath()).thenReturn(String.valueOf(i)); return node; }).collect(Collectors.toList()); - final List datanodeDetails = nodeManager.getAllNodes(); - rackMap = datanodeRackMap.entrySet().stream() + final List 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 1b8c4bbf384d..3dbe2b2bf0bb 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 healthyNodes; private final List staleNodes; private final List deadNodes; - private final Map nodeMetricMap; + private final Map nodeMetricMap = new TreeMap<>(); private final SCMNodeStat aggregateStat; private final Map>> 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 getNodes(NodeStatus status) { - return getNodes(status.getOperationalState(), status.getHealth()); + return getDatanodeDetails(status.getOperationalState(), status.getHealth()); } /** @@ -261,7 +261,16 @@ public List getNodes(NodeStatus status) { * @return List of Datanodes that are Heartbeating SCM. */ @Override - public List getNodes( + public List getNodes( + HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodestate) { + final List details = getDatanodeDetails(opState, nodestate); + if (details == null) { + return null; + } + return details.stream().map(this::getDatanodeInfo).collect(Collectors.toList()); + } + + private List 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 nodes = getNodes(opState, nodestate); + List 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 getAllNodes() { + public List getAllNodes() { // mock storage reports for TestDiskBalancer - List healthyNodesWithInfo = new ArrayList<>(); + List healthyNodesWithInfo = new ArrayList<>(); for (Map.Entry entry: nodeMetricMap.entrySet()) { NodeStatus nodeStatus = NodeStatus.inServiceHealthy(); @@ -399,7 +408,7 @@ public Map getNodeStats() { public List getMostOrLeastUsedDatanodes( boolean mostUsed) { List 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; } @@ -904,9 +915,9 @@ public List> 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 e47b979b4750..89c7c0c698e1 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 getNodes(NodeStatus nodeStatus) { } @Override - public List getNodes( + public List getNodes( NodeOperationalState opState, HddsProtos.NodeState health) { return null; } @@ -212,8 +211,8 @@ public int getNodeCount(NodeOperationalState opState, } @Override - public List getAllNodes() { - return null; + public List 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; @@ -362,7 +355,7 @@ public List> 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 e79c446c87f8..228908a6c714 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; @@ -102,7 +101,7 @@ public class TestContainerReportHandler { @BeforeEach void setup() throws IOException, InvalidStateTransitionException { 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); @@ -635,12 +634,11 @@ private List setupECContainerForTesting( container.getReplicationType()); final int numDatanodes = container.getReplicationConfig().getRequiredNodes(); - // Register required number of datanodes with NodeManager - List 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 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 1d54e404c6b3..a202111951dc 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; @@ -313,8 +312,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 49fab0afe225..e4ca940cf561 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 addReplicasToContainer(State... replicaStates) thr // If no states are specified, replica list will be empty. Set replicas = new HashSet<>(); try (MockNodeManager nodeManager = new MockNodeManager(true, replicaStates.length)) { - List nodes = nodeManager.getAllNodes(); + List 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 ffca82e231bd..fad2cf053f26 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 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 229d9283c74b..ef6f187b040c 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 341bbedf42d9..2871af693d70 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 decomNodes = nm.getNodes( + List 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 decomNodes = nm.getNodes( + List 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 maintenanceNodes = nm.getNodes( + List 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 ca0e231a2077..e1464f3e8794 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 cfdd4c7d37d7..232554209c5e 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 inServiceDatanodes = nm.getNodes(IN_SERVICE, HddsProtos.NodeState.HEALTHY); + final List 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 6e0686e13593..b8bee566a544 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 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 getInServiceDatanodes(NodeManager nm) { private List queryAllInServiceDatanodes( DiskBalancerQuery query) throws IOException { NodeManager nm = cluster.getStorageContainerManager().getScmNodeManager(); - List inServiceDatanodes = getInServiceDatanodes(nm); + final List inServiceDatanodes = nm.getNodes(IN_SERVICE, HddsProtos.NodeState.HEALTHY); List 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 f8985a0a8671..ec40c3a6a1a4 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 deadNodesOnRecon = nodeManager.getNodes(null, DEAD); + final Set deadNodesOnRecon = nodeManager.getNodes(null, DEAD).stream() + .map(info -> info.getID()) + .collect(Collectors.toSet()); if (!deadNodesOnRecon.isEmpty()) { List scmNodes = scmClient.getNodes(); List 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 c82eabb92e73..7d6a217c0f8f 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 9c93441457dd..508d391b32e5 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 nodes = + List nodes = reconNodeManager.getNodes(DECOMMISSIONING, null); assertEquals(1, nodes.size()); assertEquals(datanodeDetails.getID(), nodes.get(0).getID());