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(); }
