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 f77e61ec87 HDDS-8119. Remove loosely related AutoCloseable from 
SendContainerOutputStream (#4368)
f77e61ec87 is described below

commit f77e61ec8789ab0659eceb34f850c5188a53b10f
Author: Doroszlai, Attila <[email protected]>
AuthorDate: Thu Mar 9 21:45:43 2023 +0100

    HDDS-8119. Remove loosely related AutoCloseable from 
SendContainerOutputStream (#4368)
---
 .../replication/GrpcContainerUploader.java          | 16 ++++++++++++----
 .../replication/SendContainerOutputStream.java      | 21 +--------------------
 .../replication/TestSendContainerOutputStream.java  |  8 ++------
 3 files changed, 15 insertions(+), 30 deletions(-)

diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java
index ae8b5764d3..a3f1cdee51 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/GrpcContainerUploader.java
@@ -55,14 +55,22 @@ public class GrpcContainerUploader implements 
ContainerUploader {
   public OutputStream startUpload(long containerId, DatanodeDetails target,
       CompletableFuture<Void> callback, CopyContainerCompression compression)
       throws IOException {
-    GrpcReplicationClient client = null;
+    GrpcReplicationClient client = createReplicationClient(target, 
compression);
     try {
-      client = createReplicationClient(target, compression);
       StreamObserver<SendContainerRequest> requestStream = client.upload(
           new SendContainerResponseStreamObserver(containerId, target,
               callback));
-      return new SendContainerOutputStream(client, requestStream, containerId,
-          GrpcReplicationService.BUFFER_SIZE, compression);
+      return new SendContainerOutputStream(requestStream, containerId,
+          GrpcReplicationService.BUFFER_SIZE, compression) {
+        @Override
+        public void close() throws IOException {
+          try {
+            super.close();
+          } finally {
+            IOUtils.close(LOG, client);
+          }
+        }
+      };
     } catch (Exception e) {
       IOUtils.close(LOG, client);
       throw e;
diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java
index 5de63958fe..b24dcc10f5 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/SendContainerOutputStream.java
@@ -18,31 +18,21 @@
 package org.apache.hadoop.ozone.container.replication;
 
 import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerRequest;
-import org.apache.hadoop.hdds.utils.IOUtils;
 import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
 import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.io.IOException;
 
 /**
  * Output stream adapter for SendContainerResponse.
  */
 class SendContainerOutputStream extends GrpcOutputStream<SendContainerRequest> 
{
 
-  private static final Logger LOG =
-      LoggerFactory.getLogger(SendContainerOutputStream.class);
-
   private final CopyContainerCompression compression;
-  private final AutoCloseable client;
 
   SendContainerOutputStream(
-      AutoCloseable client, StreamObserver<SendContainerRequest> 
streamObserver,
+      StreamObserver<SendContainerRequest> streamObserver,
       long containerId, int bufferSize, CopyContainerCompression compression) {
     super(streamObserver, containerId, bufferSize);
     this.compression = compression;
-    this.client = client;
   }
 
   @Override
@@ -55,13 +45,4 @@ class SendContainerOutputStream extends 
GrpcOutputStream<SendContainerRequest> {
         .build();
     getStreamObserver().onNext(request);
   }
-
-  @Override
-  public void close() throws IOException {
-    try {
-      super.close();
-    } finally {
-      IOUtils.close(LOG, client);
-    }
-  }
 }
diff --git 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java
 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java
index cfd5552d56..747a4d2935 100644
--- 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java
+++ 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestSendContainerOutputStream.java
@@ -21,7 +21,6 @@ import 
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContai
 import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.EnumSource;
-import org.mockito.Mock;
 
 import java.io.OutputStream;
 
@@ -35,23 +34,20 @@ import static org.mockito.Mockito.verify;
 class TestSendContainerOutputStream
     extends GrpcOutputStreamTest<SendContainerRequest> {
 
-  @Mock
-  private AutoCloseable client;
-
   TestSendContainerOutputStream() {
     super(SendContainerRequest.class);
   }
 
   @Override
   protected OutputStream createSubject() {
-    return new SendContainerOutputStream(client, getObserver(),
+    return new SendContainerOutputStream(getObserver(),
         getContainerId(), getBufferSize(), NO_COMPRESSION);
   }
 
   @ParameterizedTest
   @EnumSource
   void usesCompression(CopyContainerCompression compression) throws Exception {
-    OutputStream subject = new SendContainerOutputStream(client,
+    OutputStream subject = new SendContainerOutputStream(
         getObserver(), getContainerId(), getBufferSize(), compression);
 
     byte[] bytes = getRandomBytes(16);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to