Skip to content
Merged
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 @@ -31,7 +31,6 @@
import static org.apache.ozone.test.GenericTestUtils.setLogLevel;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.any;
Expand All @@ -40,12 +39,10 @@
import com.google.common.collect.ImmutableMap;
import java.io.IOException;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
Expand All @@ -71,7 +68,6 @@
import org.apache.hadoop.ozone.client.ObjectStore;
import org.apache.hadoop.ozone.client.OzoneBucket;
import org.apache.hadoop.ozone.client.OzoneClient;
import org.apache.hadoop.ozone.client.OzoneClientFactory;
import org.apache.hadoop.ozone.client.OzoneVolume;
import org.apache.hadoop.ozone.client.io.OzoneInputStream;
import org.apache.hadoop.ozone.container.common.interfaces.Container;
Expand All @@ -85,7 +81,6 @@
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.MockedStatic;
import org.mockito.Mockito;
Expand All @@ -105,12 +100,10 @@ class TestContainerReplication {
SCMContainerPlacementRandom.class
);

static List<Arguments> containerReplicationArguments() {
List<Arguments> arguments = new LinkedList<>();
static List<String> containerReplicationArguments() {
List<String> arguments = new LinkedList<>();
for (Class<? extends PlacementPolicy> policyClass : POLICIES) {
String canonicalName = policyClass.getCanonicalName();
arguments.add(Arguments.arguments(canonicalName, true));
arguments.add(Arguments.arguments(canonicalName, false));
arguments.add(policyClass.getCanonicalName());
}
return arguments;
}
Expand All @@ -122,27 +115,30 @@ static void setUp() {
setLogLevel(SCMContainerPlacementRandom.class, Level.DEBUG);
}

/**
* Verifies that a closed RATIS THREE container which becomes under-replicated
* after a datanode shutdown is restored to three replicas by ReplicationManager,
* and that the configured placement policy records the datanode-choose metrics.
* Runs once per placement policy in {@link #containerReplicationArguments()}.
*/
@ParameterizedTest
@MethodSource("containerReplicationArguments")
void testContainerReplication(
String placementPolicyClass, boolean legacyEnabled) throws Exception {
void testRatisContainerReReplicationAfterDatanodeShutdown(String placementPolicyClass) throws Exception {

OzoneConfiguration conf = createConfiguration(legacyEnabled);
OzoneConfiguration conf = createConfiguration();
conf.set(OZONE_SCM_CONTAINER_PLACEMENT_IMPL_KEY, placementPolicyClass);
try (MiniOzoneCluster cluster = newCluster(conf)) {
try (MiniOzoneCluster cluster = MiniOzoneCluster.newBuilder(conf).setNumDatanodes(5).build()) {
cluster.waitForClusterToBeReady();
SCMContainerPlacementMetrics metrics = cluster.getStorageContainerManager().getPlacementMetrics();
try (OzoneClient client = cluster.newClient()) {
createTestData(client);

List<OmKeyLocationInfo> keyLocations = lookupKey(cluster);
assertThat(keyLocations).isNotEmpty();
long datanodeChooseAttemptCount = metrics.getDatanodeChooseAttemptCount();
long datanodeChooseSuccessCount = metrics.getDatanodeChooseSuccessCount();
long datanodeChooseFallbackCount = metrics.getDatanodeChooseFallbackCount();
long datanodeRequestCount = metrics.getDatanodeRequestCount();

OmKeyLocationInfo keyLocation = keyLocations.get(0);
OmKeyLocationInfo keyLocation = lookupKeyFirstLocation(cluster);
long containerID = keyLocation.getContainerID();
waitForContainerClose(cluster, containerID);

Expand All @@ -151,7 +147,7 @@ void testContainerReplication(

waitForReplicaCount(containerID, 3, cluster);

Supplier<String> messageSupplier = () -> "policy=" + placementPolicyClass + " legacy=" + legacyEnabled;
Supplier<String> messageSupplier = () -> "policy=" + placementPolicyClass;
assertEquals(datanodeRequestCount + 1, metrics.getDatanodeRequestCount(), messageSupplier);
assertThat(metrics.getDatanodeChooseAttemptCount()).isGreaterThan(datanodeChooseAttemptCount);
assertEquals(datanodeChooseSuccessCount + 1, metrics.getDatanodeChooseSuccessCount(), messageSupplier);
Expand All @@ -160,14 +156,7 @@ void testContainerReplication(
}
}

private static MiniOzoneCluster newCluster(OzoneConfiguration conf)
throws IOException {
return MiniOzoneCluster.newBuilder(conf)
.setNumDatanodes(5)
.build();
}

private static OzoneConfiguration createConfiguration(boolean enableLegacy) {
private static OzoneConfiguration createConfiguration() {
OzoneConfiguration conf = new OzoneConfiguration();
conf.setTimeDuration(OZONE_SCM_STALENODE_INTERVAL, 3, TimeUnit.SECONDS);
conf.setTimeDuration(OZONE_SCM_DEADNODE_INTERVAL, 6, TimeUnit.SECONDS);
Expand All @@ -181,47 +170,23 @@ private static OzoneConfiguration createConfiguration(boolean enableLegacy) {
return conf;
}

// TODO use common helper to create test data
private void createTestData(OzoneClient client) throws IOException {
ObjectStore objectStore = client.getObjectStore();
objectStore.createVolume(VOLUME);
OzoneVolume volume = objectStore.getVolume(VOLUME);
volume.createBucket(BUCKET);

OzoneBucket bucket = volume.getBucket(BUCKET);
OzoneBucket bucket = TestDataUtil.createVolumeAndBucket(client, VOLUME, BUCKET);

TestDataUtil.createKey(bucket, KEY,
RatisReplicationConfig.getInstance(THREE),
"Hello".getBytes(UTF_8));
}

private byte[] createTestData(OzoneClient client, int size) throws IOException {
ObjectStore objectStore = client.getObjectStore();
objectStore.createVolume(VOLUME);
OzoneVolume volume = objectStore.getVolume(VOLUME);
volume.createBucket(BUCKET);
OzoneBucket bucket = volume.getBucket(BUCKET);
OzoneBucket bucket = TestDataUtil.createVolumeAndBucket(client, VOLUME, BUCKET);

byte[] b = new byte[size];
b = RandomUtils.secure().randomBytes(b.length);
byte[] b = RandomUtils.secure().randomBytes(size);
TestDataUtil.createKey(bucket, KEY,
new ECReplicationConfig("RS-3-2-1k"), b);
return b;
}

private static List<OmKeyLocationInfo> lookupKey(MiniOzoneCluster cluster)
throws IOException {
OmKeyArgs keyArgs = new OmKeyArgs.Builder()
.setVolumeName(VOLUME)
.setBucketName(BUCKET)
.setKeyName(KEY)
.build();
OmKeyInfo keyInfo = cluster.getOzoneManager().lookupKey(keyArgs);
OmKeyLocationInfoGroup locations = keyInfo.getLatestVersionLocations();
assertNotNull(locations);
return locations.getLocationList();
}

private static OmKeyLocationInfo lookupKeyFirstLocation(MiniOzoneCluster cluster)
throws IOException {
OmKeyArgs keyArgs = new OmKeyArgs.Builder()
Expand Down Expand Up @@ -278,12 +243,12 @@ private static void deleteContainer(MiniOzoneCluster cluster, DatanodeDetails dn

@Test
public void testImportedContainerIsClosed() throws Exception {
OzoneConfiguration conf = createConfiguration(false);
OzoneConfiguration conf = createConfiguration();
// create a 4 node cluster
try (MiniOzoneCluster cluster = MiniOzoneCluster.newBuilder(conf).setNumDatanodes(4).build()) {
cluster.waitForClusterToBeReady();

try (OzoneClient client = OzoneClientFactory.getRpcClient(conf)) {
try (OzoneClient client = cluster.newClient()) {
List<DatanodeDetails> allNodes =
cluster.getHddsDatanodes().stream()
.map(HddsDatanodeService::getDatanodeDetails)
Expand Down Expand Up @@ -315,9 +280,9 @@ public void testImportedContainerIsClosed() throws Exception {
@Test
@Flaky("HDDS-11087")
public void testECContainerReplication() throws Exception {
OzoneConfiguration conf = createConfiguration(false);
OzoneConfiguration conf = createConfiguration();
final Map<Integer, Integer> failedReadChunkCountMap = new ConcurrentHashMap<>();
// Overiding Config to support 1k Chunk size
// Overriding Config to support 1k Chunk size
conf.set("ozone.replication.allowed-configs", "(^((STANDALONE|RATIS)/(ONE|THREE))|(EC/(3-2|6-3|10-4)-" +
"(512|1024|2048|4096|1)k)$)");
conf.set(OZONE_SCM_CONTAINER_PLACEMENT_EC_IMPL_KEY, SCMContainerPlacementRackScatter.class.getCanonicalName());
Expand All @@ -327,20 +292,7 @@ public void testECContainerReplication() throws Exception {
// Creating Cluster with 5 Nodes
try (MiniOzoneCluster cluster = MiniOzoneCluster.newBuilder(conf).setNumDatanodes(5).build()) {
cluster.waitForClusterToBeReady();
try (OzoneClient client = OzoneClientFactory.getRpcClient(conf)) {
Set<DatanodeDetails> allNodes =
cluster.getHddsDatanodes().stream().map(HddsDatanodeService::getDatanodeDetails).collect(
Collectors.toSet());
List<DatanodeDetails> initialNodesWithData = new ArrayList<>();
// Keeping 5 DNs and stopping the 6th Node here it is kept in the var extraNodes
for (DatanodeDetails dn : allNodes) {
if (initialNodesWithData.size() < 5) {
initialNodesWithData.add(dn);
} else {
cluster.shutdownHddsDatanode(dn);
}
}

try (OzoneClient client = cluster.newClient()) {
// Creating 2 stripes with Chunk Size 1k
int size = 6 * 1024;
byte[] originalData = createTestData(client, size);
Expand All @@ -350,6 +302,12 @@ public void testECContainerReplication() throws Exception {
long containerID = keyLocation.getContainerID();
waitForContainerClose(cluster, containerID);

// The cluster has 5 datanodes and the key is written as EC RS-3-2 (3 data + 2 parity = 5
// replica indices), so every datanode now holds exactly one replica index.
List<DatanodeDetails> initialNodesWithData =
cluster.getHddsDatanodes().stream().map(HddsDatanodeService::getDatanodeDetails)
.collect(Collectors.toList());

// Forming Replica Index Map
Map<Integer, DatanodeDetails> replicaIndexMap =
initialNodesWithData.stream().map(dn -> new Object[]{dn, keyLocation.getPipeline().getReplicaIndex(dn)})
Expand Down