From 4db5cdc1d385cda466d5e9a99ef40a7d3b17b0b5 Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Wed, 8 Jul 2026 14:45:27 +0800 Subject: [PATCH 1/5] [flink] respect flink io tmp multi dirs to fix disk exhaustion --- .../fluss/flink/source/FlinkSource.java | 3 +- .../flink/tiering/source/TieringSource.java | 5 ++- .../utils/FlinkConnectorOptionsUtils.java | 36 ++++++++++++++++--- .../utils/FlinkConnectorOptionsUtilTest.java | 36 ++++++++++++++++++- 4 files changed, 72 insertions(+), 8 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java index 227db3fabc5..e65f4ff2281 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java @@ -290,7 +290,8 @@ public SourceReader createReader(SourceReaderContext conte flussConf.set( CLIENT_SCANNER_IO_TMP_DIR, - getClientScannerIoTmpDir(flussConf, context.getConfiguration())); + getClientScannerIoTmpDir( + flussConf, context.getConfiguration(), context.getIndexOfSubtask())); deserializationSchema.open( new DeserializerInitContextImpl( context.metricGroup().addGroup("deserializer"), diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java index 9251e2bfd06..61f4d8441f1 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java @@ -114,7 +114,10 @@ public SourceReader, TieringSplit> createRea elementsQueue = new FutureCompletingBlockingQueue<>(); flussConf.set( CLIENT_SCANNER_IO_TMP_DIR, - getClientScannerIoTmpDir(flussConf, sourceReaderContext.getConfiguration())); + getClientScannerIoTmpDir( + flussConf, + sourceReaderContext.getConfiguration(), + sourceReaderContext.getIndexOfSubtask())); Connection connection = ConnectionFactory.createConnection(flussConf); return new TieringSourceReader<>( elementsQueue, sourceReaderContext, connection, lakeTieringFactory); diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java index 249e1be0892..8a2a913a73a 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java @@ -29,6 +29,7 @@ import org.apache.flink.table.api.config.TableConfigOptions; import org.apache.flink.table.types.logical.RowType; +import javax.annotation.Nonnull; import javax.annotation.Nullable; import java.io.File; @@ -201,13 +202,38 @@ public static long parseTimestamp(String timestampStr, String optionKey, ZoneId public static String getClientScannerIoTmpDir( Configuration flussConf, org.apache.flink.configuration.Configuration flinkConfig) { - if (!flussConf.contains(CLIENT_SCANNER_IO_TMP_DIR)) { - if (flinkConfig.contains(TMP_DIRS)) { - // pass flink io tmp dir to fluss client. - return new File(flinkConfig.get(CoreOptions.TMP_DIRS), "/fluss").getAbsolutePath(); + return getClientScannerIoTmpDir(flussConf, flinkConfig, 0); + } + + public static String getClientScannerIoTmpDir( + Configuration flussConf, + org.apache.flink.configuration.Configuration flinkConfig, + int taskIndex) { + return flussConf + .getOptional(CLIENT_SCANNER_IO_TMP_DIR) + .orElseGet( + () -> { + String[] flinkTmpDirs = getFlinkIoTmpDirs(flinkConfig); + int idx = taskIndex % flinkTmpDirs.length; + return new File(flinkTmpDirs[idx], "/fluss").getAbsolutePath(); + }); + } + + private static String[] getFlinkIoTmpDirs( + org.apache.flink.configuration.Configuration flinkConfig) { + if (flinkConfig.contains(TMP_DIRS)) { + String[] paths = splitPaths(flinkConfig.get(CoreOptions.TMP_DIRS)); + if (paths.length > 0) { + return paths; } } - return flussConf.getString(CLIENT_SCANNER_IO_TMP_DIR); + return new String[] {System.getProperty("java.io.tmpdir")}; + } + + private static String[] splitPaths(@Nonnull String separatedPaths) { + return separatedPaths.length() > 0 + ? separatedPaths.split(",|" + File.pathSeparator) + : new String[0]; } /** Fluss startup options. * */ diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java index 0e06dc10196..e0502f0b5ec 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java @@ -22,6 +22,7 @@ import org.apache.flink.table.api.ValidationException; import org.junit.jupiter.api.Test; +import java.io.File; import java.time.ZoneId; import java.util.TimeZone; @@ -96,7 +97,7 @@ void testGetClientScannerIoTmpDir() { FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( new Configuration(), new org.apache.flink.configuration.Configuration())) - .isEqualTo(property + "/fluss"); + .isEqualTo(new File(property, "fluss").getAbsolutePath()); // only replace when flussConfig not contains CLIENT_SCANNER_IO_TMP_DIR while flinkConfig // contains TMP_DIRS. @@ -104,5 +105,38 @@ void testGetClientScannerIoTmpDir() { FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( new Configuration(), flinkConfig)) .isEqualTo("/flink_tmp_dir/fluss"); + assertThat(FlinkConnectorOptionsUtils.getClientScannerIoTmpDir(flussConfig, flinkConfig, 1)) + .isEqualTo("/fluss_tmp_dir"); + } + + @Test + void testGetClientScannerIoTmpDirWithMultipleFlinkTmpDirs() { + org.apache.flink.configuration.Configuration commaSeparatedFlinkConfig = + new org.apache.flink.configuration.Configuration() + .set(TMP_DIRS, "/flink_tmp_dir_0,/flink_tmp_dir_1"); + + assertThat( + FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( + new Configuration(), commaSeparatedFlinkConfig, 0)) + .isEqualTo(new File("/flink_tmp_dir_0", "fluss").getAbsolutePath()); + assertThat( + FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( + new Configuration(), commaSeparatedFlinkConfig, 1)) + .isEqualTo(new File("/flink_tmp_dir_1", "fluss").getAbsolutePath()); + assertThat( + FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( + new Configuration(), commaSeparatedFlinkConfig, 2)) + .isEqualTo(new File("/flink_tmp_dir_0", "fluss").getAbsolutePath()); + + org.apache.flink.configuration.Configuration pathSeparatorFlinkConfig = + new org.apache.flink.configuration.Configuration() + .set( + TMP_DIRS, + "/flink_tmp_dir_0" + File.pathSeparator + "/flink_tmp_dir_1"); + + assertThat( + FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( + new Configuration(), pathSeparatorFlinkConfig, 1)) + .isEqualTo(new File("/flink_tmp_dir_1", "fluss").getAbsolutePath()); } } From 34742ba3ff36d081815d72597fd0646478e962eb Mon Sep 17 00:00:00 2001 From: Junfan Zhang Date: Wed, 8 Jul 2026 15:42:17 +0800 Subject: [PATCH 2/5] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .../apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java index 8a2a913a73a..b2330fea6f0 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java @@ -214,8 +214,8 @@ public static String getClientScannerIoTmpDir( .orElseGet( () -> { String[] flinkTmpDirs = getFlinkIoTmpDirs(flinkConfig); - int idx = taskIndex % flinkTmpDirs.length; - return new File(flinkTmpDirs[idx], "/fluss").getAbsolutePath(); + int idx = Math.floorMod(taskIndex, flinkTmpDirs.length); + return new File(flinkTmpDirs[idx], "fluss").getAbsolutePath(); }); } From 1b069c538931327d999097d9c8e9b842063c5ad1 Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Wed, 8 Jul 2026 15:46:50 +0800 Subject: [PATCH 3/5] use flink internal built-in path spliter --- .../flink/utils/FlinkConnectorOptionsUtils.java | 17 +++-------------- 1 file changed, 3 insertions(+), 14 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java index b2330fea6f0..a9e332e0588 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java @@ -23,7 +23,6 @@ import org.apache.fluss.flink.sink.shuffle.DistributionMode; import org.apache.fluss.metadata.MergeEngineType; -import org.apache.flink.configuration.CoreOptions; import org.apache.flink.configuration.ReadableConfig; import org.apache.flink.table.api.ValidationException; import org.apache.flink.table.api.config.TableConfigOptions; @@ -42,7 +41,6 @@ import java.util.Optional; import java.util.stream.Collectors; -import static org.apache.flink.configuration.CoreOptions.TMP_DIRS; import static org.apache.fluss.config.ConfigOptions.CLIENT_SCANNER_IO_TMP_DIR; import static org.apache.fluss.flink.FlinkConnectorOptions.SCAN_SPLIT_ASSIGNMENT_BATCH_SIZE; import static org.apache.fluss.flink.FlinkConnectorOptions.SCAN_STARTUP_MODE; @@ -213,23 +211,14 @@ public static String getClientScannerIoTmpDir( .getOptional(CLIENT_SCANNER_IO_TMP_DIR) .orElseGet( () -> { - String[] flinkTmpDirs = getFlinkIoTmpDirs(flinkConfig); + String[] flinkTmpDirs = + org.apache.flink.configuration.ConfigurationUtils + .parseTempDirectories(flinkConfig); int idx = Math.floorMod(taskIndex, flinkTmpDirs.length); return new File(flinkTmpDirs[idx], "fluss").getAbsolutePath(); }); } - private static String[] getFlinkIoTmpDirs( - org.apache.flink.configuration.Configuration flinkConfig) { - if (flinkConfig.contains(TMP_DIRS)) { - String[] paths = splitPaths(flinkConfig.get(CoreOptions.TMP_DIRS)); - if (paths.length > 0) { - return paths; - } - } - return new String[] {System.getProperty("java.io.tmpdir")}; - } - private static String[] splitPaths(@Nonnull String separatedPaths) { return separatedPaths.length() > 0 ? separatedPaths.split(",|" + File.pathSeparator) From f97b9c4c0fac6657f476b0b353a94f719b91d46c Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Wed, 8 Jul 2026 15:48:43 +0800 Subject: [PATCH 4/5] remove --- .../fluss/flink/utils/FlinkConnectorOptionsUtils.java | 7 ------- 1 file changed, 7 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java index a9e332e0588..3873821e56c 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java @@ -28,7 +28,6 @@ import org.apache.flink.table.api.config.TableConfigOptions; import org.apache.flink.table.types.logical.RowType; -import javax.annotation.Nonnull; import javax.annotation.Nullable; import java.io.File; @@ -219,12 +218,6 @@ public static String getClientScannerIoTmpDir( }); } - private static String[] splitPaths(@Nonnull String separatedPaths) { - return separatedPaths.length() > 0 - ? separatedPaths.split(",|" + File.pathSeparator) - : new String[0]; - } - /** Fluss startup options. * */ public static class StartupOptions { public ScanStartupMode startupMode; From 67e3b0dde54d08df83641e53c8e89bb1caa92e90 Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Wed, 22 Jul 2026 14:59:47 +0800 Subject: [PATCH 5/5] address comment by yuxia --- .../fluss/flink/source/FlinkSource.java | 8 ++--- .../flink/tiering/source/TieringSource.java | 10 +++--- .../utils/FlinkConnectorOptionsUtils.java | 21 ++++-------- .../utils/FlinkConnectorOptionsUtilTest.java | 34 +++++++++---------- 4 files changed, 32 insertions(+), 41 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java index e65f4ff2281..b6cc0049c06 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/FlinkSource.java @@ -288,10 +288,10 @@ public SourceReader createReader(SourceReaderContext conte FlinkSourceReaderMetrics flinkSourceReaderMetrics = new FlinkSourceReaderMetrics(context.metricGroup()); - flussConf.set( + Configuration readerConf = new Configuration(flussConf); + readerConf.set( CLIENT_SCANNER_IO_TMP_DIR, - getClientScannerIoTmpDir( - flussConf, context.getConfiguration(), context.getIndexOfSubtask())); + getClientScannerIoTmpDir(readerConf, context.getConfiguration())); deserializationSchema.open( new DeserializerInitContextImpl( context.metricGroup().addGroup("deserializer"), @@ -301,7 +301,7 @@ public SourceReader createReader(SourceReaderContext conte return new FlinkSourceReader<>( elementsQueue, - flussConf, + readerConf, tablePath, sourceOutputType, context, diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java index 61f4d8441f1..9f616ca1e36 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/source/TieringSource.java @@ -112,13 +112,11 @@ public SourceReader, TieringSplit> createRea SourceReaderContext sourceReaderContext) { FutureCompletingBlockingQueue>> elementsQueue = new FutureCompletingBlockingQueue<>(); - flussConf.set( + Configuration readerConf = new Configuration(flussConf); + readerConf.set( CLIENT_SCANNER_IO_TMP_DIR, - getClientScannerIoTmpDir( - flussConf, - sourceReaderContext.getConfiguration(), - sourceReaderContext.getIndexOfSubtask())); - Connection connection = ConnectionFactory.createConnection(flussConf); + getClientScannerIoTmpDir(readerConf, sourceReaderContext.getConfiguration())); + Connection connection = ConnectionFactory.createConnection(readerConf); return new TieringSourceReader<>( elementsQueue, sourceReaderContext, connection, lakeTieringFactory); } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java index 3873821e56c..cb14c1abe48 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java @@ -23,6 +23,7 @@ import org.apache.fluss.flink.sink.shuffle.DistributionMode; import org.apache.fluss.metadata.MergeEngineType; +import org.apache.flink.configuration.ConfigurationUtils; import org.apache.flink.configuration.ReadableConfig; import org.apache.flink.table.api.ValidationException; import org.apache.flink.table.api.config.TableConfigOptions; @@ -199,23 +200,15 @@ public static long parseTimestamp(String timestampStr, String optionKey, ZoneId public static String getClientScannerIoTmpDir( Configuration flussConf, org.apache.flink.configuration.Configuration flinkConfig) { - return getClientScannerIoTmpDir(flussConf, flinkConfig, 0); - } - - public static String getClientScannerIoTmpDir( - Configuration flussConf, - org.apache.flink.configuration.Configuration flinkConfig, - int taskIndex) { return flussConf .getOptional(CLIENT_SCANNER_IO_TMP_DIR) .orElseGet( - () -> { - String[] flinkTmpDirs = - org.apache.flink.configuration.ConfigurationUtils - .parseTempDirectories(flinkConfig); - int idx = Math.floorMod(taskIndex, flinkTmpDirs.length); - return new File(flinkTmpDirs[idx], "fluss").getAbsolutePath(); - }); + () -> + new File( + ConfigurationUtils.getRandomTempDirectory( + flinkConfig), + "fluss") + .getAbsolutePath()); } /** Fluss startup options. * */ diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java index e0502f0b5ec..b2f12ba9093 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java @@ -24,6 +24,8 @@ import java.io.File; import java.time.ZoneId; +import java.util.HashSet; +import java.util.Set; import java.util.TimeZone; import static org.apache.flink.configuration.CoreOptions.TMP_DIRS; @@ -105,28 +107,24 @@ void testGetClientScannerIoTmpDir() { FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( new Configuration(), flinkConfig)) .isEqualTo("/flink_tmp_dir/fluss"); - assertThat(FlinkConnectorOptionsUtils.getClientScannerIoTmpDir(flussConfig, flinkConfig, 1)) - .isEqualTo("/fluss_tmp_dir"); } @Test - void testGetClientScannerIoTmpDirWithMultipleFlinkTmpDirs() { + void testGetClientScannerIoTmpDirForMultipleReadersOnSameTaskManager() { org.apache.flink.configuration.Configuration commaSeparatedFlinkConfig = new org.apache.flink.configuration.Configuration() .set(TMP_DIRS, "/flink_tmp_dir_0,/flink_tmp_dir_1"); - assertThat( - FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( - new Configuration(), commaSeparatedFlinkConfig, 0)) - .isEqualTo(new File("/flink_tmp_dir_0", "fluss").getAbsolutePath()); - assertThat( - FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( - new Configuration(), commaSeparatedFlinkConfig, 1)) - .isEqualTo(new File("/flink_tmp_dir_1", "fluss").getAbsolutePath()); - assertThat( - FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( - new Configuration(), commaSeparatedFlinkConfig, 2)) - .isEqualTo(new File("/flink_tmp_dir_0", "fluss").getAbsolutePath()); + Set selectedDirectories = new HashSet<>(); + for (int i = 0; i < 100; i++) { + selectedDirectories.add( + FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( + new Configuration(), commaSeparatedFlinkConfig)); + } + assertThat(selectedDirectories) + .containsExactlyInAnyOrder( + new File("/flink_tmp_dir_0", "fluss").getAbsolutePath(), + new File("/flink_tmp_dir_1", "fluss").getAbsolutePath()); org.apache.flink.configuration.Configuration pathSeparatorFlinkConfig = new org.apache.flink.configuration.Configuration() @@ -136,7 +134,9 @@ void testGetClientScannerIoTmpDirWithMultipleFlinkTmpDirs() { assertThat( FlinkConnectorOptionsUtils.getClientScannerIoTmpDir( - new Configuration(), pathSeparatorFlinkConfig, 1)) - .isEqualTo(new File("/flink_tmp_dir_1", "fluss").getAbsolutePath()); + new Configuration(), pathSeparatorFlinkConfig)) + .isIn( + new File("/flink_tmp_dir_0", "fluss").getAbsolutePath(), + new File("/flink_tmp_dir_1", "fluss").getAbsolutePath()); } }