This is an automated email from the ASF dual-hosted git repository.

roman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git

commit 5855f5354985cabea6f01b6b2effaaa8cfcbee55
Author: Roman Khachatryan <[email protected]>
AuthorDate: Wed May 29 21:46:11 2024 +0200

    [FLINK-35501] Replace inheritance with encapsulation for 
RocksDBStateDataTransfer*
---
 .../state/RocksDBKeyedStateBackendBuilder.java     |  4 +-
 .../streaming/state/RocksDBStateDataTransfer.java  | 47 ----------------
 .../state/RocksDBStateDataTransferHelper.java      | 63 ++++++++++++++++++++++
 .../streaming/state/RocksDBStateDownloader.java    | 22 ++++++--
 .../streaming/state/RocksDBStateUploader.java      | 20 +++++--
 .../RocksDBIncrementalRestoreOperation.java        |  4 +-
 6 files changed, 105 insertions(+), 55 deletions(-)

diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBKeyedStateBackendBuilder.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBKeyedStateBackendBuilder.java
index 1a63950305d..e0b1359060b 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBKeyedStateBackendBuilder.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBKeyedStateBackendBuilder.java
@@ -585,7 +585,9 @@ public class RocksDBKeyedStateBackendBuilder<K> extends 
AbstractKeyedStateBacken
         RocksDBSnapshotStrategyBase<K, ?> checkpointSnapshotStrategy;
         RocksDBStateUploader stateUploader =
                 injectRocksDBStateUploader == null
