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]