This is an automated email from the ASF dual-hosted git repository.
adoroszlai pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 1a8396cb560 HDDS-16100. Parameterize TestChunkInputStream (#11117)
1a8396cb560 is described below
commit 1a8396cb560890a17397d0dfc036b015669ab6c6
Author: Doroszlai, Attila <[email protected]>
AuthorDate: Fri Aug 28 08:00:30 2026 +0200
HDDS-16100. Parameterize TestChunkInputStream (#11117)
---
.../hdds/scm/storage/DomainSocketFactory.java | 4 +
.../org/apache/ozone/test/GenericTestUtils.java | 7 +-
.../ozone/client/rpc/read/InputStreamTests.java | 43 +++++++--
.../client/rpc/read/TestChunkInputStream.java | 71 ++++++++++++++-
.../client/rpc/read/TestLocalChunkInputStream.java | 100 +++++++--------------
.../rpc/read/TestStreamBlockInputStream.java | 69 +++++++-------
6 files changed, 179 insertions(+), 115 deletions(-)
diff --git
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/DomainSocketFactory.java
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/DomainSocketFactory.java
index 6ccaffe0fe7..fa4ffae9631 100644
---
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/DomainSocketFactory.java
+++
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/DomainSocketFactory.java
@@ -133,6 +133,10 @@ public static DomainSocketFactory
getInstance(ConfigurationSource conf) {
return instance;
}
+ public static boolean isNativeLibraryLoaded() {
+ return nativeLibraryLoaded;
+ }
+
private DomainSocketFactory(ConfigurationSource conf) {
OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class);
boolean shortCircuitEnabled = clientConfig.isShortCircuitEnabled();
diff --git
a/hadoop-hdds/test-utils/src/main/java/org/apache/ozone/test/GenericTestUtils.java
b/hadoop-hdds/test-utils/src/main/java/org/apache/ozone/test/GenericTestUtils.java
index accb3595b85..2e6bc601ded 100644
---
a/hadoop-hdds/test-utils/src/main/java/org/apache/ozone/test/GenericTestUtils.java
+++
b/hadoop-hdds/test-utils/src/main/java/org/apache/ozone/test/GenericTestUtils.java
@@ -228,7 +228,7 @@ public static <K, V> Map<V, K> getReverseMap(Map<K,
List<V>> map) {
/**
* Class to capture logs for doing assertions.
*/
- public abstract static class LogCapturer {
+ public abstract static class LogCapturer implements AutoCloseable {
private final StringWriter sw = new StringWriter();
public static LogCapturer captureLogs(Logger logger) {
@@ -265,6 +265,11 @@ protected StringWriter writer() {
public void clearOutput() {
writer().getBuffer().setLength(0);
}
+
+ @Override
+ public void close() {
+ stopCapturing();
+ }
}
@Deprecated
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/InputStreamTests.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/InputStreamTests.java
index 9ef3153b5f3..c16ca259e05 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/InputStreamTests.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/InputStreamTests.java
@@ -22,6 +22,7 @@
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL;
import static
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL;
+import java.io.File;
import java.time.Duration;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@@ -31,19 +32,25 @@
import org.apache.hadoop.hdds.conf.StorageUnit;
import org.apache.hadoop.hdds.scm.OzoneClientConfig;
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
+import org.apache.hadoop.hdds.scm.XceiverClientGrpc;
+import org.apache.hadoop.hdds.scm.XceiverClientShortCircuit;
import
org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.scm.storage.BlockInputStream;
+import org.apache.hadoop.hdds.scm.storage.LocalChunkInputStream;
import org.apache.hadoop.hdds.utils.IOUtils;
import org.apache.hadoop.ozone.ClientConfigForTesting;
import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.OzoneConfigKeys;
import org.apache.hadoop.ozone.container.OzoneTestHelper;
import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion;
+import org.apache.ozone.test.GenericTestUtils;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.io.TempDir;
+import org.slf4j.event.Level;
-// TODO remove this class, set config as default in integration tests
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
abstract class InputStreamTests {
@@ -55,7 +62,10 @@ abstract class InputStreamTests {
private MiniOzoneCluster cluster;
- protected MiniOzoneCluster newCluster() throws Exception {
+ @TempDir
+ private File dir;
+
+ private MiniOzoneCluster newCluster() throws Exception {
OzoneConfiguration conf = new OzoneConfiguration();
OzoneClientConfig config = conf.getObject(OzoneClientConfig.class);
@@ -74,7 +84,6 @@ protected MiniOzoneCluster newCluster() throws Exception {
conf.getObject(ReplicationManagerConfiguration.class);
repConf.setInterval(Duration.ofSeconds(1));
conf.setFromObject(repConf);
- setCustomizedProperties(conf);
ClientConfigForTesting.newBuilder(StorageUnit.BYTES)
.setBlockSize(BLOCK_SIZE)
@@ -83,8 +92,13 @@ protected MiniOzoneCluster newCluster() throws Exception {
.setStreamBufferMaxSize(MAX_FLUSH_SIZE)
.applyTo(conf);
+ int datanodeCount = getDatanodeCount();
+ if (datanodeCount == 1) {
+ enableShortCircuitRead(dir, conf);
+ }
+
return MiniOzoneCluster.newBuilder(conf)
- .setNumDatanodes(getDatanodeCount())
+ .setNumDatanodes(datanodeCount)
.build();
}
@@ -96,9 +110,6 @@ int getDatanodeCount() {
return 5;
}
- void setCustomizedProperties(OzoneConfiguration configuration) {
- }
-
ReplicationConfig getRepConfig() {
return RatisReplicationConfig.getInstance(THREE);
}
@@ -135,4 +146,22 @@ private void closeContainers() {
}
});
}
+
+ protected static void enableShortCircuitRead(File dir, OzoneConfiguration
configuration) {
+ useShortCircuitRead(configuration, true);
+ configuration.set(OzoneClientConfig.OZONE_DOMAIN_SOCKET_PATH, new
File(dir, "ozone-socket").getAbsolutePath());
+ }
+
+ protected static void useShortCircuitRead(OzoneConfiguration conf, boolean
use) {
+ OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class);
+ clientConfig.setShortCircuit(use);
+ conf.setFromObject(clientConfig);
+ }
+
+ protected static void debugShortCircuitRead() {
+ GenericTestUtils.setLogLevel(XceiverClientShortCircuit.LOG, Level.DEBUG);
+ GenericTestUtils.setLogLevel(XceiverClientGrpc.LOG, Level.DEBUG);
+ GenericTestUtils.setLogLevel(LocalChunkInputStream.LOG, Level.DEBUG);
+ GenericTestUtils.setLogLevel(BlockInputStream.LOG, Level.DEBUG);
+ }
}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestChunkInputStream.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestChunkInputStream.java
index 4db70817f7e..1dff849d11d 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestChunkInputStream.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestChunkInputStream.java
@@ -17,44 +17,111 @@
package org.apache.hadoop.ozone.client.rpc.read;
+import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.ONE;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assumptions.assumeTrue;
import java.io.IOException;
import java.nio.ByteBuffer;
import org.apache.commons.io.IOUtils;
+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.scm.XceiverClientGrpc;
+import org.apache.hadoop.hdds.scm.XceiverClientShortCircuit;
import org.apache.hadoop.hdds.scm.storage.BlockInputStream;
import org.apache.hadoop.hdds.scm.storage.ChunkInputStream;
+import org.apache.hadoop.hdds.scm.storage.DomainSocketFactory;
+import org.apache.hadoop.hdds.scm.storage.LocalChunkInputStream;
import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.OzoneClientFactory;
import org.apache.hadoop.ozone.client.io.KeyInputStream;
import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion;
import org.apache.hadoop.ozone.container.keyvalue.ContainerLayoutTestInfo;
import org.apache.hadoop.ozone.om.BucketForTesting;
+import org.apache.ozone.test.GenericTestUtils.LogCapturer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.params.Parameter;
+import org.junit.jupiter.params.ParameterizedClass;
+import org.junit.jupiter.params.provider.ValueSource;
/**
* Tests {@link ChunkInputStream}.
*/
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@ParameterizedClass
+@ValueSource(booleans = {true, false})
class TestChunkInputStream extends InputStreamTests {
+ @Parameter
+ private boolean useShortCircuitRead;
+
+ private LogCapturer localChunkInputStreamLog;
+ private LogCapturer shortCircuitClientLog;
+ private LogCapturer blockInputStreamLog;
+ private LogCapturer grpcClientLog;
+
+ @BeforeEach
+ void startCapturing() {
+ localChunkInputStreamLog =
LogCapturer.captureLogs(LocalChunkInputStream.LOG);
+ shortCircuitClientLog =
LogCapturer.captureLogs(XceiverClientShortCircuit.LOG);
+ blockInputStreamLog = LogCapturer.captureLogs(BlockInputStream.LOG);
+ grpcClientLog = LogCapturer.captureLogs(XceiverClientGrpc.LOG);
+ }
+
+ @AfterEach
+ void stopCapturing() {
+ org.apache.hadoop.hdds.utils.IOUtils.closeQuietly(
+ blockInputStreamLog, grpcClientLog, localChunkInputStreamLog,
shortCircuitClientLog);
+ }
+
+ @Override
+ int getDatanodeCount() {
+ return getRepConfig().getRequiredNodes();
+ }
+
+ @Override
+ ReplicationConfig getRepConfig() {
+ return RatisReplicationConfig.getInstance(ONE);
+ }
+
/**
* Run the tests as a single test method to avoid needing a new mini-cluster
* for each test.
*/
@ContainerLayoutTestInfo.ContainerTest
void testAll(ContainerLayoutVersion layout) throws Exception {
- try (OzoneClient client = getCluster().newClient()) {
- updateConfig(layout);
+ if (useShortCircuitRead) {
+
assumeTrue(DomainSocketFactory.getInstance(getCluster().getConf()).isServiceReady());
+ }
+ updateConfig(layout);
+
+ OzoneConfiguration clientConfig = new
OzoneConfiguration(getCluster().getConf());
+ useShortCircuitRead(clientConfig, useShortCircuitRead);
+ debugShortCircuitRead();
+
+ try (OzoneClient client = OzoneClientFactory.getRpcClient(clientConfig)) {
BucketForTesting bucket = BucketForTesting.newBuilder(client).build();
testChunkReadBuffers(bucket);
testBufferRelease(bucket);
testCloseReleasesBuffers(bucket);
}
+
+ assertEquals(useShortCircuitRead, localChunkInputStreamLog.getOutput()
+ .contains("LocalChunkInputStream is created"));
+ assertEquals(useShortCircuitRead, shortCircuitClientLog.getOutput()
+ .contains("XceiverClientShortCircuit is created"));
+ assertEquals(useShortCircuitRead, blockInputStreamLog.getOutput()
+ .contains("Get the FileInputStream of block"));
+ assertEquals(!useShortCircuitRead, grpcClientLog.getOutput()
+ .contains("XceiverClientGrpc is created"));
}
/**
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestLocalChunkInputStream.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestLocalChunkInputStream.java
index ccf366beaaa..9ff0d5c165d 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestLocalChunkInputStream.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestLocalChunkInputStream.java
@@ -18,36 +18,32 @@
package org.apache.hadoop.ozone.client.rpc.read;
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.ONE;
+import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
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.assertNull;
-import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assumptions.assumeTrue;
-import java.io.File;
import java.io.IOException;
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.scm.OzoneClientConfig;
import org.apache.hadoop.hdds.scm.XceiverClientGrpc;
import org.apache.hadoop.hdds.scm.XceiverClientShortCircuit;
import org.apache.hadoop.hdds.scm.storage.BlockInputStream;
import org.apache.hadoop.hdds.scm.storage.DomainSocketFactory;
import org.apache.hadoop.hdds.scm.storage.LocalChunkInputStream;
+import org.apache.hadoop.hdds.utils.IOUtils;
import org.apache.hadoop.ozone.client.OzoneClient;
import org.apache.hadoop.ozone.client.io.KeyInputStream;
-import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion;
import
org.apache.hadoop.ozone.container.common.transport.server.XceiverServerSpi;
-import org.apache.hadoop.ozone.container.keyvalue.ContainerLayoutTestInfo;
import org.apache.hadoop.ozone.om.BucketForTesting;
-import org.apache.ozone.test.GenericTestUtils;
+import org.apache.ozone.test.GenericTestUtils.LogCapturer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance;
-import org.junit.jupiter.api.io.TempDir;
-import org.slf4j.event.Level;
/**
* Tests {@link LocalChunkInputStream}.
@@ -58,75 +54,46 @@
* to intellij run configuration.
*/
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
-public class TestLocalChunkInputStream extends TestChunkInputStream {
+class TestLocalChunkInputStream extends InputStreamTests {
- @TempDir
- private File dir;
+ private LogCapturer shortCircuitClientLog;
+ private LogCapturer grpcClientLog;
+ @BeforeAll
@Override
- int getDatanodeCount() {
- return 1;
+ void setup() throws Exception {
+ assumeTrue(DomainSocketFactory.isNativeLibraryLoaded());
+ super.setup();
+
assumeTrue(DomainSocketFactory.getInstance(getCluster().getConf()).isServiceReady());
}
- @Override
- void setCustomizedProperties(OzoneConfiguration configuration) {
- OzoneClientConfig clientConfig =
configuration.getObject(OzoneClientConfig.class);
- clientConfig.setShortCircuit(true);
- configuration.setFromObject(clientConfig);
- configuration.set(OzoneClientConfig.OZONE_DOMAIN_SOCKET_PATH,
- new File(dir, "ozone-socket").getAbsolutePath());
- GenericTestUtils.setLogLevel(XceiverClientShortCircuit.LOG, Level.DEBUG);
- GenericTestUtils.setLogLevel(XceiverClientGrpc.LOG, Level.DEBUG);
- GenericTestUtils.setLogLevel(LocalChunkInputStream.LOG, Level.DEBUG);
- GenericTestUtils.setLogLevel(BlockInputStream.LOG, Level.DEBUG);
+ @BeforeEach
+ void startCapturing() {
+ shortCircuitClientLog =
LogCapturer.captureLogs(XceiverClientShortCircuit.LOG);
+ grpcClientLog = LogCapturer.captureLogs(XceiverClientGrpc.LOG);
}
- @Override
- ReplicationConfig getRepConfig() {
- return RatisReplicationConfig.getInstance(ONE);
+ @AfterEach
+ void stopCapturing() {
+ IOUtils.closeQuietly(grpcClientLog, shortCircuitClientLog);
}
-
- /**
- * Run the tests as a single test method to avoid needing a new mini-cluster
- * for each test.
- */
- @ContainerLayoutTestInfo.ContainerTest
@Override
- void testAll(ContainerLayoutVersion layout) throws Exception {
- try (OzoneClient client = getCluster().newClient()) {
- updateConfig(layout);
-
assumeTrue(DomainSocketFactory.getInstance(getCluster().getConf()).isServiceReady());
+ int getDatanodeCount() {
+ return getRepConfig().getRequiredNodes();
+ }
- BucketForTesting bucket = BucketForTesting.newBuilder(client).build();
- GenericTestUtils.LogCapturer logCapturer1 =
- GenericTestUtils.LogCapturer.captureLogs(LocalChunkInputStream.LOG);
- GenericTestUtils.LogCapturer logCapturer2 =
-
GenericTestUtils.LogCapturer.captureLogs(XceiverClientShortCircuit.LOG);
- GenericTestUtils.LogCapturer logCapturer3 =
- GenericTestUtils.LogCapturer.captureLogs(BlockInputStream.LOG);
- GenericTestUtils.LogCapturer logCapturer4 =
- GenericTestUtils.LogCapturer.captureLogs(XceiverClientGrpc.LOG);
- testChunkReadBuffers(bucket);
- testBufferRelease(bucket);
- testCloseReleasesBuffers(bucket);
- assertTrue(logCapturer1.getOutput().contains("LocalChunkInputStream is
created"));
- assertTrue(logCapturer2.getOutput().contains("XceiverClientShortCircuit
is created"));
- assertTrue((logCapturer3.getOutput().contains("Get the FileInputStream
of block")));
- assertFalse(logCapturer4.getOutput().contains("XceiverClientGrpc is
created"));
- }
+ @Override
+ ReplicationConfig getRepConfig() {
+ return RatisReplicationConfig.getInstance(ONE);
}
@Test
void testFallbackToGrpc() throws Exception {
try (OzoneClient client = getCluster().newClient()) {
-
assumeTrue(DomainSocketFactory.getInstance(getCluster().getConf()).isServiceReady());
-
BucketForTesting bucket = BucketForTesting.newBuilder(client).build();
- GenericTestUtils.LogCapturer logCapturer1 =
-
GenericTestUtils.LogCapturer.captureLogs(XceiverClientShortCircuit.LOG);
- GenericTestUtils.LogCapturer logCapturer2 =
- GenericTestUtils.LogCapturer.captureLogs(XceiverClientGrpc.LOG);
+
+ debugShortCircuitRead();
// create key
String keyName = getNewKeyName();
@@ -137,7 +104,7 @@ void testFallbackToGrpc() throws Exception {
(BlockInputStream)keyInputStream.getPartStreams().get(0);
block0Stream.initialize();
assertNotNull(block0Stream.getBlockFileInputStream());
-
assertTrue(logCapturer1.getOutput().contains("XceiverClientShortCircuit is
created"));
+
assertThat(shortCircuitClientLog.getOutput()).contains("XceiverClientShortCircuit
is created");
// stop XceiverServerDomainSocket server before client sends the
second getBlockRequest to server
XceiverServerSpi server = getCluster().getHddsDatanodes().get(0)
@@ -147,8 +114,9 @@ void testFallbackToGrpc() throws Exception {
try {
block1Stream.initialize();
} catch (IOException e) {
- assertTrue(e.getMessage().contains("DomainSocket stream is not
open"));
- assertTrue(logCapturer1.getOutput().contains("ReceiveResponseTask is
closed due to java.io.EOFException"));
+ assertThat(e.getMessage()).contains("DomainSocket stream is not
open");
+ assertThat(shortCircuitClientLog.getOutput())
+ .contains("ReceiveResponseTask is closed due to
java.io.EOFException");
}
assertNull(block1Stream.getBlockFileInputStream());
// read whole key through Grpc channel
@@ -156,7 +124,7 @@ void testFallbackToGrpc() throws Exception {
int readLen = keyInputStream.read(data);
assertEquals(dataLength, readLen);
assertArrayEquals(inputData, data);
- assertTrue(logCapturer2.getOutput().contains("XceiverClientGrpc is
created"));
+ assertThat(grpcClientLog.getOutput()).contains("XceiverClientGrpc is
created");
}
}
}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
index 2d4ce5095aa..5f4c7f34be5 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
@@ -29,7 +29,6 @@
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
import org.apache.hadoop.hdds.scm.OzoneClientConfig;
import org.apache.hadoop.hdds.scm.storage.StreamBlockInputStream;
-import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.client.OzoneClient;
import org.apache.hadoop.ozone.client.OzoneClientFactory;
import org.apache.hadoop.ozone.client.io.KeyInputStream;
@@ -76,16 +75,12 @@ public class TestStreamBlockInputStream extends
InputStreamTests {
@Test
void testReadKey() throws Exception {
- try (MiniOzoneCluster cluster = newCluster()) {
- cluster.waitForClusterToBeReady();
- LOG.info("cluster ready");
- OzoneConfiguration conf = cluster.getConf();
-
- runTestReadKey(DATA_LENGTH, false, conf);
- for (int i = 0; i < 2; i++) {
- final int keyLength = DATA_LENGTH +
ThreadLocalRandom.current().nextInt(DATA_LENGTH);
- runTestReadKey(keyLength, true, conf);
- }
+ OzoneConfiguration conf = getCluster().getConf();
+
+ runTestReadKey(DATA_LENGTH, false, conf);
+ for (int i = 0; i < 2; i++) {
+ final int keyLength = DATA_LENGTH +
ThreadLocalRandom.current().nextInt(DATA_LENGTH);
+ runTestReadKey(keyLength, true, conf);
}
}
@@ -196,34 +191,30 @@ void testAllWithoutPreRead() throws Exception {
}
void runTestAll(boolean preRead) throws Exception {
- try (MiniOzoneCluster cluster = newCluster()) {
- cluster.waitForClusterToBeReady();
-
- OzoneConfiguration conf = cluster.getConf();
- OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class);
- clientConfig.setStreamReadBlock(true);
- if (!preRead) {
- clientConfig.setStreamReadPreReadSize(0);
- }
- OzoneConfiguration copy = new OzoneConfiguration(conf);
- copy.setFromObject(clientConfig);
- String keyName = getNewKeyName();
- try (OzoneClient client = OzoneClientFactory.getRpcClient(copy)) {
- bucket = BucketForTesting.newBuilder(client).build();
- inputData = bucket.writeRandomBytes(keyName, DATA_LENGTH);
- testReadKeyFully(keyName);
- testSeek(keyName);
- testReadEmptyBlock();
- }
- keyName = getNewKeyName();
- clientConfig.setChecksumType(ContainerProtos.ChecksumType.NONE);
- copy.setFromObject(clientConfig);
- try (OzoneClient client = OzoneClientFactory.getRpcClient(copy)) {
- bucket = BucketForTesting.newBuilder(client).build();
- inputData = bucket.writeRandomBytes(keyName, DATA_LENGTH);
- testReadKeyFully(keyName);
- testSeek(keyName);
- }
+ OzoneConfiguration conf = getCluster().getConf();
+ OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class);
+ clientConfig.setStreamReadBlock(true);
+ if (!preRead) {
+ clientConfig.setStreamReadPreReadSize(0);
+ }
+ OzoneConfiguration copy = new OzoneConfiguration(conf);
+ copy.setFromObject(clientConfig);
+ String keyName = getNewKeyName();
+ try (OzoneClient client = OzoneClientFactory.getRpcClient(copy)) {
+ bucket = BucketForTesting.newBuilder(client).build();
+ inputData = bucket.writeRandomBytes(keyName, DATA_LENGTH);
+ testReadKeyFully(keyName);
+ testSeek(keyName);
+ testReadEmptyBlock();
+ }
+ keyName = getNewKeyName();
+ clientConfig.setChecksumType(ContainerProtos.ChecksumType.NONE);
+ copy.setFromObject(clientConfig);
+ try (OzoneClient client = OzoneClientFactory.getRpcClient(copy)) {
+ bucket = BucketForTesting.newBuilder(client).build();
+ inputData = bucket.writeRandomBytes(keyName, DATA_LENGTH);
+ testReadKeyFully(keyName);
+ testSeek(keyName);
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]