smengcl commented on code in PR #11146: URL: https://github.com/apache/ozone/pull/11146#discussion_r4100707068
########## hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java: ########## @@ -0,0 +1,545 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.ozone.container.replication; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.util.Collections.singleton; +import static java.util.Collections.singletonList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_NODE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.CLOSED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.DEAD; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DATANODE_ADMIN_MONITOR_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.getDNHostAndPort; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachHealthState; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachOpState; +import static org.apache.hadoop.hdds.scm.pipeline.MockPipeline.createPipeline; +import static org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls.createContainer; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForContainerClose; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForReplicaCount; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; +import java.util.stream.Collectors; +import org.apache.hadoop.hdds.HddsConfigKeys; +import org.apache.hadoop.hdds.client.BlockID; +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.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.DatanodeID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.scm.ScmConfigKeys; +import org.apache.hadoop.hdds.scm.XceiverClientFactory; +import org.apache.hadoop.hdds.scm.XceiverClientManager; +import org.apache.hadoop.hdds.scm.XceiverClientSpi; +import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerInfo; +import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration; +import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.DataTestUtil; +import org.apache.hadoop.ozone.HddsDatanodeService; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.UniformDatanodesFactory; +import org.apache.hadoop.ozone.client.OzoneBucket; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.OzoneKeyDetails; +import org.apache.hadoop.ozone.container.ContainerTestHelper; +import org.apache.hadoop.ozone.container.common.interfaces.Container; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.volume.HddsVolume; +import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.dn.DatanodeTestUtils; +import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; +import org.apache.ozone.test.GenericTestUtils; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.TestMethodOrder; +import org.junit.jupiter.api.parallel.Execution; +import org.junit.jupiter.api.parallel.ExecutionMode; +import org.junit.jupiter.api.parallel.ResourceLock; + +/** + * Integration tests for per-volume push replication thread pools (HDDS-15412). + */ +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +@Execution(ExecutionMode.SAME_THREAD) +@ResourceLock("MiniOzoneCluster") +class TestPerVolumePushReplication { + + private static final AtomicLong CONTAINER_ID = new AtomicLong(1_000_000L); + private static final int DATA_VOLUMES = 2; + private static final int DATANODE_COUNT = 7; + private static final String VOLUME = "vol1"; + private static final String BUCKET = "bucket1"; + private static final RatisReplicationConfig RATIS_THREE = + RatisReplicationConfig.getInstance(THREE); + private static final ECReplicationConfig EC_REP = new ECReplicationConfig(3, 2); + + private MiniOzoneCluster cluster; + private XceiverClientFactory clientFactory; + private OzoneBucket bucket; + + @BeforeAll + void setUp() throws Exception { + OzoneConfiguration conf = createSharedConfig(); + cluster = newCluster(conf, DATANODE_COUNT); + cluster.waitForClusterToBeReady(); + clientFactory = new XceiverClientManager(conf); + try (OzoneClient client = cluster.newClient()) { + bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET); + } + } + + @AfterAll + void tearDown() { + IOUtils.closeQuietly(clientFactory, cluster); + } + + @Order(1) + @Test + void testPushAndScmReplicationWithPerVolumeEnabled() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(0); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + long containerId = createClosedContainer(clientFactory, source, 0L); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerId, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerId); + getContainer(cluster, target, containerId); + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + DataTestUtil.createKey(bucket, "pushKey1", RATIS_THREE, "data".getBytes(UTF_8)); + OzoneKeyDetails keyDetails = bucket.getKey("pushKey1"); + long scmContainerId = keyDetails.getOzoneKeyLocations().get(0).getContainerID(); + waitForContainerClose(cluster, scmContainerId); + + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + Set<ContainerReplica> replicas = + containerManager.getContainerReplicas(ContainerID.valueOf(scmContainerId)); + DatanodeDetails replicaDn = replicas.iterator().next().getDatanodeDetails(); + cluster.shutdownHddsDatanode(replicaDn); + waitForReplicaCount(scmContainerId, 3, cluster); + } + + @Order(2) + @Test + void testHealthyVolumeReplicationAfterVolumeFailure() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(1); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + MutableVolumeSet volSet = sourceDn.getDatanodeStateMachine().getContainer().getVolumeSet(); + HddsVolume vol0 = (HddsVolume) volSet.getVolumesList().get(0); + HddsVolume vol1 = (HddsVolume) volSet.getVolumesList().get(1); + + long containerOnVol0 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol0, 0L); + long containerOnVol1 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol1, 0L); + + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + triggerAndWaitForVolumeFailure(volSet, vol0); + waitForVolumePoolState(sourceDn, vol0, vol1); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerOnVol1, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerOnVol1); + getContainer(cluster, target, containerOnVol1); + assertEquals(1, volSet.getFailedVolumesList().size()); + + // Task routes via global pool fallback (HDDS-15327); replication fails on bad volume. + ReplicateContainerCommand failedVolCmd = + ReplicateContainerCommand.toTarget(containerOnVol0, target); + ReplicationSupervisor supervisor = + sourceDn.getDatanodeStateMachine().getSupervisor(); + long previousFailures = supervisor.getReplicationFailureCount(); + queuePushAndWaitForFailure(cluster, failedVolCmd, source, target, + containerOnVol0, supervisor, previousFailures); + assertTrue(supervisor.getReplicationFailureCount() >= previousFailures + 1); + + DatanodeTestUtils.restoreBadVolume(vol0); + } + + @Order(3) + @Test + void testDecommissionWithPerVolumePools() throws Exception { + try (ContainerOperationClient scmClient = new ContainerOperationClient(cluster.getConf())) { + StorageContainerManager scm = cluster.getStorageContainerManager(); + NodeManager nm = scm.getScmNodeManager(); + ContainerManager cm = scm.getContainerManager(); + PipelineManager pm = scm.getPipelineManager(); + + generateData(bucket, 20, "decomKey", RATIS_THREE); + generateData(bucket, 20, "decomEcKey", EC_REP); + + ContainerInfo ratisContainer = waitForKeyContainer(bucket, cm, "decomKey0", 3); + ContainerInfo ecContainer = waitForKeyContainer(bucket, cm, "decomEcKey0", 5); + Pipeline ratisPipeline = pm.getPipeline(ratisContainer.getPipelineID()); + Pipeline ecPipeline = pm.getPipeline(ecContainer.getPipelineID()); + + DatanodeID dnId = ratisPipeline.getNodes().stream() + .filter(node -> ecPipeline.getNodes().contains(node)) + .findFirst() + .orElseThrow(() -> new AssertionError("no intersecting datanode found")) + .getID(); + DatanodeDetails toDecommission = nm.getNode(dnId); + + scmClient.decommissionNodes(singletonList(getDNHostAndPort(toDecommission)), false); + waitForDnToReachOpState(nm, toDecommission, DECOMMISSIONED); + + waitForContainerReplicas(cm, ratisContainer, 4); + waitForContainerReplicas(cm, ecContainer, 6); + + cluster.shutdownHddsDatanode(toDecommission); + waitForDnToReachHealthState(nm, toDecommission, DEAD); + + waitForContainerReplicas(cm, ratisContainer, 3); + waitForContainerReplicas(cm, ecContainer, 5); + + DataTestUtil.createKey(bucket, "sanityKey", RATIS_THREE, + "still healthy".getBytes(StandardCharsets.UTF_8)); + } + } + + private HddsDatanodeService selectHealthyDatanode(int indexAmongHealthy) { + List<HddsDatanodeService> healthy = cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .collect(Collectors.toList()); + if (indexAmongHealthy >= healthy.size()) { + throw new AssertionError("not enough healthy datanodes: requested index " + + indexAmongHealthy + ", found " + healthy.size()); + } + return healthy.get(indexAmongHealthy); + } + + private DatanodeDetails selectOtherHealthyNode(DatanodeDetails source) { + return cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .map(HddsDatanodeService::getDatanodeDetails) + .filter(dn -> !dn.equals(source)) + .findAny() + .orElseThrow(() -> new AssertionError("no target datanode found")); + } + + private boolean isHealthyDatanode(HddsDatanodeService datanode) { + if (datanode.getDatanodeDetails().getPersistedOpState() != IN_SERVICE) { + return false; + } + MutableVolumeSet volumeSet = + datanode.getDatanodeStateMachine().getContainer().getVolumeSet(); + return volumeSet.getFailedVolumesList().isEmpty() + && volumeSet.getVolumesList().size() == DATA_VOLUMES; + } + + private static void assertVolumePools(HddsDatanodeService dn, + int expectedVolumeCount, int expectedPoolSize) { + VolumeReplicationThreadPools pools = + dn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + List<? extends StorageVolume> volumes = + dn.getDatanodeStateMachine().getContainer().getVolumeSet().getVolumesList(); + assertEquals(expectedVolumeCount, volumes.size()); + for (StorageVolume volume : volumes) { + String volumeRoot = volume.getStorageDir().getPath(); + assertTrue(pools.hasPool(volumeRoot), "missing pool for " + volumeRoot); + assertEquals(expectedPoolSize, pools.getPoolSize(volumeRoot)); + } + } + + private static void waitForVolumePoolState(HddsDatanodeService sourceDn, + HddsVolume failedVolume, HddsVolume healthyVolume) + throws TimeoutException, InterruptedException { + String failedPath = failedVolume.getStorageDir().getPath(); + String healthyPath = healthyVolume.getStorageDir().getPath(); + GenericTestUtils.waitFor((BooleanSupplier) () -> { + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + return pools != null + && !pools.hasPool(failedPath) + && pools.hasPool(healthyPath); + }, 100, 60000); + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + assertFalse(pools.hasPool(failedPath)); + assertTrue(pools.hasPool(healthyPath)); + } + + private static void queuePushAndWaitForContainer(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + hasContainer(cluster, target, containerId), 100, 30000); + } + + private static void queuePushAndWaitForFailure(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId, ReplicationSupervisor supervisor, + long previousFailureCount) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + supervisor.getReplicationFailureCount() >= previousFailureCount + 1 + && !hasContainer(cluster, target, containerId), + 100, 30000); + } + + private static void queueReplicationCommand(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source) throws IOException { + DatanodeStateMachine stateMachine = cluster.getHddsDatanode(source).getDatanodeStateMachine(); + StateContext context = stateMachine.getContext(); + context.getTermOfLeaderSCM().ifPresent(cmd::setTerm); + context.addCommand(cmd); + } + + private static boolean hasContainer(MiniOzoneCluster cluster, + DatanodeDetails datanode, long containerId) { + try { + return cluster.getHddsDatanode(datanode).getDatanodeStateMachine().getContainer() + .getContainerSet().getContainer(containerId) != null; + } catch (IOException e) { + return false; + } + } + + private static long findOrCreateContainerOnVolume(MiniOzoneCluster cluster, + XceiverClientFactory clientFactory, DatanodeDetails dn, HddsVolume targetVolume, + long dataSize) throws Exception { + for (int attempt = 0; attempt < 30; attempt++) { + long containerId = createClosedContainer(clientFactory, dn, dataSize); + Container<?> container = getContainer(cluster, dn, containerId); + if (targetVolume.equals(container.getContainerData().getVolume())) { + return containerId; + } + } + throw new AssertionError("Could not place container on volume " + targetVolume); + } + + private static long createClosedContainer(XceiverClientFactory clientFactory, + DatanodeDetails dn, long dataSize) throws Exception { + long containerId = CONTAINER_ID.incrementAndGet(); + try (XceiverClientSpi client = clientFactory.acquireClient(createPipeline(singleton(dn)))) { + if (dataSize <= 0) { + createContainer(client, containerId, null, CLOSED, 0); + return containerId; + } + createContainer(client, containerId, null); Review Comment: Addressed in 66bf80e1221. The `dataSize` parameter and the dead branch are gone, so `createClosedContainer` is now create, close, return the id. ########## hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java: ########## @@ -0,0 +1,545 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.ozone.container.replication; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.util.Collections.singleton; +import static java.util.Collections.singletonList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_NODE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.CLOSED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.DEAD; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DATANODE_ADMIN_MONITOR_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.getDNHostAndPort; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachHealthState; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachOpState; +import static org.apache.hadoop.hdds.scm.pipeline.MockPipeline.createPipeline; +import static org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls.createContainer; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForContainerClose; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForReplicaCount; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; +import java.util.stream.Collectors; +import org.apache.hadoop.hdds.HddsConfigKeys; +import org.apache.hadoop.hdds.client.BlockID; +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.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.DatanodeID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.scm.ScmConfigKeys; +import org.apache.hadoop.hdds.scm.XceiverClientFactory; +import org.apache.hadoop.hdds.scm.XceiverClientManager; +import org.apache.hadoop.hdds.scm.XceiverClientSpi; +import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerInfo; +import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration; +import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.DataTestUtil; +import org.apache.hadoop.ozone.HddsDatanodeService; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.UniformDatanodesFactory; +import org.apache.hadoop.ozone.client.OzoneBucket; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.OzoneKeyDetails; +import org.apache.hadoop.ozone.container.ContainerTestHelper; +import org.apache.hadoop.ozone.container.common.interfaces.Container; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.volume.HddsVolume; +import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.dn.DatanodeTestUtils; +import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; +import org.apache.ozone.test.GenericTestUtils; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.TestMethodOrder; +import org.junit.jupiter.api.parallel.Execution; +import org.junit.jupiter.api.parallel.ExecutionMode; +import org.junit.jupiter.api.parallel.ResourceLock; + +/** + * Integration tests for per-volume push replication thread pools (HDDS-15412). + */ +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +@Execution(ExecutionMode.SAME_THREAD) +@ResourceLock("MiniOzoneCluster") +class TestPerVolumePushReplication { + + private static final AtomicLong CONTAINER_ID = new AtomicLong(1_000_000L); + private static final int DATA_VOLUMES = 2; + private static final int DATANODE_COUNT = 7; + private static final String VOLUME = "vol1"; + private static final String BUCKET = "bucket1"; + private static final RatisReplicationConfig RATIS_THREE = + RatisReplicationConfig.getInstance(THREE); + private static final ECReplicationConfig EC_REP = new ECReplicationConfig(3, 2); + + private MiniOzoneCluster cluster; + private XceiverClientFactory clientFactory; + private OzoneBucket bucket; + + @BeforeAll + void setUp() throws Exception { + OzoneConfiguration conf = createSharedConfig(); + cluster = newCluster(conf, DATANODE_COUNT); + cluster.waitForClusterToBeReady(); + clientFactory = new XceiverClientManager(conf); + try (OzoneClient client = cluster.newClient()) { + bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET); + } + } + + @AfterAll + void tearDown() { + IOUtils.closeQuietly(clientFactory, cluster); + } + + @Order(1) + @Test + void testPushAndScmReplicationWithPerVolumeEnabled() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(0); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + long containerId = createClosedContainer(clientFactory, source, 0L); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerId, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerId); + getContainer(cluster, target, containerId); + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + DataTestUtil.createKey(bucket, "pushKey1", RATIS_THREE, "data".getBytes(UTF_8)); + OzoneKeyDetails keyDetails = bucket.getKey("pushKey1"); + long scmContainerId = keyDetails.getOzoneKeyLocations().get(0).getContainerID(); + waitForContainerClose(cluster, scmContainerId); + + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + Set<ContainerReplica> replicas = + containerManager.getContainerReplicas(ContainerID.valueOf(scmContainerId)); + DatanodeDetails replicaDn = replicas.iterator().next().getDatanodeDetails(); + cluster.shutdownHddsDatanode(replicaDn); + waitForReplicaCount(scmContainerId, 3, cluster); + } + + @Order(2) + @Test + void testHealthyVolumeReplicationAfterVolumeFailure() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(1); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + MutableVolumeSet volSet = sourceDn.getDatanodeStateMachine().getContainer().getVolumeSet(); + HddsVolume vol0 = (HddsVolume) volSet.getVolumesList().get(0); + HddsVolume vol1 = (HddsVolume) volSet.getVolumesList().get(1); + + long containerOnVol0 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol0, 0L); + long containerOnVol1 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol1, 0L); + + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + triggerAndWaitForVolumeFailure(volSet, vol0); + waitForVolumePoolState(sourceDn, vol0, vol1); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerOnVol1, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerOnVol1); + getContainer(cluster, target, containerOnVol1); + assertEquals(1, volSet.getFailedVolumesList().size()); + + // Task routes via global pool fallback (HDDS-15327); replication fails on bad volume. + ReplicateContainerCommand failedVolCmd = + ReplicateContainerCommand.toTarget(containerOnVol0, target); + ReplicationSupervisor supervisor = + sourceDn.getDatanodeStateMachine().getSupervisor(); + long previousFailures = supervisor.getReplicationFailureCount(); + queuePushAndWaitForFailure(cluster, failedVolCmd, source, target, + containerOnVol0, supervisor, previousFailures); + assertTrue(supervisor.getReplicationFailureCount() >= previousFailures + 1); + + DatanodeTestUtils.restoreBadVolume(vol0); + } + + @Order(3) + @Test + void testDecommissionWithPerVolumePools() throws Exception { + try (ContainerOperationClient scmClient = new ContainerOperationClient(cluster.getConf())) { + StorageContainerManager scm = cluster.getStorageContainerManager(); + NodeManager nm = scm.getScmNodeManager(); + ContainerManager cm = scm.getContainerManager(); + PipelineManager pm = scm.getPipelineManager(); + + generateData(bucket, 20, "decomKey", RATIS_THREE); + generateData(bucket, 20, "decomEcKey", EC_REP); + + ContainerInfo ratisContainer = waitForKeyContainer(bucket, cm, "decomKey0", 3); + ContainerInfo ecContainer = waitForKeyContainer(bucket, cm, "decomEcKey0", 5); + Pipeline ratisPipeline = pm.getPipeline(ratisContainer.getPipelineID()); + Pipeline ecPipeline = pm.getPipeline(ecContainer.getPipelineID()); + + DatanodeID dnId = ratisPipeline.getNodes().stream() + .filter(node -> ecPipeline.getNodes().contains(node)) + .findFirst() + .orElseThrow(() -> new AssertionError("no intersecting datanode found")) + .getID(); + DatanodeDetails toDecommission = nm.getNode(dnId); + + scmClient.decommissionNodes(singletonList(getDNHostAndPort(toDecommission)), false); + waitForDnToReachOpState(nm, toDecommission, DECOMMISSIONED); + + waitForContainerReplicas(cm, ratisContainer, 4); + waitForContainerReplicas(cm, ecContainer, 6); + + cluster.shutdownHddsDatanode(toDecommission); + waitForDnToReachHealthState(nm, toDecommission, DEAD); + + waitForContainerReplicas(cm, ratisContainer, 3); + waitForContainerReplicas(cm, ecContainer, 5); + + DataTestUtil.createKey(bucket, "sanityKey", RATIS_THREE, + "still healthy".getBytes(StandardCharsets.UTF_8)); + } + } + + private HddsDatanodeService selectHealthyDatanode(int indexAmongHealthy) { + List<HddsDatanodeService> healthy = cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .collect(Collectors.toList()); + if (indexAmongHealthy >= healthy.size()) { + throw new AssertionError("not enough healthy datanodes: requested index " + + indexAmongHealthy + ", found " + healthy.size()); + } + return healthy.get(indexAmongHealthy); + } + + private DatanodeDetails selectOtherHealthyNode(DatanodeDetails source) { + return cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .map(HddsDatanodeService::getDatanodeDetails) + .filter(dn -> !dn.equals(source)) + .findAny() + .orElseThrow(() -> new AssertionError("no target datanode found")); + } + + private boolean isHealthyDatanode(HddsDatanodeService datanode) { + if (datanode.getDatanodeDetails().getPersistedOpState() != IN_SERVICE) { + return false; + } + MutableVolumeSet volumeSet = + datanode.getDatanodeStateMachine().getContainer().getVolumeSet(); + return volumeSet.getFailedVolumesList().isEmpty() + && volumeSet.getVolumesList().size() == DATA_VOLUMES; + } + + private static void assertVolumePools(HddsDatanodeService dn, + int expectedVolumeCount, int expectedPoolSize) { + VolumeReplicationThreadPools pools = + dn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + List<? extends StorageVolume> volumes = + dn.getDatanodeStateMachine().getContainer().getVolumeSet().getVolumesList(); + assertEquals(expectedVolumeCount, volumes.size()); + for (StorageVolume volume : volumes) { + String volumeRoot = volume.getStorageDir().getPath(); + assertTrue(pools.hasPool(volumeRoot), "missing pool for " + volumeRoot); + assertEquals(expectedPoolSize, pools.getPoolSize(volumeRoot)); + } + } + + private static void waitForVolumePoolState(HddsDatanodeService sourceDn, + HddsVolume failedVolume, HddsVolume healthyVolume) + throws TimeoutException, InterruptedException { + String failedPath = failedVolume.getStorageDir().getPath(); + String healthyPath = healthyVolume.getStorageDir().getPath(); + GenericTestUtils.waitFor((BooleanSupplier) () -> { + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + return pools != null + && !pools.hasPool(failedPath) + && pools.hasPool(healthyPath); + }, 100, 60000); + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + assertFalse(pools.hasPool(failedPath)); + assertTrue(pools.hasPool(healthyPath)); + } + + private static void queuePushAndWaitForContainer(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + hasContainer(cluster, target, containerId), 100, 30000); + } + + private static void queuePushAndWaitForFailure(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId, ReplicationSupervisor supervisor, + long previousFailureCount) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + supervisor.getReplicationFailureCount() >= previousFailureCount + 1 + && !hasContainer(cluster, target, containerId), + 100, 30000); + } + + private static void queueReplicationCommand(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source) throws IOException { + DatanodeStateMachine stateMachine = cluster.getHddsDatanode(source).getDatanodeStateMachine(); + StateContext context = stateMachine.getContext(); + context.getTermOfLeaderSCM().ifPresent(cmd::setTerm); Review Comment: Not addressed in 66bf80e1221, deliberately, so leaving this thread open to track it. Reusing `TestContainerReplication#queueAndWaitForCompletion`, `CONTAINER_ID` and `createNewClosedContainer` means widening their visibility and lifting them somewhere shared such as `OzoneTestHelper`, which edits a test class this PR does not otherwise touch. The duplication is real and worth removing, but I would rather do it as a follow up where the shared move is reviewable on its own instead of mixed into a new test class. Happy to do it here instead if you prefer. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
