From e277dc2e398c6defd95db7aa352fbaa513e7ef6f Mon Sep 17 00:00:00 2001 From: David Wang Date: Fri, 14 Aug 2026 18:04:23 +1000 Subject: [PATCH 1/4] file-size option in Paimon source. --- .../paimon/flink/FlinkConnectorOptions.java | 38 ++++++++ .../flink/source/FlinkSourceBuilder.java | 48 ++++++++++ .../assigners/PreAssignSplitAssigner.java | 37 ++++++++ .../apache/paimon/flink/FileStoreITCase.java | 72 ++++++++++++++ .../flink/source/FlinkSourceBuilderTest.java | 94 +++++++++++++++++++ 5 files changed, 289 insertions(+) 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..6383447ce56d 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,47 @@ 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]); + 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 +508,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().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 From cbeadcb656ab79d95fb1f5e728860b4eacd882c5 Mon Sep 17 00:00:00 2001 From: David Wang Date: Mon, 24 Aug 2026 22:15:44 +1000 Subject: [PATCH 2/4] Use Flink 2 compatible task info API in test --- .../src/test/java/org/apache/paimon/flink/FileStoreITCase.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 6383447ce56d..40015fecf3c7 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 @@ -528,7 +528,7 @@ private static class SubtaskAndPayloadSize extends RichMapFunction @Override public Row map(RowData value) { return Row.of( - getRuntimeContext().getIndexOfThisSubtask(), + getRuntimeContext().getTaskInfo().getIndexOfThisSubtask(), value.getInt(0), value.getString(1).toString().length()); } From b609301f258914f8dfc3ee8f8512ef9b869e6d0c Mon Sep 17 00:00:00 2001 From: David Wang Date: Tue, 25 Aug 2026 08:48:50 +1000 Subject: [PATCH 3/4] Document file-size split weight option --- docs/generated/flink_connector_configuration.html | 6 ++++++ 1 file changed, 6 insertions(+) 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) From 2c3980c1f5537765e8f0c5489e892d720f66ebbe Mon Sep 17 00:00:00 2001 From: David Wang Date: Tue, 25 Aug 2026 12:16:55 +1000 Subject: [PATCH 4/4] Clarify file-size split assignment test --- .../src/test/java/org/apache/paimon/flink/FileStoreITCase.java | 1 + 1 file changed, 1 insertion(+) 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 40015fecf3c7..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 @@ -236,6 +236,7 @@ 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);