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 9708f9fd65751296b3b0377c964207b630958259
Author: Roman Khachatryan <[email protected]>
AuthorDate: Wed May 29 22:30:50 2024 +0200

    [FLINK-35501] Use common IO thread pool for RocksDB data transfer
    
    Currently, each RocksDB state backend creates an executor backed by a
    thread pool.
    
    This makes it difficult to control the total number of threads per TM
    because it might have at least one task per slot and theoretically, many
    state backends per task (because of chaining).
    
    Additionally, using a common thread pool allows to indirectly control
    the load on the underlying DFS (e.g. the total number of requests to S3
    from a TM).
    
    This change allows to control the total number of data transfers per TM.
---
 docs/layouts/shortcodes/generated/expert_rocksdb_section.html    | 2 +-
 docs/layouts/shortcodes/generated/rocksdb_configuration.html     | 2 +-
 .../contrib/streaming/state/RocksDBKeyedStateBackendBuilder.java | 7 ++++---
 .../org/apache/flink/contrib/streaming/state/RocksDBOptions.java | 6 +++++-
 .../contrib/streaming/state/RocksDBStateDataTransferHelper.java  | 9 +++++++++
 .../state/restore/RocksDBIncrementalRestoreOperation.java        | 9 +++++++--
 .../streaming/state/snapshot/RocksDBSnapshotStrategyBase.java    | 2 +-
 .../state/snapshot/RocksIncrementalSnapshotStrategy.java         | 3 ++-
 .../state/snapshot/RocksNativeFullSnapshotStrategy.java          | 3 ++-
 .../contrib/streaming/state/EmbeddedRocksDBStateBackendTest.java | 2 +-
 10 files changed, 33 insertions(+), 12 deletions(-)

diff --git a/docs/layouts/shortcodes/generated/expert_rocksdb_section.html 
b/docs/layouts/shortcodes/generated/expert_rocksdb_section.html
index 47e4edce3ce..eac1d574e26 100644
--- a/docs/layouts/shortcodes/generated/expert_rocksdb_section.html
+++ b/docs/layouts/shortcodes/generated/expert_rocksdb_section.html
@@ -12,7 +12,7 @@
             
<td><h5>state.backend.rocksdb.checkpoint.transfer.thread.num</h5></td>
             <td style="word-wrap: break-word;">4</td>
             <td>Integer</td>
-            <td>The number of threads (per stateful operator) used to transfer 
(download and upload) files in RocksDBStateBackend.</td>
+            <td>The number of threads (per stateful operator) used to transfer 
(download and upload) files in RocksDBStateBackend.If negative, the common (TM) 
IO thread pool is used (see cluster.io-pool.size)</td>
         </tr>
         <tr>
             <td><h5>state.backend.rocksdb.localdir</h5></td>
diff --git a/docs/layouts/shortcodes/generated/rocksdb_configuration.html 
b/docs/layouts/shortcodes/generated/rocksdb_configuration.html
index b6ad67234e0..e9cca6e415a 100644
--- a/docs/layouts/shortcodes/generated/rocksdb_configuration.html
+++ b/docs/layouts/shortcodes/generated/rocksdb_configuration.html
@@ -12,7 +12,7 @@
             
<td><h5>state.backend.rocksdb.checkpoint.transfer.thread.num</h5></td>
             <td style="word-wrap: break-word;">4</td>
             <td>Integer</td>
-            <td>The number of threads (per stateful operator) used to transfer 
(download and upload) files in RocksDBStateBackend.</td>
+            <td>The number of threads (per stateful operator) used to transfer 
(download and upload) files in RocksDBStateBackend.If negative, the common (TM) 
IO thread pool is used (see cluster.io-pool.size)</td>
         </tr>
         <tr>
             <td><h5>state.backend.rocksdb.localdir</h5></td>
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 e0b1359060b..bcdb32f4e79 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
@@ -536,7 +536,8 @@ public class RocksDBKeyedStateBackendBuilder<K> extends 
AbstractKeyedStateBacken
                     overlapFractionThreshold,
                     useIngestDbRestoreMode,
                     incrementalRestoreAsyncCompactAfterRescale,
