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(
