diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAndAdminContainerCLI.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAndAdminContainerCLI.java index 89ce3172c917..714703a5da6f 100644 --- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAndAdminContainerCLI.java +++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAndAdminContainerCLI.java @@ -60,6 +60,7 @@ import org.apache.hadoop.hdds.scm.container.ContainerHealthState; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport; import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager; import org.apache.hadoop.hdds.scm.node.NodeManager; @@ -235,11 +236,11 @@ void testMissingContainer() throws Exception { void testNodesInDecommissionOrMaintenance( NodeOperationalState initialState, NodeOperationalState finalState, boolean isMaintenance) throws Exception { - Pipeline pipeline = - scmClient.getContainerWithPipeline(containerIdR3).getPipeline(); + OzoneTestHelper.waitForStableReplicaCount(containerIdR3, 3, cluster); List details = - pipeline.getNodes().stream() + scmContainerManager.getContainerReplicas(ContainerID.valueOf(containerIdR3)).stream() + .map(ContainerReplica::getDatanodeDetails) .filter(d -> d.getPersistedOpState().equals(IN_SERVICE)) .collect(Collectors.toList()); @@ -272,7 +273,7 @@ void testNodesInDecommissionOrMaintenance( // a new replica-copy is made to another node. // For maintenance, there is no replica-copy in this case. if (!isMaintenance) { - OzoneTestHelper.waitForReplicaCount(containerIdR3, 4, cluster); + OzoneTestHelper.waitForStableReplicaCount(containerIdR3, 4, cluster); } compareRMReportToReconResponse(underReplicatedState); @@ -299,7 +300,7 @@ void testNodesInDecommissionOrMaintenance( // There will be a replica copy for both maintenance and decommission. // maintenance 3 -> 4, decommission 4 -> 5. int expectedReplicaNum = isMaintenance ? 4 : 5; - OzoneTestHelper.waitForReplicaCount(containerIdR3, expectedReplicaNum, cluster); + OzoneTestHelper.waitForStableReplicaCount(containerIdR3, expectedReplicaNum, cluster); compareRMReportToReconResponse(underReplicatedState); compareRMReportToReconResponse(overReplicatedState); @@ -316,6 +317,8 @@ void testNodesInDecommissionOrMaintenance( NodeTestUtil.waitForDnToReachPersistedOpState(nodeToGoOffline1, IN_SERVICE); NodeTestUtil.waitForDnToReachPersistedOpState(nodeToGoOffline2, IN_SERVICE); + OzoneTestHelper.waitForStableReplicaCount(containerIdR3, 3, cluster); + compareRMReportToReconResponse(underReplicatedState); compareRMReportToReconResponse(overReplicatedState); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/OzoneTestHelper.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/OzoneTestHelper.java index 9558be282560..198d02106a3d 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/OzoneTestHelper.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/OzoneTestHelper.java @@ -51,6 +51,8 @@ import org.apache.hadoop.hdds.scm.container.ContainerManager; import org.apache.hadoop.hdds.scm.container.ContainerNotFoundException; import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ContainerHealthResult; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager; import org.apache.hadoop.hdds.scm.events.SCMEvents; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.scm.pipeline.PipelineNotFoundException; @@ -465,6 +467,32 @@ public static void waitForReplicaCount(long containerID, int count, 200, 30000); } + /** + * Like {@link #waitForReplicaCount(long, int, MiniOzoneCluster)}, but only returns once + * replication has quiesced: the container is healthy with {@code count} replicas and no pending + * add or delete ops. Checking the current replication health avoids treating an empty pending-op + * list as settled before ReplicationManager has evaluated the latest replica or node state. + */ + public static void waitForStableReplicaCount(long containerID, int count, MiniOzoneCluster cluster) + throws TimeoutException, InterruptedException { + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + ReplicationManager replicationManager = cluster.getStorageContainerManager() + .getReplicationManager(); + ContainerID cid = ContainerID.valueOf(containerID); + GenericTestUtils.waitFor(() -> { + try { + ContainerInfo container = containerManager.getContainer(cid); + Set replicas = containerManager.getContainerReplicas(cid); + return replicas.size() == count + && replicationManager.getPendingReplicationOps(cid).isEmpty() + && replicationManager.getContainerReplicationHealth(container, replicas).getHealthState() + == ContainerHealthResult.HealthState.HEALTHY; + } catch (ContainerNotFoundException e) { + return false; + } + }, 200, 30000); + } + /** * Wait until SCM reports exactly {@code count} replicas for the container and every replica is in {@code state}. * Unlike {@link #waitForContainerStateInSCM}, which checks the container's aggregate state (it flips as soon as the