diff --git a/docs/generated/flink_connector_configuration.html b/docs/generated/flink_connector_configuration.html
index 3d7c5149a480..a3a8780a1ced 100644
--- a/docs/generated/flink_connector_configuration.html
+++ b/docs/generated/flink_connector_configuration.html
@@ -212,6 +212,12 @@
Enum |
The mode used by StaticFileStoreSplitEnumerator to assign splits.
Possible values:- "fair": Distribute splits evenly when batch reading to prevent a few tasks from reading all.
- "preemptive": Distribute splits preemptively according to the consumption speed of the task.
|
+
+ scan.split-enumerator.weight-mode |
+ row-count |
+ Enum |
+ The weight metric used by StaticFileStoreSplitEnumerator. 'row-count' balances by split row count. 'file-size' only works with 'scan.split-enumerator.mode' = 'fair', balances by total data file size for DataSplit, and falls back to row count otherwise.
Possible values:- "row-count": Balance splits by row count.
- "file-size": Balance splits by total data file size for DataSplit and fall back to row count otherwise. Only works with fair assign mode.
|
+
scan.watermark.alignment.group |
(none) |
diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
index 2eb488dd7d3d..7bb4f5c468c2 100644
--- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
+++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
@@ -181,6 +181,16 @@ public class FlinkConnectorOptions {
.withDescription(
"The mode used by StaticFileStoreSplitEnumerator to assign splits.");
+ public static final ConfigOption SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE =
+ key("scan.split-enumerator.weight-mode")
+ .enumType(SplitWeightMode.class)
+ .defaultValue(SplitWeightMode.ROW_COUNT)
+ .withDescription(
+ "The weight metric used by StaticFileStoreSplitEnumerator. "
+ + "'row-count' balances by split row count. "
+ + "'file-size' only works with 'scan.split-enumerator.mode' = 'fair', "
+ + "balances by total data file size for DataSplit, and falls back to row count otherwise.");
+
/* Sink writer allocate segments from managed memory. */
public static final ConfigOption SINK_USE_MANAGED_MEMORY =
ConfigOptions.key("sink.use-managed-memory-allocator")
@@ -682,6 +692,34 @@ public InlineElement getDescription() {
}
}
+ /**
+ * Split weight mode for {@link org.apache.paimon.flink.source.StaticFileStoreSplitEnumerator}.
+ */
+ public enum SplitWeightMode implements DescribedEnum {
+ ROW_COUNT("row-count", "Balance splits by row count."),
+ FILE_SIZE(
+ "file-size",
+ "Balance splits by total data file size for DataSplit and fall back to row count otherwise. Only works with fair assign mode.");
+
+ private final String value;
+ private final String description;
+
+ SplitWeightMode(String value, String description) {
+ this.value = value;
+ this.description = description;
+ }
+
+ @Override
+ public String toString() {
+ return value;
+ }
+
+ @Override
+ public InlineElement getDescription() {
+ return text(description);
+ }
+ }
+
/**
* Split assign mode for {@link org.apache.paimon.flink.source.StaticFileStoreSplitEnumerator}.
*/
diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/FlinkSourceBuilder.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/FlinkSourceBuilder.java
index eb3da3e16358..3857e7b9c1b7 100644
--- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/FlinkSourceBuilder.java
+++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/FlinkSourceBuilder.java
@@ -38,6 +38,8 @@
import org.apache.paimon.table.source.PostponeMergePlan;
import org.apache.paimon.table.source.PostponeMergeReadBuilder;
import org.apache.paimon.table.source.ReadBuilder;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.utils.SerializableFunction;
import org.apache.paimon.utils.StringUtils;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
@@ -219,6 +221,7 @@ private ReadBuilder createReadBuilder(@Nullable org.apache.paimon.types.RowType
private DataStream buildStaticFileSource() {
Options options = Options.fromMap(table.options());
+ validateSplitWeightMode(options);
return toDataStream(
new StaticFileStoreSource(
createReadBuilder(projectedRowType()),
@@ -227,10 +230,55 @@ private DataStream buildStaticFileSource() {
options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE),
dynamicPartitionFilteringInfo,
outerProject(),
+ splitWeightFunc(options),
+ null,
options.get(CoreOptions.BLOB_AS_DESCRIPTOR),
skipPreloadTargetSnapshot));
}
+ private static SerializableFunction splitWeightFunc(
+ Options options) {
+ if (isFileSizeWeightMode(options)) {
+ return FlinkSourceBuilder::splitFileSizeOrRowCount;
+ }
+ switch (options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE)) {
+ case ROW_COUNT:
+ return split -> split.split().rowCount();
+ default:
+ throw new UnsupportedOperationException(
+ "Unsupported split weight mode "
+ + options.get(
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE));
+ }
+ }
+
+ private static void validateSplitWeightMode(Options options) {
+ checkArgument(
+ !isFileSizeWeightMode(options)
+ || options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE)
+ == FlinkConnectorOptions.SplitAssignMode.FAIR,
+ "'%s' = '%s' only works with '%s' = '%s'.",
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key(),
+ FlinkConnectorOptions.SplitWeightMode.FILE_SIZE,
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE.key(),
+ FlinkConnectorOptions.SplitAssignMode.FAIR);
+ }
+
+ private static boolean isFileSizeWeightMode(Options options) {
+ return options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE)
+ == FlinkConnectorOptions.SplitWeightMode.FILE_SIZE;
+ }
+
+ @VisibleForTesting
+ static long splitFileSizeOrRowCount(FileStoreSourceSplit sourceSplit) {
+ Split split = sourceSplit.split();
+ if (split instanceof DataSplit) {
+ return ((DataSplit) split)
+ .dataFiles().stream().mapToLong(file -> file.fileSize()).sum();
+ }
+ return split.rowCount();
+ }
+
private @Nullable DataStream buildPostponeMergeSource() {
FileStoreTable fileStoreTable = (FileStoreTable) table;
if (fileStoreTable.coreOptions().startupMode() == CoreOptions.StartupMode.COMPACTED_FULL) {
diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/assigners/PreAssignSplitAssigner.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/assigners/PreAssignSplitAssigner.java
index 24ac4a291164..53f4a5914e63 100644
--- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/assigners/PreAssignSplitAssigner.java
+++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/assigners/PreAssignSplitAssigner.java
@@ -28,6 +28,8 @@
import org.apache.flink.api.connector.source.SplitEnumeratorContext;
import org.apache.flink.table.connector.source.DynamicFilteringData;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import javax.annotation.Nullable;
@@ -52,6 +54,8 @@
*/
public class PreAssignSplitAssigner implements SplitAssigner {
+ private static final Logger LOG = LoggerFactory.getLogger(PreAssignSplitAssigner.class);
+
/** Default batch splits size to avoid exceed `akka.framesize`. */
private final int splitBatchSize;
@@ -145,9 +149,42 @@ public PreAssignSplitAssigner(
this.groupFunc = groupFunc;
this.pendingSplitAssignment =
createBatchFairSplitAssignment(splits, parallelism, this.weightFunc, groupFunc);
+ logSplitAssignmentSummary(
+ this.pendingSplitAssignment, parallelism, splits.size(), this.weightFunc);
this.numberOfPendingSplits = new AtomicInteger(splits.size());
}
+ private static void logSplitAssignmentSummary(
+ Map> assignment,
+ int parallelism,
+ int totalSplits,
+ SerializableFunction weightFunc) {
+ if (!LOG.isInfoEnabled()) {
+ return;
+ }
+
+ long totalWeight = 0L;
+ List splitCounts = new ArrayList<>(parallelism);
+ List assignedWeights = new ArrayList<>(parallelism);
+ for (int i = 0; i < parallelism; i++) {
+ Collection assignedSplits =
+ assignment.getOrDefault(i, new LinkedList<>());
+ long assignedWeight = assignedSplits.stream().mapToLong(weightFunc::apply).sum();
+ splitCounts.add(assignedSplits.size());
+ assignedWeights.add(assignedWeight);
+ totalWeight += assignedWeight;
+ }
+
+ LOG.info(
+ "Created FAIR split assignment summary: parallelism={}, totalSplits={}, "
+ + "totalWeight={}, splitCountsPerSubtask={}, assignedWeightsPerSubtask={}",
+ parallelism,
+ totalSplits,
+ totalWeight,
+ splitCounts,
+ assignedWeights);
+ }
+
@Override
public List getNext(int subtask, @Nullable String hostname) {
// The following batch assignment operation is for two purposes:
diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FileStoreITCase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FileStoreITCase.java
index 0d9f512fbee9..7c7618b5c4b3 100644
--- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FileStoreITCase.java
+++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FileStoreITCase.java
@@ -19,6 +19,8 @@
package org.apache.paimon.flink;
import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericRow;
import org.apache.paimon.flink.sink.FixedBucketSink;
import org.apache.paimon.flink.sink.FlinkSinkBuilder;
import org.apache.paimon.flink.source.ContinuousFileStoreSource;
@@ -32,11 +34,14 @@
import org.apache.paimon.schema.SchemaManager;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.FileStoreTableFactory;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchTableWrite;
import org.apache.paimon.utils.BlockingIterator;
import org.apache.paimon.utils.FailingFileIO;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
+import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.api.connector.source.Boundedness;
import org.apache.flink.api.dag.Transformation;
import org.apache.flink.streaming.api.datastream.DataStream;
@@ -226,6 +231,48 @@ public void testNonPartitioned() throws Exception {
assertThat(results).containsExactlyInAnyOrder(expected);
}
+ @TestTemplate
+ public void testFileSizeSplitWeightModeForBoundedSource() throws Exception {
+ assumeTrue(isBatch);
+
+ FileStoreTable table = buildFileStoreTable(new int[0], new int[0]);
+ // Use equal row counts with skewed payload sizes to verify byte-aware assignment.
+ writeSingleRecordFile(table, 1, repeat("a", 8), 1);
+ writeSingleRecordFile(table, 2, repeat("b", 8), 2);
+ writeSingleRecordFile(table, 3, repeat("c", 32 * 1024), 3);
+ writeSingleRecordFile(table, 4, repeat("d", 32 * 1024), 4);
+
+ Map options = new HashMap<>();
+ options.put(CoreOptions.SOURCE_SPLIT_TARGET_SIZE.key(), "1 B");
+ options.put(CoreOptions.SOURCE_SPLIT_OPEN_FILE_COST.key(), "1 B");
+ options.put(FlinkConnectorOptions.SCAN_PARALLELISM.key(), "2");
+ options.put(
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key(),
+ FlinkConnectorOptions.SplitWeightMode.FILE_SIZE.toString());
+ table = table.copy(options);
+
+ List results =
+ executeAndCollectRow(
+ new FlinkSourceBuilder(table)
+ .sourceBounded(true)
+ .env(env)
+ .build()
+ .map(new SubtaskAndPayloadSize())
+ .setParallelism(2));
+
+ Map largePayloadSubtasks = new HashMap<>();
+ for (Row row : results) {
+ int subtask = (int) row.getField(0);
+ int payloadSize = (int) row.getField(2);
+ if (payloadSize > 1024) {
+ largePayloadSubtasks.put((int) row.getField(1), subtask);
+ }
+ }
+
+ assertThat(largePayloadSubtasks).hasSize(2);
+ assertThat(largePayloadSubtasks.values()).containsExactlyInAnyOrder(0, 1);
+ }
+
@TestTemplate
public void testOverwrite() throws Exception {
assumeTrue(isBatch);
@@ -462,6 +509,32 @@ private void sinkAndValidate(
assertThat(iterator.collect(expected.length)).containsExactlyInAnyOrder(expected);
}
+ private static void writeSingleRecordFile(FileStoreTable table, int v, String p, int k)
+ throws Exception {
+ try (BatchTableWrite write = table.newBatchWriteBuilder().newWrite();
+ BatchTableCommit commit = table.newBatchWriteBuilder().newCommit()) {
+ write.write(GenericRow.of(v, BinaryString.fromString(p), k));
+ commit.commit(write.prepareCommit());
+ }
+ }
+
+ private static String repeat(String value, int count) {
+ char[] chars = new char[count];
+ Arrays.fill(chars, value.charAt(0));
+ return new String(chars);
+ }
+
+ private static class SubtaskAndPayloadSize extends RichMapFunction {
+
+ @Override
+ public Row map(RowData value) {
+ return Row.of(
+ getRuntimeContext().getTaskInfo().getIndexOfThisSubtask(),
+ value.getInt(0),
+ value.getString(1).toString().length());
+ }
+ }
+
public FileStoreTable buildFileStoreTable(int[] partitions, int[] primaryKey) throws Exception {
return buildFileStoreTable(isBatch, getTempDirPath(), partitions, primaryKey);
}
diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/FlinkSourceBuilderTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/FlinkSourceBuilderTest.java
index 4f8a097e4724..17c58c1de000 100644
--- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/FlinkSourceBuilderTest.java
+++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/FlinkSourceBuilderTest.java
@@ -22,9 +22,13 @@
import org.apache.paimon.catalog.CatalogContext;
import org.apache.paimon.catalog.CatalogFactory;
import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.flink.FlinkConnectorOptions;
import org.apache.paimon.flink.source.operator.MonitorSource;
+import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.schema.Schema;
import org.apache.paimon.table.Table;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.Split;
import org.apache.paimon.types.DataTypes;
import org.apache.flink.api.dag.Transformation;
@@ -38,6 +42,11 @@
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Path;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.OptionalLong;
import static org.apache.paimon.flink.LogicalTypeConversion.toLogicalType;
import static org.assertj.core.api.Assertions.assertThat;
@@ -92,6 +101,91 @@ private Table createTable(
return catalog.getTable(identifier);
}
+ @Test
+ public void testSplitFileSizeOrRowCountUsesDataSplitFileSize() {
+ FileStoreSourceSplit split =
+ new FileStoreSourceSplit(
+ "split-1",
+ DataSplit.builder()
+ .withSnapshot(1L)
+ .withPartition(org.apache.paimon.data.BinaryRow.EMPTY_ROW)
+ .withBucket(0)
+ .withBucketPath("bucket-0")
+ .withDataFiles(
+ Arrays.asList(
+ dataFile("file-1", 10L, 1L),
+ dataFile("file-2", 25L, 1000L)))
+ .build());
+
+ assertThat(FlinkSourceBuilder.splitFileSizeOrRowCount(split)).isEqualTo(35L);
+ }
+
+ @Test
+ public void testSplitFileSizeOrRowCountFallsBackToRowCount() {
+ FileStoreSourceSplit split = new FileStoreSourceSplit("split-1", new TestSplit(123L));
+
+ assertThat(FlinkSourceBuilder.splitFileSizeOrRowCount(split)).isEqualTo(123L);
+ }
+
+ private static DataFileMeta dataFile(String fileName, long fileSize, long rowCount) {
+ return DataFileMeta.forAppend(
+ fileName,
+ fileSize,
+ rowCount,
+ null,
+ 0L,
+ 0L,
+ 0L,
+ Collections.emptyList(),
+ null,
+ null,
+ null,
+ null,
+ null,
+ null);
+ }
+
+ private static class TestSplit implements Split {
+
+ private final long rowCount;
+
+ private TestSplit(long rowCount) {
+ this.rowCount = rowCount;
+ }
+
+ @Override
+ public long rowCount() {
+ return rowCount;
+ }
+
+ @Override
+ public OptionalLong mergedRowCount() {
+ return OptionalLong.of(rowCount);
+ }
+ }
+
+ @Test
+ public void testFileSizeWeightModeOnlyWorksWithFairAssignMode() throws Exception {
+ Table table = createTable("file_size_preemptive", false, 2, false);
+ Map options = new HashMap<>();
+ options.put(
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key(),
+ FlinkConnectorOptions.SplitWeightMode.FILE_SIZE.toString());
+ options.put(
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE.key(),
+ FlinkConnectorOptions.SplitAssignMode.PREEMPTIVE.toString());
+ table = table.copy(options);
+
+ StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
+ FlinkSourceBuilder builder = new FlinkSourceBuilder(table).env(env).sourceBounded(true);
+
+ assertThatThrownBy(builder::build)
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key())
+ .hasMessageContaining(
+ FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE.key());
+ }
+
@Test
public void testUnawareBucket() throws Exception {
// pk table && bucket-append-ordered is true