Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -290,7 +290,8 @@ public SourceReader<OUT, SourceSplitBase> createReader(SourceReaderContext conte

flussConf.set(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Just a personal question, might be missing something here — since flussConf is a member field of the source, calling set seems to mutate it in place, and this config eventually gets passed down to the Connection and shared by the downstream scanners. So I'm wondering whether multiple readers reusing the same instance could end up sharing the same config. Would it be safer to copy it first before setting?
Curious how you see it @luoyuxia ?

@zuston zuston Jul 15, 2026

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

aha. this logic is aligned with the previous impl. I think it's OK to share

CLIENT_SCANNER_IO_TMP_DIR,
getClientScannerIoTmpDir(flussConf, context.getConfiguration()));
getClientScannerIoTmpDir(
flussConf, context.getConfiguration(), context.getIndexOfSubtask()));
deserializationSchema.open(
new DeserializerInitContextImpl(
context.metricGroup().addGroup("deserializer"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,10 @@ public SourceReader<TableBucketWriteResult<WriteResult>, 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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. * */
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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))
Comment on lines 107 to +108
.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());
}
}
Loading