-                        ? new RocksDBStateUploader(numberOfTransferingThreads)
+                        ? new RocksDBStateUploader(
+                                RocksDBStateDataTransferHelper.forThreadNum(
+                                        numberOfTransferingThreads))
                         : injectRocksDBStateUploader;
         if (enableIncrementalCheckpointing) {
             checkpointSnapshotStrategy =
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDataTransfer.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDataTransfer.java
deleted file mode 100644
index 246332671ce..00000000000
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDataTransfer.java
+++ /dev/null
@@ -1,47 +0,0 @@
-/*
- * 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.flink.contrib.streaming.state;
-
-import org.apache.flink.util.concurrent.ExecutorThreadFactory;
-
-import java.io.Closeable;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
-
-import static 
org.apache.flink.util.concurrent.Executors.newDirectExecutorService;
-
-/** Data transfer base class for {@link RocksDBKeyedStateBackend}. */
-class RocksDBStateDataTransfer implements Closeable {
-
-    protected final ExecutorService executorService;
-
-    RocksDBStateDataTransfer(int threadNum) {
-        if (threadNum > 1) {
-            executorService =
-                    Executors.newFixedThreadPool(
-                            threadNum, new 
ExecutorThreadFactory("Flink-RocksDBStateDataTransfer"));
-        } else {
-            executorService = newDirectExecutorService();
-        }
-    }
-
-    @Override
-    public void close() {
-        executorService.shutdownNow();
-    }
-}
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDataTransferHelper.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDataTransferHelper.java
new file mode 100644
index 00000000000..70f31b52616
--- /dev/null
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDataTransferHelper.java
@@ -0,0 +1,63 @@
+/*
+ * 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.flink.contrib.streaming.state;
+
+import org.apache.flink.util.concurrent.ExecutorThreadFactory;
+
+import java.io.Closeable;
+import java.io.IOException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+import static org.apache.flink.util.Preconditions.checkNotNull;
+import static 
org.apache.flink.util.concurrent.Executors.newDirectExecutorService;
+
+/** Data transfer helper for {@link RocksDBKeyedStateBackend}. */
+public final class RocksDBStateDataTransferHelper implements Closeable {
+
+    private final ExecutorService executorService;
+    private final Closeable closeable;
+
+    public static RocksDBStateDataTransferHelper forThreadNum(int threadNum) {
+        ExecutorService executorService = getExecutorService(threadNum);
+        return new RocksDBStateDataTransferHelper(executorService, 
executorService::shutdownNow);
+    }
+
+    private static ExecutorService getExecutorService(int threadNum) {
+        if (threadNum > 1) {
+            return Executors.newFixedThreadPool(
+                    threadNum, new 
ExecutorThreadFactory("Flink-RocksDBStateDataTransferHelper"));
+        } else {
+            return newDirectExecutorService();
+        }
+    }
+
+    RocksDBStateDataTransferHelper(ExecutorService executorService, Closeable 
closeable) {
+        this.executorService = checkNotNull(executorService);
+        this.closeable = checkNotNull(closeable);
+    }
+
+    @Override
+    public void close() throws IOException {
+        closeable.close();
+    }
+
+    public ExecutorService getExecutorService() {
+        return executorService;
+    }
+}
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDownloader.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDownloader.java
index 850bda6bad7..6c319cb9358 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDownloader.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateDownloader.java
@@ -17,6 +17,7 @@
 
 package org.apache.flink.contrib.streaming.state;
 
+import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.core.fs.CloseableRegistry;
 import org.apache.flink.core.fs.FSDataInputStream;
 import org.apache.flink.runtime.state.StreamStateHandle;
@@ -27,6 +28,7 @@ import org.apache.flink.util.IOUtils;
 import org.apache.flink.util.concurrent.FutureUtils;
 import org.apache.flink.util.function.ThrowingRunnable;
 
+import java.io.Closeable;
 import java.io.IOException;
 import java.io.OutputStream;
 import java.nio.file.Files;
@@ -38,10 +40,16 @@ import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
 /** Help class for downloading RocksDB state files. */
-public class RocksDBStateDownloader extends RocksDBStateDataTransfer {
+public class RocksDBStateDownloader implements Closeable {
+    private final RocksDBStateDataTransferHelper transfer;
 
+    @VisibleForTesting
     public RocksDBStateDownloader(int restoringThreadNum) {
-        super(restoringThreadNum);
+        this(RocksDBStateDataTransferHelper.forThreadNum(restoringThreadNum));
+    }
+
+    public RocksDBStateDownloader(RocksDBStateDataTransferHelper transfer) {
+        this.transfer = transfer;
     }
 
     /**
@@ -116,7 +124,10 @@ public class RocksDBStateDownloader extends 
RocksDBStateDataTransfer {
                                                                             
remoteFileHandle,
                                                                             
closeableRegistry));
                                                 }))
-                .map(runnable -> CompletableFuture.runAsync(runnable, 
executorService));
+                .map(
+                        runnable ->
+                                CompletableFuture.runAsync(
+                                        runnable, 
transfer.getExecutorService()));
     }
 
     /** Copies the file from a single state handle to the given path. */
@@ -155,4 +166,9 @@ public class RocksDBStateDownloader extends 
RocksDBStateDataTransfer {
             throw new IOException(ex);
         }
     }
+
+    @Override
+    public void close() throws IOException {
+        this.transfer.close();
+    }
 }
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateUploader.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateUploader.java
index 983450fad75..bb89729cef9 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateUploader.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBStateUploader.java
@@ -18,6 +18,7 @@
 
 package org.apache.flink.contrib.streaming.state;
 
+import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.core.fs.CloseableRegistry;
 import org.apache.flink.runtime.state.CheckpointStateOutputStream;
 import org.apache.flink.runtime.state.CheckpointStreamFactory;
@@ -33,6 +34,7 @@ import org.apache.flink.util.function.CheckedSupplier;
 
 import javax.annotation.Nonnull;
 
+import java.io.Closeable;
 import java.io.IOException;
 import java.io.InputStream;
 import java.nio.file.Files;
@@ -44,11 +46,18 @@ import java.util.concurrent.ExecutionException;
 import java.util.stream.Collectors;
 
 /** Help class for uploading RocksDB state files. */
-public class RocksDBStateUploader extends RocksDBStateDataTransfer {
+public class RocksDBStateUploader implements Closeable {
     private static final int READ_BUFFER_SIZE = 16 * 1024;
 
+    private final RocksDBStateDataTransferHelper transfer;
+
+    @VisibleForTesting
     public RocksDBStateUploader(int numberOfSnapshottingThreads) {
-        super(numberOfSnapshottingThreads);
+        
this(RocksDBStateDataTransferHelper.forThreadNum(numberOfSnapshottingThreads));
+    }
+
+    public RocksDBStateUploader(RocksDBStateDataTransferHelper transfer) {
+        this.transfer = transfer;
     }
 
     /**
@@ -114,7 +123,7 @@ public class RocksDBStateUploader extends 
RocksDBStateDataTransfer {
                                                                 stateScope,
                                                                 
closeableRegistry,
                                                                 
tmpResourcesRegistry)),
-                                        executorService))
+                                        transfer.getExecutorService()))
                 .collect(Collectors.toList());
     }
 
@@ -170,4 +179,9 @@ public class RocksDBStateUploader extends 
RocksDBStateDataTransfer {
             }
         }
     }
+
+    @Override
+    public void close() throws IOException {
+        this.transfer.close();
+    }
 }
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/restore/RocksDBIncrementalRestoreOperation.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/restore/RocksDBIncrementalRestoreOperation.java
index b5276abc848..8ee45b372a8 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/restore/RocksDBIncrementalRestoreOperation.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/restore/RocksDBIncrementalRestoreOperation.java
@@ -24,6 +24,7 @@ import 
org.apache.flink.contrib.streaming.state.RocksDBIncrementalCheckpointUtil
 import 
org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackend.RocksDbKvStateInfo;
 import org.apache.flink.contrib.streaming.state.RocksDBNativeMetricOptions;
 import org.apache.flink.contrib.streaming.state.RocksDBOperationUtils;
+import org.apache.flink.contrib.streaming.state.RocksDBStateDataTransferHelper;
 import org.apache.flink.contrib.streaming.state.RocksDBStateDownloader;
 import org.apache.flink.contrib.streaming.state.RocksDBWriteBatchWrapper;
 import org.apache.flink.contrib.streaming.state.RocksIteratorWrapper;
@@ -721,7 +722,8 @@ public class RocksDBIncrementalRestoreOperation<K> 
implements RocksDBRestoreOper
                 operatorIdentifier,
                 keyGroupRange.prettyPrintInterval());
         try (RocksDBStateDownloader rocksDBStateDownloader =
-                new RocksDBStateDownloader(numberOfTransferringThreads)) {
+                new RocksDBStateDownloader(
+                        
RocksDBStateDataTransferHelper.forThreadNum(numberOfTransferringThreads))) {
             rocksDBStateDownloader.transferAllStateDataToDirectory(
                     downloadSpecs, cancelStreamRegistry);
             logger.info(

Reply via email to