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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<DatanodeDetails> details =
pipeline.getNodes().stream()
scmContainerManager.getContainerReplicas(ContainerID.valueOf(containerIdR3)).stream()
.map(ContainerReplica::getDatanodeDetails)
.filter(d -> d.getPersistedOpState().equals(IN_SERVICE))
.collect(Collectors.toList());

Expand Down Expand Up @@ -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);
Expand All @@ -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);
Expand All @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<ContainerReplica> 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
Expand Down
Loading