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