Copilot commented on code in PR #3617:
URL: https://github.com/apache/fluss/pull/3617#discussion_r3542046414


##########
fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java:
##########
@@ -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];
     }

Review Comment:
   splitPaths() returns raw split tokens without trimming or filtering empty 
entries. A value like "/tmp1, /tmp2" (or trailing separators) can yield paths 
with leading whitespace or empty strings, which may resolve to unintended 
locations. Trim and drop empty tokens before returning.



##########
fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtils.java:
##########
@@ -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();
+                        });

Review Comment:
   Path joining uses a child path starting with a leading separator ("/fluss"), 
which can be treated as an absolute path on some platforms and can ignore the 
selected Flink tmp dir. Use a relative child name ("fluss") and use 
Math.floorMod to avoid negative modulo edge cases.



##########
fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/FlinkConnectorOptionsUtilTest.java:
##########
@@ -96,13 +97,46 @@ 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.
         assertThat(
                         FlinkConnectorOptionsUtils.getClientScannerIoTmpDir(
                                 new Configuration(), flinkConfig))
                 .isEqualTo("/flink_tmp_dir/fluss");
+        
assertThat(FlinkConnectorOptionsUtils.getClientScannerIoTmpDir(flussConfig, 
flinkConfig, 1))

Review Comment:
   This assertion block hard-codes a Unix-style expected path and also 
introduces a very long line that may violate formatting checks. Prefer building 
the expected path with File for portability and wrap the call for readability.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to