-                    rescalingUseDeleteFilesInRange);
+                    rescalingUseDeleteFilesInRange,
+                    ioExecutor);
         } else if (priorityQueueConfig.getPriorityQueueStateType()
                 == EmbeddedRocksDBStateBackend.PriorityQueueStateType.HEAP) {
             return new RocksDBHeapTimersFullRestoreOperation<>(
@@ -586,8 +587,8 @@ public class RocksDBKeyedStateBackendBuilder<K> extends 
AbstractKeyedStateBacken
         RocksDBStateUploader stateUploader =
                 injectRocksDBStateUploader == null
                         ? new RocksDBStateUploader(
-                                RocksDBStateDataTransferHelper.forThreadNum(
-                                        numberOfTransferingThreads))
+                                
RocksDBStateDataTransferHelper.forThreadNumIfSpecified(
+                                        numberOfTransferingThreads, 
ioExecutor))
                         : injectRocksDBStateUploader;
         if (enableIncrementalCheckpointing) {
             checkpointSnapshotStrategy =
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBOptions.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBOptions.java
index e0fd2303939..e6e0c51a82c 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBOptions.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/RocksDBOptions.java
@@ -27,6 +27,7 @@ import org.apache.flink.configuration.MemorySize;
 import org.apache.flink.configuration.description.Description;
 import org.apache.flink.configuration.description.TextElement;
 
+import static 
org.apache.flink.configuration.ClusterOptions.CLUSTER_IO_EXECUTOR_POOL_SIZE;
 import static 
org.apache.flink.contrib.streaming.state.EmbeddedRocksDBStateBackend.PriorityQueueStateType.ROCKSDB;
 import static 
org.apache.flink.contrib.streaming.state.PredefinedOptions.DEFAULT;
 import static 
org.apache.flink.contrib.streaming.state.PredefinedOptions.FLASH_SSD_OPTIMIZED;
@@ -86,7 +87,10 @@ public class RocksDBOptions {
                     .intType()
                     .defaultValue(4)
                     .withDescription(
-                            "The number of threads (per stateful operator) 
used to transfer (download and upload) files in RocksDBStateBackend.");
+                            "The number of threads (per stateful operator) 
used to transfer (download and upload) files in RocksDBStateBackend."
+                                    + "If negative, the common (TM) IO thread 
pool is used (see "
+                                    + CLUSTER_IO_EXECUTOR_POOL_SIZE.key()
+                                    + ")");
 
     /** The predefined settings for RocksDB DBOptions and ColumnFamilyOptions 
by Flink community. */
     @Documentation.Section(Documentation.Sections.EXPERT_ROCKSDB)
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
index 70f31b52616..9e65b8bba91 100644
--- 
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
@@ -38,6 +38,15 @@ public final class RocksDBStateDataTransferHelper implements 
Closeable {
         return new RocksDBStateDataTransferHelper(executorService, 
executorService::shutdownNow);
     }
 
+    public static RocksDBStateDataTransferHelper forThreadNumIfSpecified(
+            int threadNum, ExecutorService executorService) {
+        if (threadNum >= 0) {
+            return forThreadNum(threadNum);
+        } else {
+            return new RocksDBStateDataTransferHelper(executorService, () -> 
{});
+        }
+    }
+
     private static ExecutorService getExecutorService(int threadNum) {
         if (threadNum > 1) {
             return Executors.newFixedThreadPool(
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 8ee45b372a8..2a9166d0ea0 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
@@ -129,6 +129,8 @@ public class RocksDBIncrementalRestoreOperation<K> 
implements RocksDBRestoreOper
 
     private final boolean useDeleteFilesInRange;
 
+    private final ExecutorService ioExecutor;
+
     public RocksDBIncrementalRestoreOperation(
             String operatorIdentifier,
             KeyGroupRange keyGroupRange,
@@ -152,7 +154,8 @@ public class RocksDBIncrementalRestoreOperation<K> 
implements RocksDBRestoreOper
             double overlapFractionThreshold,
             boolean useIngestDbRestoreMode,
             boolean asyncCompactAfterRescale,
-            boolean useDeleteFilesInRange) {
+            boolean useDeleteFilesInRange,
+            ExecutorService ioExecutor) {
         this.rocksHandle =
                 new RocksDBHandle(
                         kvStateInformation,
@@ -183,6 +186,7 @@ public class RocksDBIncrementalRestoreOperation<K> 
implements RocksDBRestoreOper
         this.useIngestDbRestoreMode = false;
         this.asyncCompactAfterRescale = false;
         this.useDeleteFilesInRange = useDeleteFilesInRange;
+        this.ioExecutor = ioExecutor;
     }
 
     /**
@@ -723,7 +727,8 @@ public class RocksDBIncrementalRestoreOperation<K> 
implements RocksDBRestoreOper
                 keyGroupRange.prettyPrintInterval());
         try (RocksDBStateDownloader rocksDBStateDownloader =
                 new RocksDBStateDownloader(
-                        
RocksDBStateDataTransferHelper.forThreadNum(numberOfTransferringThreads))) {
+                        RocksDBStateDataTransferHelper.forThreadNumIfSpecified(
+                                numberOfTransferringThreads, ioExecutor))) {
             rocksDBStateDownloader.transferAllStateDataToDirectory(
                     downloadSpecs, cancelStreamRegistry);
             logger.info(
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksDBSnapshotStrategyBase.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksDBSnapshotStrategyBase.java
index 6c3948a6303..f6bb541cf0d 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksDBSnapshotStrategyBase.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksDBSnapshotStrategyBase.java
@@ -303,7 +303,7 @@ public abstract class RocksDBSnapshotStrategyBase<K, R 
extends SnapshotResources
     }
 
     @Override
-    public abstract void close();
+    public abstract void close() throws IOException;
 
     /** Common operation in native rocksdb snapshot result supplier. */
     protected abstract class RocksDBSnapshotOperation
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksIncrementalSnapshotStrategy.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksIncrementalSnapshotStrategy.java
index 48d8bb35aa2..e29518cdd4b 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksIncrementalSnapshotStrategy.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksIncrementalSnapshotStrategy.java
@@ -47,6 +47,7 @@ import javax.annotation.Nonnegative;
 import javax.annotation.Nonnull;
 
 import java.io.File;
+import java.io.IOException;
 import java.nio.file.Path;
 import java.util.ArrayList;
 import java.util.Collection;
@@ -186,7 +187,7 @@ public class RocksIncrementalSnapshotStrategy<K>
     }
 
     @Override
-    public void close() {
+    public void close() throws IOException {
         stateUploader.close();
     }
 
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksNativeFullSnapshotStrategy.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksNativeFullSnapshotStrategy.java
index 3fc45862dc8..f0b2065635d 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksNativeFullSnapshotStrategy.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/main/java/org/apache/flink/contrib/streaming/state/snapshot/RocksNativeFullSnapshotStrategy.java
@@ -44,6 +44,7 @@ import javax.annotation.Nonnegative;
 import javax.annotation.Nonnull;
 
 import java.io.File;
+import java.io.IOException;
 import java.nio.file.Path;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -135,7 +136,7 @@ public class RocksNativeFullSnapshotStrategy<K>
     }
 
     @Override
-    public void close() {
+    public void close() throws IOException {
         stateUploader.close();
     }
 
diff --git 
a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/contrib/streaming/state/EmbeddedRocksDBStateBackendTest.java
 
b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/contrib/streaming/state/EmbeddedRocksDBStateBackendTest.java
index d58f388ab44..a8277fa07fc 100644
--- 
a/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/contrib/streaming/state/EmbeddedRocksDBStateBackendTest.java
+++ 
b/flink-state-backends/flink-statebackend-rocksdb/src/test/java/org/apache/flink/contrib/streaming/state/EmbeddedRocksDBStateBackendTest.java
@@ -785,7 +785,7 @@ public class EmbeddedRocksDBStateBackendTest
         assertThat(keyedStateBackend.isDisposed()).isTrue();
     }
 
-    private void verifyRocksDBStateUploaderClosed() {
+    private void verifyRocksDBStateUploaderClosed() throws IOException {
         if (enableIncrementalCheckpointing) {
             verify(rocksDBStateUploader, times(1)).close();
         }

Reply via email to