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]

Reply via email to