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 227db3fabc..e65f4ff228 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 649b758704..4e34bdfd59 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 249e1be089..3873821e56 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; @@ -41,7 +40,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; @@ -201,13 +199,23 @@ 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 flussConf.getString(CLIENT_SCANNER_IO_TMP_DIR); + 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(); + }); } /** 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 0e06dc1019..e0502f0b5e 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()); } }