This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new d27a54c152 [flink][core] Introduce file-size option in Paimon source 
(#9368)
d27a54c152 is described below

commit d27a54c152db1ff529e727f9e6832afba12bcda7
Author: dwangatt <[email protected]>
AuthorDate: Mon Sep 21 13:29:55 2026 +1000

    [flink][core] Introduce file-size option in Paimon source (#9368)
---
 docs/generated/flink_connector_configuration.html  |   6 ++
 .../java/org/apache/paimon/utils/BinPacking.java   |   4 +-
 .../org/apache/paimon/utils/BinPackingTest.java    |  18 +++-
 .../apache/paimon/flink/FlinkConnectorOptions.java |  38 +++++++
 .../paimon/flink/source/FlinkSourceBuilder.java    |   2 +
 .../paimon/flink/source/SplitWeightUtils.java      |  82 +++++++++++++++
 .../paimon/flink/source/SystemTableSource.java     |  39 ++++++++
 .../source/assigners/PreAssignSplitAssigner.java   |  37 +++++++
 .../org/apache/paimon/flink/FileStoreITCase.java   |  73 ++++++++++++++
 .../paimon/flink/source/DataTableSourceTest.java   |  97 ++++++++++++++++++
 .../flink/source/FlinkSourceBuilderTest.java       | 110 +++++++++++++++++++++
 11 files changed, 501 insertions(+), 5 deletions(-)

diff --git a/docs/generated/flink_connector_configuration.html 
b/docs/generated/flink_connector_configuration.html
index 742267c00d..37638d7f99 100644
--- a/docs/generated/flink_connector_configuration.html
+++ b/docs/generated/flink_connector_configuration.html
@@ -212,6 +212,12 @@ under the License.
             <td><p>Enum</p></td>
             <td>The mode used by StaticFileStoreSplitEnumerator to assign 
splits.<br /><br />Possible values:<ul><li>"fair": Distribute splits evenly 
when batch reading to prevent a few tasks from reading 
all.</li><li>"preemptive": Distribute splits preemptively according to the 
consumption speed of the task.</li></ul></td>
         </tr>
+        <tr>
+            <td><h5>scan.split-enumerator.weight-mode</h5></td>
+            <td style="word-wrap: break-word;">row-count</td>
+            <td><p>Enum</p></td>
+            <td>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.<br /><br />Possible 
values:<ul><li>"row-count": Balance splits by row count.</li><li>"file-size": 
Balance splits by total data file size for DataSplit and fall back to row count 
otherwise. Only works with fair assign mode.< [...]
+        </tr>
         <tr>
             <td><h5>scan.watermark.alignment.group</h5></td>
             <td style="word-wrap: break-word;">(none)</td>
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/utils/BinPacking.java 
b/paimon-common/src/main/java/org/apache/paimon/utils/BinPacking.java
index e7419835d7..35c979f2df 100644
--- a/paimon-common/src/main/java/org/apache/paimon/utils/BinPacking.java
+++ b/paimon-common/src/main/java/org/apache/paimon/utils/BinPacking.java
@@ -57,13 +57,13 @@ public class BinPacking {
         return packed;
     }
 
-    /** A bin packing implementation for fixed bin number. */
+    /** A bin packing implementation for fixed bin number using longest 
processing time first. */
     public static <T> List<List<T>> packForFixedBinNumber(
             Iterable<T> items, Function<T, Long> weightFunc, int binNumber) {
         // 1. sort items first
         List<T> sorted = new ArrayList<>();
         items.forEach(sorted::add);
-        sorted.sort(comparingLong(weightFunc::apply));
+        sorted.sort(comparingLong(weightFunc::apply).reversed());
 
         // 2. packing
         PriorityQueue<FixedNumberBin<T>> bins = new PriorityQueue<>();
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/utils/BinPackingTest.java 
b/paimon-common/src/test/java/org/apache/paimon/utils/BinPackingTest.java
index ec850b8c2d..4deafc3636 100644
--- a/paimon-common/src/test/java/org/apache/paimon/utils/BinPackingTest.java
+++ b/paimon-common/src/test/java/org/apache/paimon/utils/BinPackingTest.java
@@ -20,7 +20,9 @@ package org.apache.paimon.utils;
 
 import org.junit.jupiter.api.Test;
 
+import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.List;
 
 import static org.assertj.core.api.Assertions.assertThat;
@@ -33,9 +35,19 @@ public class BinPackingTest {
         List<List<Integer>> pack =
                 BinPacking.packForFixedBinNumber(
                         Arrays.asList(1, 5, 1, 2, 3, 6, 2), 
Integer::longValue, 3);
-        assertThat(pack)
-                .containsExactlyInAnyOrder(
-                        Arrays.asList(1, 3), Arrays.asList(2, 5), 
Arrays.asList(1, 2, 6));
+        assertThat(pack.stream().mapToInt(bin -> bin.stream().mapToInt(i -> 
i).sum()))
+                .containsExactlyInAnyOrder(6, 7, 7);
+    }
+
+    @Test
+    public void testPackForFixedBinNumberAssignsLargestItemsFirst() {
+        List<Integer> items = new ArrayList<>(Collections.nCopies(100, 1));
+        items.add(100);
+
+        List<List<Integer>> pack = BinPacking.packForFixedBinNumber(items, 
Integer::longValue, 2);
+
+        assertThat(pack.stream().mapToInt(bin -> bin.stream().mapToInt(i -> 
i).sum()))
+                .containsExactlyInAnyOrder(100, 100);
     }
 
     @Test
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 2eb488dd7d..7bb4f5c468 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<SplitWeightMode> 
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<Boolean> SINK_USE_MANAGED_MEMORY =
             ConfigOptions.key("sink.use-managed-memory-allocator")
@@ -682,6 +692,34 @@ public class FlinkConnectorOptions {
         }
     }
 
+    /**
+     * 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 ec09badf9e..4f8a10fd9f 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
@@ -227,6 +227,8 @@ public class FlinkSourceBuilder {
                         
options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE),
                         dynamicPartitionFilteringInfo,
                         outerProject(),
+                        SplitWeightUtils.splitWeightFunc(options),
+                        null,
                         options.get(CoreOptions.BLOB_AS_DESCRIPTOR),
                         skipPreloadTargetSnapshot));
     }
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SplitWeightUtils.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SplitWeightUtils.java
new file mode 100644
index 0000000000..c1357552f2
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SplitWeightUtils.java
@@ -0,0 +1,82 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.flink.source;
+
+import org.apache.paimon.annotation.VisibleForTesting;
+import org.apache.paimon.flink.FlinkConnectorOptions;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.QueryAuthSplit;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.utils.SerializableFunction;
+
+import static org.apache.paimon.utils.Preconditions.checkArgument;
+
+/** Utilities for parsing and applying split weight options. */
+final class SplitWeightUtils {
+
+    private SplitWeightUtils() {}
+
+    static SerializableFunction<FileStoreSourceSplit, Long> 
splitWeightFunc(Options options) {
+        return splitWeightFunc(
+                options, 
options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE));
+    }
+
+    static SerializableFunction<FileStoreSourceSplit, Long> splitWeightFunc(
+            Options options, FlinkConnectorOptions.SplitAssignMode 
splitAssignMode) {
+        validateSplitWeightMode(options, splitAssignMode);
+        switch 
(options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE)) {
+            case FILE_SIZE:
+                return SplitWeightUtils::splitFileSizeOrRowCount;
+            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, FlinkConnectorOptions.SplitAssignMode 
splitAssignMode) {
+        checkArgument(
+                
options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE)
+                                != 
FlinkConnectorOptions.SplitWeightMode.FILE_SIZE
+                        || splitAssignMode == 
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);
+    }
+
+    @VisibleForTesting
+    static long splitFileSizeOrRowCount(FileStoreSourceSplit sourceSplit) {
+        Split split = sourceSplit.split();
+        while (split instanceof QueryAuthSplit) {
+            split = ((QueryAuthSplit) split).split();
+        }
+        if (split instanceof DataSplit) {
+            return ((DataSplit) split)
+                    .dataFiles().stream().mapToLong(file -> 
file.fileSize()).sum();
+        }
+        return split.rowCount();
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SystemTableSource.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SystemTableSource.java
index f6cba0e696..d4fbd333a2 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SystemTableSource.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/SystemTableSource.java
@@ -28,6 +28,7 @@ import org.apache.paimon.predicate.Predicate;
 import org.apache.paimon.table.DataTable;
 import org.apache.paimon.table.Table;
 import org.apache.paimon.table.source.ReadBuilder;
+import org.apache.paimon.utils.SerializableFunction;
 
 import org.apache.flink.api.common.eventtime.WatermarkStrategy;
 import org.apache.flink.api.connector.source.Boundedness;
@@ -47,6 +48,7 @@ public class SystemTableSource extends FlinkTableSource {
     private final boolean unbounded;
     private final int splitBatchSize;
     private final FlinkConnectorOptions.SplitAssignMode splitAssignMode;
+    @Nullable private final SerializableFunction<FileStoreSourceSplit, Long> 
splitWeightFunc;
     private final ObjectIdentifier tableIdentifier;
 
     public SystemTableSource(Table table, boolean unbounded, ObjectIdentifier 
tableIdentifier) {
@@ -55,6 +57,10 @@ public class SystemTableSource extends FlinkTableSource {
         Options options = Options.fromMap(table.options());
         this.splitBatchSize = 
options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_BATCH_SIZE);
         this.splitAssignMode = 
options.get(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE);
+        this.splitWeightFunc =
+                usesStaticSource(table, unbounded)
+                        ? SplitWeightUtils.splitWeightFunc(options)
+                        : null;
         this.tableIdentifier = tableIdentifier;
     }
 
@@ -67,10 +73,36 @@ public class SystemTableSource extends FlinkTableSource {
             int splitBatchSize,
             FlinkConnectorOptions.SplitAssignMode splitAssignMode,
             ObjectIdentifier tableIdentifier) {
+        this(
+                table,
+                unbounded,
+                predicate,
+                projectFields,
+                limit,
+                splitBatchSize,
+                splitAssignMode,
+                usesStaticSource(table, unbounded)
+                        ? SplitWeightUtils.splitWeightFunc(
+                                Options.fromMap(table.options()), 
splitAssignMode)
+                        : null,
+                tableIdentifier);
+    }
+
+    private SystemTableSource(
+            Table table,
+            boolean unbounded,
+            @Nullable Predicate predicate,
+            @Nullable int[][] projectFields,
+            @Nullable Long limit,
+            int splitBatchSize,
+            FlinkConnectorOptions.SplitAssignMode splitAssignMode,
+            @Nullable SerializableFunction<FileStoreSourceSplit, Long> 
splitWeightFunc,
+            ObjectIdentifier tableIdentifier) {
         super(table, predicate, projectFields, limit);
         this.unbounded = unbounded;
         this.splitBatchSize = splitBatchSize;
         this.splitAssignMode = splitAssignMode;
+        this.splitWeightFunc = splitWeightFunc;
         this.tableIdentifier = tableIdentifier;
     }
 
@@ -113,6 +145,8 @@ public class SystemTableSource extends FlinkTableSource {
                             splitAssignMode,
                             null,
                             rowData,
+                            splitWeightFunc,
+                            null,
                             Boolean.parseBoolean(
                                     table.options()
                                             .getOrDefault(
@@ -148,6 +182,7 @@ public class SystemTableSource extends FlinkTableSource {
                 limit,
                 splitBatchSize,
                 splitAssignMode,
+                splitWeightFunc,
                 tableIdentifier);
     }
 
@@ -161,6 +196,10 @@ public class SystemTableSource extends FlinkTableSource {
         return unbounded;
     }
 
+    private static boolean usesStaticSource(Table table, boolean unbounded) {
+        return !unbounded || !(table instanceof DataTable);
+    }
+
     private static boolean isUnordered(Table table) {
         if (!table.primaryKeys().isEmpty()) {
             return false;
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 c0077e89ee..dbf1b65d6c 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
@@ -25,6 +25,8 @@ import org.apache.paimon.utils.SerializableFunction;
 
 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;
 
@@ -49,6 +51,8 @@ import static 
org.apache.paimon.flink.utils.TableScanUtils.getSnapshotId;
  */
 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 class PreAssignSplitAssigner implements 
SplitAssigner {
         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<Integer, LinkedList<FileStoreSourceSplit>> assignment,
+            int parallelism,
+            int totalSplits,
+            SerializableFunction<FileStoreSourceSplit, Long> weightFunc) {
+        if (!LOG.isInfoEnabled()) {
+            return;
+        }
+
+        long totalWeight = 0L;
+        List<Integer> splitCounts = new ArrayList<>(parallelism);
+        List<Long> assignedWeights = new ArrayList<>(parallelism);
+        for (int i = 0; i < parallelism; i++) {
+            Collection<FileStoreSourceSplit> 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<FileStoreSourceSplit> 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 1c7015168b..f5a4e0ca8e 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.FileSystemSchemaManager;
 import org.apache.paimon.schema.Schema;
 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 class FileStoreITCase extends AbstractTestBase {
         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<String, String> 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<Row> results =
+                executeAndCollectRow(
+                        new FlinkSourceBuilder(table)
+                                .sourceBounded(true)
+                                .env(env)
+                                .build()
+                                .map(new SubtaskAndPayloadSize())
+                                .setParallelism(2));
+
+        Map<Integer, Integer> 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 @@ public class FileStoreITCase extends AbstractTestBase {
         
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<RowData, Row> {
+
+        @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/DataTableSourceTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/DataTableSourceTest.java
index e3cbb43078..0aca92741e 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/DataTableSourceTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/DataTableSourceTest.java
@@ -24,6 +24,7 @@ import org.apache.paimon.flink.PaimonDataStreamScanProvider;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.schema.FileSystemSchemaManager;
 import org.apache.paimon.schema.Schema;
 import org.apache.paimon.schema.SchemaManager;
@@ -32,9 +33,11 @@ import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.FileStoreTableFactory;
 import org.apache.paimon.table.sink.InnerTableWrite;
 import org.apache.paimon.table.sink.TableCommitImpl;
+import org.apache.paimon.table.source.DataSplit;
 import org.apache.paimon.table.system.AuditLogTable;
 import org.apache.paimon.table.system.ReadOptimizedTable;
 import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.utils.SerializableFunction;
 
 import org.apache.paimon.shade.guava30.com.google.common.collect.ImmutableMap;
 
@@ -214,6 +217,89 @@ class DataTableSourceTest {
         assertThat(sourceStream1.getParallelism()).isEqualTo(3);
     }
 
+    @Test
+    void testBoundedSystemTableUsesFileSizeWeightModeAfterCopy() throws 
Exception {
+        FileStoreTable fileStoreTable =
+                createTable(
+                        ImmutableMap.of(
+                                "bucket",
+                                "1",
+                                "bucket-key",
+                                "a",
+                                
FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE.key(),
+                                
FlinkConnectorOptions.SplitAssignMode.FAIR.toString(),
+                                
FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key(),
+                                
FlinkConnectorOptions.SplitWeightMode.FILE_SIZE.toString()));
+        SystemTableSource tableSource =
+                new SystemTableSource(
+                                new ReadOptimizedTable(fileStoreTable),
+                                false,
+                                ObjectIdentifier.of("cat", "db", "table$ro"))
+                        .copy();
+
+        StaticFileStoreSource source = staticSource(tableSource);
+        java.lang.reflect.Field weightFuncField =
+                
StaticFileStoreSource.class.getDeclaredField("splitWeightFunc");
+        weightFuncField.setAccessible(true);
+        @SuppressWarnings("unchecked")
+        SerializableFunction<FileStoreSourceSplit, Long> weightFunc =
+                (SerializableFunction<FileStoreSourceSplit, Long>) 
weightFuncField.get(source);
+        FileStoreSourceSplit split =
+                new FileStoreSourceSplit(
+                        "split-1",
+                        DataSplit.builder()
+                                .withSnapshot(1L)
+                                
.withPartition(org.apache.paimon.data.BinaryRow.EMPTY_ROW)
+                                .withBucket(0)
+                                .withBucketPath("bucket-0")
+                                .withDataFiles(
+                                        Collections.singletonList(
+                                                DataFileMeta.forAppend(
+                                                        "file-1",
+                                                        35L,
+                                                        1000L,
+                                                        null,
+                                                        0L,
+                                                        0L,
+                                                        0L,
+                                                        
Collections.emptyList(),
+                                                        null,
+                                                        null,
+                                                        null,
+                                                        null,
+                                                        null,
+                                                        null)))
+                                .build());
+
+        assertThat(weightFunc.apply(split)).isEqualTo(35L);
+    }
+
+    @Test
+    void testBoundedSystemTableRejectsFileSizeWithPreemptiveMode() throws 
Exception {
+        FileStoreTable fileStoreTable =
+                createTable(
+                        ImmutableMap.of(
+                                "bucket",
+                                "1",
+                                "bucket-key",
+                                "a",
+                                
FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE.key(),
+                                
FlinkConnectorOptions.SplitAssignMode.PREEMPTIVE.toString(),
+                                
FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key(),
+                                
FlinkConnectorOptions.SplitWeightMode.FILE_SIZE.toString()));
+
+        assertThatThrownBy(
+                        () ->
+                                new SystemTableSource(
+                                        new ReadOptimizedTable(fileStoreTable),
+                                        false,
+                                        ObjectIdentifier.of("cat", "db", 
"table$ro")))
+                .isInstanceOf(IllegalArgumentException.class)
+                
.hasMessageContaining(FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_WEIGHT_MODE.key())
+                .hasMessageContaining(
+                        
FlinkConnectorOptions.SCAN_SPLIT_ENUMERATOR_ASSIGN_MODE.key());
+    }
+
     @Test
     public void testSystemTableSourceUnorderedForBucketUnawareTable() throws 
Exception {
         // bucket = -1 (BUCKET_UNAWARE append-only table) wrapped in a system 
table should produce
@@ -294,6 +380,17 @@ class DataTableSourceTest {
                         });
     }
 
+    private StaticFileStoreSource staticSource(SystemTableSource tableSource) {
+        PaimonDataStreamScanProvider runtimeProvider = 
runtimeProvider(tableSource);
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.createLocalEnvironment();
+        DataStream<RowData> sourceStream =
+                runtimeProvider.produceDataStream(s -> Optional.empty(), env);
+        return (StaticFileStoreSource)
+                
((org.apache.flink.streaming.api.transformations.SourceTransformation<?, ?, ?>)
+                                sourceStream.getTransformation())
+                        .getSource();
+    }
+
     private FileStoreTable createTable(Map<String, String> options) throws 
Exception {
         FileIO fileIO = LocalFileIO.create();
         Path tablePath = new Path(path.toString());
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 d2f708f2d1..dd93f1bab9 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
@@ -24,7 +24,9 @@ import org.apache.paimon.catalog.CatalogFactory;
 import org.apache.paimon.catalog.Identifier;
 import org.apache.paimon.catalog.TableQueryAuthResult;
 import org.apache.paimon.data.BinaryRow;
+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.FallbackReadFileStoreTable;
 import org.apache.paimon.table.Table;
@@ -53,7 +55,11 @@ import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.ValueSource;
 
 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;
@@ -112,6 +118,110 @@ public class FlinkSourceBuilderTest {
         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(SplitWeightUtils.splitFileSizeOrRowCount(split)).isEqualTo(35L);
+    }
+
+    @Test
+    public void testSplitFileSizeOrRowCountUnwrapsQueryAuthSplit() {
+        DataSplit dataSplit =
+                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();
+        FileStoreSourceSplit split =
+                new FileStoreSourceSplit("split-1", new 
QueryAuthSplit(dataSplit, null));
+
+        
assertThat(SplitWeightUtils.splitFileSizeOrRowCount(split)).isEqualTo(35L);
+    }
+
+    @Test
+    public void testSplitFileSizeOrRowCountFallsBackToRowCount() {
+        FileStoreSourceSplit split = new FileStoreSourceSplit("split-1", new 
TestSplit(123L));
+
+        
assertThat(SplitWeightUtils.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<String, String> 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

Reply via email to