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 930fdfaafc [flink][spark] Complete postpone merge-on-read support 
(#8849)
930fdfaafc is described below

commit 930fdfaafc5bc4c2d6f3ae4adc56bcafb679a9f4
Author: Jingsong Lee <[email protected]>
AuthorDate: Sun Jul 26 20:55:41 2026 +0800

    [flink][spark] Complete postpone merge-on-read support (#8849)
---
 .../org/apache/paimon/table/PostponeUtils.java     |  43 +++-
 .../paimon/table/source/PostponeMergePlan.java     |  10 +
 .../table/source/PostponeMergeReadBuilder.java     |  59 ++++-
 .../apache/paimon/table/source/TableScanTest.java  |  33 ++-
 .../paimon/flink/source/BaseDataTableSource.java   |  19 +-
 .../paimon/flink/source/DataTableSource.java       |   9 +-
 .../paimon/flink/source/FlinkSourceBuilder.java    | 119 +++++++++
 .../flink/source/PostponeMergeInputOperator.java   | 145 +++++++++++
 .../flink/source/PostponeMergeInputSource.java     | 127 ++++++++++
 .../paimon/flink/source/PostponeMergeOnRead.java   | 160 ++++++++++++
 .../paimon/flink/source/PostponeMergeOperator.java | 279 +++++++++++++++++++++
 .../paimon/flink/PostponeBucketTableITCase.java    | 103 +++++++-
 .../paimon/flink/source/DataTableSourceTest.java   |  99 ++++++++
 .../flink/source/FlinkSourceBuilderTest.java       |  22 ++
 .../org/apache/paimon/spark/PaimonBaseScan.scala   |  16 +-
 .../apache/paimon/spark/PaimonScanBuilder.scala    |   2 +-
 .../apache/paimon/spark/PostponeMergeOnRead.scala  |  18 +-
 .../paimon/spark/execution/PaimonStrategy.scala    |  12 +-
 .../paimon/spark/sql/PostponeBucketTableTest.scala | 138 +++++++++-
 19 files changed, 1368 insertions(+), 45 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/PostponeUtils.java 
b/paimon-core/src/main/java/org/apache/paimon/table/PostponeUtils.java
index 0463bb3f2d..5335830988 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/PostponeUtils.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/PostponeUtils.java
@@ -55,6 +55,7 @@ import java.util.Set;
 import static org.apache.paimon.CoreOptions.BUCKET;
 import static 
org.apache.paimon.CoreOptions.COMMIT_STRICT_MODE_LAST_SAFE_SNAPSHOT;
 import static org.apache.paimon.CoreOptions.WRITE_ONLY;
+import static org.apache.paimon.utils.Preconditions.checkArgument;
 
 /** Utils for postpone table. */
 public class PostponeUtils {
@@ -137,12 +138,21 @@ public class PostponeUtils {
             @Nullable PartitionPredicate partitionFilter) {
         Map<BinaryRow, Integer> knownNumBuckets =
                 getKnownNumBuckets(table, snapshotId, partitionFilter);
-        Long targetRowNumPerBucket =
-                
table.coreOptions().postponeTargetRowNumPerBucket().orElse(null);
         Map<BinaryRow, Long> postponeRowCounts =
-                targetRowNumPerBucket == null
+                
!table.coreOptions().postponeTargetRowNumPerBucket().isPresent()
                         ? Collections.emptyMap()
                         : getPostponeRowCounts(table, snapshotId, 
partitionFilter);
+        return createPostponeBucketAssigner(
+                table, knownNumBuckets, postponeRowCounts, defaultParallelism);
+    }
+
+    private static PostponeBucketAssigner createPostponeBucketAssigner(
+            FileStoreTable table,
+            Map<BinaryRow, Integer> knownNumBuckets,
+            Map<BinaryRow, Long> postponeRowCounts,
+            int defaultParallelism) {
+        Long targetRowNumPerBucket =
+                
table.coreOptions().postponeTargetRowNumPerBucket().orElse(null);
         int defaultBucketNum =
                 table.coreOptions()
                                 .toConfiguration()
@@ -158,6 +168,14 @@ public class PostponeUtils {
             long snapshotId,
             int defaultParallelism,
             @Nullable PartitionPredicate partitionFilter) {
+        return createPostponeBucketRouter(
+                table,
+                createPostponeBucketAssigner(
+                        table, snapshotId, defaultParallelism, 
partitionFilter));
+    }
+
+    private static PostponeBucketRouter createPostponeBucketRouter(
+            FileStoreTable table, PostponeBucketAssigner bucketAssigner) {
         List<String> trimmedPrimaryKeys = table.schema().trimmedPrimaryKeys();
         int[] bucketKeyMapping =
                 table.schema().bucketKeys().stream()
@@ -174,8 +192,7 @@ public class PostponeUtils {
                         
PrimaryKeyTableUtils.PrimaryKeyFieldsExtractor.EXTRACTOR.keyFields(
                                 table.schema()));
         return new PostponeBucketRouter(
-                createPostponeBucketAssigner(
-                        table, snapshotId, defaultParallelism, 
partitionFilter),
+                bucketAssigner,
                 keyType,
                 table.schema().logicalBucketKeyType(),
                 bucketKeyMapping,
@@ -376,6 +393,11 @@ public class PostponeUtils {
                     postponeRowCounts,
                     defaultBucketNum);
         }
+
+        private PostponeBucketAssigner withDefaultBucketNum(int 
newDefaultBucketNum) {
+            return new PostponeBucketAssigner(
+                    knownNumBuckets, targetRowNumPerBucket, postponeRowCounts, 
newDefaultBucketNum);
+        }
     }
 
     /** Snapshot-bound routing metadata for postpone records. */
@@ -418,6 +440,17 @@ public class PostponeUtils {
         public int numBuckets(BinaryRow partition) {
             return bucketAssigner.assign(partition);
         }
+
+        public PostponeBucketRouter withDefaultBucketNum(int 
newDefaultBucketNum) {
+            checkArgument(
+                    newDefaultBucketNum > 0, "Default postpone bucket number 
must be positive.");
+            return new PostponeBucketRouter(
+                    bucketAssigner.withDefaultBucketNum(newDefaultBucketNum),
+                    keyType,
+                    bucketKeyType,
+                    bucketKeyMapping,
+                    bucketFunctionType);
+        }
     }
 
     private static final class PostponeFile {
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergePlan.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergePlan.java
index b2e6bc460e..a65e415ee0 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergePlan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergePlan.java
@@ -92,6 +92,16 @@ public final class PostponeMergePlan implements 
TableScan.Plan {
         return numPotentialBuckets;
     }
 
+    PostponeMergePlan withDefaultBucketNum(int newDefaultBucketNum) {
+        return new PostponeMergePlan(
+                realSplits,
+                postponeSplits,
+                bucketRouter.withDefaultBucketNum(newDefaultBucketNum),
+                keyType,
+                resultReadType,
+                mergeReadType);
+    }
+
     private static long numPotentialBuckets(
             List<DataSplit> realSplits,
             List<DataSplit> postponeSplits,
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergeReadBuilder.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergeReadBuilder.java
index d634df8bac..75eda5624e 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergeReadBuilder.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/PostponeMergeReadBuilder.java
@@ -87,23 +87,56 @@ public final class PostponeMergeReadBuilder implements 
Serializable {
             return Optional.empty();
         }
 
-        SnapshotReader postponeReader =
-                table.newSnapshotReader()
-                        .withSnapshot(snapshot)
-                        .withBucket(BucketMode.POSTPONE_BUCKET);
-        if (partitionFilter != null) {
-            postponeReader.withPartitionFilter(partitionFilter);
+        PostponeMergeReadBuilder builder = new PostponeMergeReadBuilder(table, 
snapshot);
+        builder.withPartitionFilter(partitionFilter);
+        if (!builder.hasPostponeFiles()) {
+            return Optional.empty();
         }
-        if (!postponeReader.readFileIterator().hasNext()) {
+
+        validateReadMode(table);
+        return Optional.of(builder);
+    }
+
+    /**
+     * Creates a builder bound to the snapshot selected now, including when 
that snapshot contains
+     * no postpone files.
+     *
+     * <p>Execution engines which enable merge-on-read should use this method 
to avoid probing one
+     * snapshot and then falling back to an ordinary source which may select a 
newer snapshot.
+     */
+    public static Optional<PostponeMergeReadBuilder> createSnapshotBound(
+            FileStoreTable table, @Nullable PartitionPredicate 
partitionFilter) {
+        checkArgument(
+                table.bucketMode() == BucketMode.POSTPONE_MODE && 
!table.primaryKeys().isEmpty(),
+                "Postpone merge read requires a primary-key postpone bucket 
table.");
+
+        // Let the ordinary compacted-full scanner select its snapshot.
+        if (table.coreOptions().startupMode() == 
CoreOptions.StartupMode.COMPACTED_FULL) {
             return Optional.empty();
         }
 
         validateReadMode(table);
+        Snapshot snapshot = TimeTravelUtil.tryTravelOrLatest(table);
+        if (snapshot == null) {
+            return Optional.empty();
+        }
+
         PostponeMergeReadBuilder builder = new PostponeMergeReadBuilder(table, 
snapshot);
         builder.withPartitionFilter(partitionFilter);
         return Optional.of(builder);
     }
 
+    private boolean hasPostponeFiles() {
+        SnapshotReader postponeReader =
+                table.newSnapshotReader()
+                        .withSnapshot(snapshot)
+                        .withBucket(BucketMode.POSTPONE_BUCKET);
+        if (partitionFilter != null) {
+            postponeReader.withPartitionFilter(partitionFilter);
+        }
+        return postponeReader.readFileIterator().hasNext();
+    }
+
     private PostponeMergeReadBuilder withPartitionFilter(
             @Nullable PartitionPredicate partitionPredicate) {
         if (partitionPredicate == null) {
@@ -189,6 +222,18 @@ public final class PostponeMergeReadBuilder implements 
Serializable {
         return plan;
     }
 
+    /** Rebuilds only the routing metadata of an existing plan with a new 
default bucket number. */
+    public PostponeMergePlan reroute(PostponeMergePlan plan, int 
newDefaultBucketNum) {
+        checkArgument(newDefaultBucketNum > 0, "Default postpone bucket number 
must be positive.");
+        if (table.coreOptions()
+                .toConfiguration()
+                .contains(CoreOptions.POSTPONE_DEFAULT_BUCKET_NUM)) {
+            return plan;
+        }
+        defaultBucketNum = newDefaultBucketNum;
+        return plan.withDefaultBucketNum(defaultBucketNum);
+    }
+
     @Nullable
     public String readProtectionTagName() {
         return readProtectionTagName;
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java 
b/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
index 972713a5d0..e7dd850d63 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/source/TableScanTest.java
@@ -404,6 +404,9 @@ public class TableScanTest extends ScannerTestBase {
         // The option must not change an ordinary Core scan.
         assertThat(postponeTable.newScan().plan().splits()).hasSize(3);
         
assertThat(postponeTable.newScan().withLimit(1).plan().splits()).hasSize(1);
+        long realOnlySnapshotId = 
postponeTable.snapshotManager().latestSnapshotId();
+        PostponeMergeReadBuilder realOnlyBuilder =
+                PostponeMergeReadBuilder.createSnapshotBound(postponeTable, 
null).get();
         assertThat(PostponeMergeReadBuilder.create(postponeTable, 
null)).isEmpty();
 
         StreamTableWrite postponeWrite = postponeTable.newWrite(commitUser);
@@ -413,10 +416,18 @@ public class TableScanTest extends ScannerTestBase {
         postponeWrite.close();
         postponeCommit.close();
 
+        // A snapshot-bound builder must not pick up postpone files committed 
after its selection.
+        PostponeMergePlan realOnlyPlan = realOnlyBuilder.plan();
+        assertThat(realOnlyPlan.postponeSplits()).isEmpty();
+        assertThat(realOnlyPlan.realSplits())
+                .allSatisfy(split -> 
assertThat(split.snapshotId()).isEqualTo(realOnlySnapshotId));
+
         FileStoreTable compactedFullTable =
                 postponeTable.copy(
                         Collections.singletonMap(CoreOptions.SCAN_MODE.key(), 
"compacted-full"));
         assertThat(PostponeMergeReadBuilder.create(compactedFullTable, 
null)).isEmpty();
+        
assertThat(PostponeMergeReadBuilder.createSnapshotBound(compactedFullTable, 
null))
+                .isEmpty();
 
         // A postpone record may update any value, so real-file value 
statistics are unsafe.
         Predicate valueFilter = builder.equal(2, 100L);
@@ -562,15 +573,27 @@ public class TableScanTest extends ScannerTestBase {
         postponeWrite.close();
         postponeCommit.close();
 
-        PostponeMergePlan plan =
-                PostponeMergeReadBuilder.create(postponeTable, null)
-                        .get()
-                        .withDefaultBucketNum(4)
-                        .plan();
+        PostponeMergeReadBuilder readBuilder =
+                PostponeMergeReadBuilder.create(postponeTable, 
null).get().withDefaultBucketNum(1);
+        PostponeMergePlan initialPlan = readBuilder.plan();
+        assertThat(initialPlan.numPotentialBuckets()).isEqualTo(3);
+
+        PostponeMergePlan plan = readBuilder.reroute(initialPlan, 4);
 
         // Partition 1 uses its known bucket, partition 2 is real-only, and 
the new partition 3
         // may route to any of the four default buckets.
         assertThat(plan.numPotentialBuckets()).isEqualTo(6);
+
+        FileStoreTable explicitDefaultTable =
+                postponeTable.copy(
+                        Collections.singletonMap(
+                                CoreOptions.POSTPONE_DEFAULT_BUCKET_NUM.key(), 
"2"));
+        PostponeMergeReadBuilder explicitDefaultBuilder =
+                PostponeMergeReadBuilder.create(explicitDefaultTable, 
null).get();
+        PostponeMergePlan explicitDefaultPlan = explicitDefaultBuilder.plan();
+        assertThat(explicitDefaultPlan.numPotentialBuckets()).isEqualTo(4);
+        assertThat(explicitDefaultBuilder.reroute(explicitDefaultPlan, 
8).numPotentialBuckets())
+                .isEqualTo(4);
     }
 
     @Test
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/BaseDataTableSource.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/BaseDataTableSource.java
index a49f87c63f..58e5f38d4a 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/BaseDataTableSource.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/BaseDataTableSource.java
@@ -201,10 +201,12 @@ public abstract class BaseDataTableSource extends 
FlinkTableSource
         return new PaimonDataStreamScanProvider(
                 !unbounded,
                 env ->
-                        sourceBuilder
-                                .sourceParallelism(inferSourceParallelism(env))
-                                .env(env)
-                                .build(),
+                        PostponeMergeOnRead.usesCustomSource(table)
+                                ? sourceBuilder.env(env).build()
+                                : sourceBuilder
+                                        
.sourceParallelism(inferSourceParallelism(env))
+                                        .env(env)
+                                        .build(),
                 tableIdentifier.asSummaryString(),
                 table);
     }
@@ -247,6 +249,11 @@ public abstract class BaseDataTableSource extends 
FlinkTableSource
                             + table.getClass().getName());
         }
 
+        if (PostponeMergeOnRead.configured(table)) {
+            throw new UnsupportedOperationException(
+                    "Option 'postpone.merge-on-read' is not supported for 
lookup reads.");
+        }
+
         if (limit != null) {
             throw new RuntimeException(
                     "Limit push down should not happen in Lookup source, but 
it is " + limit);
@@ -334,6 +341,10 @@ public abstract class BaseDataTableSource extends 
FlinkTableSource
             return false;
         }
 
+        if (PostponeMergeOnRead.configured(table)) {
+            return false;
+        }
+
         if (!(table instanceof FileStoreTable)) {
             return false;
         }
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java
index 9271f2f5f9..b5d5636406 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/DataTableSource.java
@@ -150,7 +150,9 @@ public class DataTableSource extends BaseDataTableSource
     @Override
     public List<String> listAcceptedFilterFields() {
         // note that streaming query doesn't support dynamic filtering
-        return unbounded ? Collections.emptyList() : table.partitionKeys();
+        return unbounded || PostponeMergeOnRead.configured(table)
+                ? Collections.emptyList()
+                : table.partitionKeys();
     }
 
     @Override
@@ -160,6 +162,11 @@ public class DataTableSource extends BaseDataTableSource
                 "Cannot apply dynamic filtering to Paimon table '%s' when 
streaming reading.",
                 table.name());
 
+        checkState(
+                !PostponeMergeOnRead.configured(table),
+                "Cannot apply dynamic filtering to Paimon table '%s' when 
postpone merge-on-read is enabled.",
+                table.name());
+
         checkState(
                 !table.partitionKeys().isEmpty(),
                 "Cannot apply dynamic filtering to non-partitioned Paimon 
table '%s'.",
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 b80b52be77..a49ba7ae41 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
@@ -27,12 +27,16 @@ import org.apache.paimon.flink.sink.FlinkSink;
 import org.apache.paimon.flink.source.align.AlignedContinuousFileStoreSource;
 import org.apache.paimon.flink.source.operator.MonitorSource;
 import org.apache.paimon.flink.utils.TableScanUtils;
+import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.options.Options;
 import org.apache.paimon.partition.PartitionPredicate;
 import org.apache.paimon.predicate.Predicate;
 import org.apache.paimon.table.BucketMode;
 import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.Table;
+import org.apache.paimon.table.source.DataSplit;
+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.utils.StringUtils;
 
@@ -40,6 +44,7 @@ import 
org.apache.flink.api.common.eventtime.WatermarkStrategy;
 import org.apache.flink.api.common.functions.MapFunction;
 import org.apache.flink.api.common.typeinfo.TypeInformation;
 import org.apache.flink.api.connector.source.Source;
+import org.apache.flink.configuration.Configuration;
 import org.apache.flink.streaming.api.CheckpointingMode;
 import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.datastream.DataStreamSource;
@@ -57,6 +62,7 @@ import org.apache.flink.types.Row;
 
 import javax.annotation.Nullable;
 
+import java.util.Collections;
 import java.util.List;
 import java.util.Optional;
 
@@ -65,6 +71,7 @@ import static 
org.apache.paimon.flink.FlinkConnectorOptions.SOURCE_OPERATOR_UID_
 import static org.apache.paimon.flink.FlinkConnectorOptions.generateCustomUid;
 import static org.apache.paimon.flink.LogicalTypeConversion.toLogicalType;
 import static 
org.apache.paimon.flink.utils.ParallelismUtils.forwardParallelism;
+import static org.apache.paimon.options.OptionsUtils.PAIMON_PREFIX;
 import static org.apache.paimon.utils.Preconditions.checkArgument;
 import static org.apache.paimon.utils.Preconditions.checkState;
 
@@ -224,6 +231,107 @@ public class FlinkSourceBuilder {
                         skipPreloadTargetSnapshot));
     }
 
+    private @Nullable DataStream<RowData> buildPostponeMergeSource() {
+        FileStoreTable fileStoreTable = (FileStoreTable) table;
+        if (fileStoreTable.coreOptions().startupMode() == 
CoreOptions.StartupMode.COMPACTED_FULL) {
+            return null;
+        }
+
+        Optional<PostponeMergeReadBuilder> optionalBuilder =
+                PostponeMergeReadBuilder.createSnapshotBound(fileStoreTable, 
partitionPredicate);
+        if (!optionalBuilder.isPresent()) {
+            return env.fromCollection(Collections.<RowData>emptyList(), 
produceTypeInfo())
+                    .name(sourceName)
+                    .forceNonParallel();
+        }
+
+        int baseParallelism = basePostponeMergeParallelism();
+        PostponeMergeReadBuilder readBuilder =
+                optionalBuilder.get().withDefaultBucketNum(baseParallelism);
+        org.apache.paimon.types.RowType readType = projectedRowType();
+        if (readType != null) {
+            readBuilder.withReadType(readType);
+        }
+        if (predicate != null) {
+            readBuilder.withFilter(predicate);
+        }
+        PostponeMergePlan plan = readBuilder.plan();
+        int mergeParallelism = inferPostponeMergeParallelism(plan, 
baseParallelism);
+        if (mergeParallelism != baseParallelism) {
+            plan = readBuilder.reroute(plan, mergeParallelism);
+        }
+        return PostponeMergeOnRead.build(
+                env,
+                sourceName,
+                fileStoreTable,
+                readBuilder,
+                plan,
+                mergeParallelism,
+                produceTypeInfo(),
+                outerProject(),
+                limit,
+                conf.get(CoreOptions.BLOB_AS_DESCRIPTOR));
+    }
+
+    private int basePostponeMergeParallelism() {
+        if (parallelism != null) {
+            return parallelism;
+        }
+
+        Integer configuredParallelism = 
conf.get(FlinkConnectorOptions.SCAN_PARALLELISM);
+        if (configuredParallelism != null) {
+            return configuredParallelism;
+        }
+        return Math.max(1, env.getParallelism());
+    }
+
+    private int inferPostponeMergeParallelism(PostponeMergePlan plan, int 
baseParallelism) {
+        if (parallelism != null
+                || conf.get(FlinkConnectorOptions.SCAN_PARALLELISM) != null
+                || env.getParallelism() > 0
+                || !inferScanParallelism()) {
+            return baseParallelism;
+        }
+
+        long inferred =
+                Math.max(
+                        plan.numPotentialBuckets(),
+                        (long) plan.realSplits().size() + 
plan.postponeSplits().size());
+        long totalFileSize = 0L;
+        for (org.apache.paimon.table.source.Split split : plan.splits()) {
+            for (DataFileMeta file : ((DataSplit) split).dataFiles()) {
+                if (file.fileSize() > Long.MAX_VALUE - totalFileSize) {
+                    totalFileSize = Long.MAX_VALUE;
+                    break;
+                }
+                totalFileSize += file.fileSize();
+            }
+            if (totalFileSize == Long.MAX_VALUE) {
+                break;
+            }
+        }
+        long splitTargetSize = ((FileStoreTable) 
table).coreOptions().splitTargetSize();
+        long sizeParallelism = totalFileSize == 0 ? 1 : (totalFileSize - 1) / 
splitTargetSize + 1;
+        inferred = Math.max(inferred, sizeParallelism);
+        if (limit != null && limit > 0) {
+            inferred = Math.min(inferred, limit);
+        }
+        inferred =
+                Math.min(
+                        Math.max(1L, inferred),
+                        
conf.get(FlinkConnectorOptions.INFER_SCAN_MAX_PARALLELISM));
+        return (int) inferred;
+    }
+
+    private boolean inferScanParallelism() {
+        Configuration envConfig = (Configuration) env.getConfiguration();
+        String flinkOption = PAIMON_PREFIX + 
FlinkConnectorOptions.INFER_SCAN_PARALLELISM.key();
+        if (envConfig.containsKey(flinkOption)) {
+            return Boolean.parseBoolean(envConfig.toMap().get(flinkOption));
+        }
+        return conf.get(FlinkConnectorOptions.INFER_SCAN_PARALLELISM);
+    }
+
     private DataStream<RowData> buildContinuousFileSource() {
         return toDataStream(
                 new ContinuousFileStoreSource(
@@ -325,6 +433,17 @@ public class FlinkSourceBuilder {
                             + " too many snapshots that could pose a risk to 
the file system.");
         }
 
+        if (PostponeMergeOnRead.configured(table)) {
+            if (!sourceBounded) {
+                throw new UnsupportedOperationException(
+                        "Option 'postpone.merge-on-read' is only supported for 
batch reads.");
+            }
+            DataStream<RowData> mergeSource = buildPostponeMergeSource();
+            if (mergeSource != null) {
+                return mergeSource;
+            }
+        }
+
         if (sourceBounded) {
             if 
(conf.get(FlinkConnectorOptions.SCAN_DEDICATED_SPLIT_GENERATION)) {
                 return buildDedicatedSplitGenSource(true);
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeInputOperator.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeInputOperator.java
new file mode 100644
index 0000000000..43c9aefb14
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeInputOperator.java
@@ -0,0 +1,145 @@
+/*
+ * 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.KeyValue;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.serializer.InternalRowSerializer;
+import org.apache.paimon.disk.IOManager;
+import org.apache.paimon.reader.RecordReaderIterator;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.PostponeUtils.PostponeBucketRouter;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.PostponeMergeRead;
+import org.apache.paimon.table.source.PostponeMergeReadBuilder;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.table.source.SplitSerializer;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.SerializationUtils;
+
+import org.apache.flink.streaming.api.operators.AbstractStreamOperator;
+import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+
+import static 
org.apache.paimon.flink.source.PostponeMergeOnRead.POSTPONE_RECORD_KIND;
+import static 
org.apache.paimon.flink.source.PostponeMergeOnRead.REAL_SPLIT_KIND;
+
+/** Reads postpone splits and emits serializable records carrying their target 
buckets. */
+final class PostponeMergeInputOperator extends 
AbstractStreamOperator<InternalRow>
+        implements OneInputStreamOperator<Split, InternalRow> {
+
+    private static final long serialVersionUID = 1L;
+
+    private final PostponeMergeReadBuilder readBuilder;
+    private final RowType keyType;
+    private final RowType mergeReadType;
+    private final PostponeBucketRouter bucketRouter;
+
+    private transient IOManager ioManager;
+    private transient PostponeMergeRead read;
+    private transient InternalRowSerializer keySerializer;
+    private transient InternalRowSerializer valueSerializer;
+
+    PostponeMergeInputOperator(
+            PostponeMergeReadBuilder readBuilder,
+            RowType keyType,
+            RowType mergeReadType,
+            PostponeBucketRouter bucketRouter) {
+        this.readBuilder = readBuilder;
+        this.keyType = keyType;
+        this.mergeReadType = mergeReadType;
+        this.bucketRouter = bucketRouter;
+    }
+
+    @Override
+    public void open() throws Exception {
+        super.open();
+        ioManager =
+                IOManager.create(
+                        getContainingTask()
+                                .getEnvironment()
+                                .getIOManager()
+                                .getSpillingDirectoriesPaths());
+        read = readBuilder.newRead().withIOManager(ioManager);
+        keySerializer = new InternalRowSerializer(keyType);
+        valueSerializer = new InternalRowSerializer(mergeReadType);
+    }
+
+    @Override
+    public void processElement(StreamRecord<Split> element) throws Exception {
+        DataSplit split = (DataSplit) element.getValue();
+        if (split.bucket() == BucketMode.POSTPONE_BUCKET) {
+            emitPostponeRecords(split);
+        } else {
+            emitRealSplit(split);
+        }
+    }
+
+    private void emitRealSplit(DataSplit split) throws Exception {
+        GenericRow carrier =
+                GenericRow.of(
+                        
SerializationUtils.serializeBinaryRow(split.partition()),
+                        split.bucket(),
+                        REAL_SPLIT_KIND,
+                        SplitSerializer.serialize(split),
+                        null,
+                        0L,
+                        (byte) 0,
+                        null);
+        output.collect(new StreamRecord<>(carrier));
+    }
+
+    private void emitPostponeRecords(DataSplit split) throws Exception {
+        byte[] partition = 
SerializationUtils.serializeBinaryRow(split.partition());
+        long writerLocalOrder = 0L;
+        try (RecordReaderIterator<KeyValue> records =
+                new RecordReaderIterator<>(read.createPostponeReader(split))) {
+            while (records.hasNext()) {
+                KeyValue keyValue = records.next();
+                BinaryRow key = keySerializer.toBinaryRow(keyValue.key());
+                BinaryRow value = 
valueSerializer.toBinaryRow(keyValue.value());
+                GenericRow carrier =
+                        GenericRow.of(
+                                partition,
+                                bucketRouter.bucket(split.partition(), key),
+                                POSTPONE_RECORD_KIND,
+                                null,
+                                SerializationUtils.serializeBinaryRow(key),
+                                writerLocalOrder,
+                                keyValue.valueKind().toByteValue(),
+                                SerializationUtils.serializeBinaryRow(value));
+                output.collect(new StreamRecord<>(carrier));
+                writerLocalOrder = Math.addExact(writerLocalOrder, 1L);
+            }
+        }
+    }
+
+    @Override
+    public void close() throws Exception {
+        try {
+            if (ioManager != null) {
+                ioManager.close();
+            }
+        } finally {
+            super.close();
+        }
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeInputSource.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeInputSource.java
new file mode 100644
index 0000000000..092af16778
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeInputSource.java
@@ -0,0 +1,127 @@
+/*
+ * 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.data.BinaryRow;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.DeletionFile;
+import org.apache.paimon.table.source.PostponeMergePlan;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.utils.Pair;
+
+import org.apache.flink.api.connector.source.Boundedness;
+import org.apache.flink.api.connector.source.ReaderOutput;
+import org.apache.flink.api.connector.source.SourceReader;
+import org.apache.flink.api.connector.source.SourceReaderContext;
+import org.apache.flink.core.io.InputStatus;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+/** Emits real-bucket markers and writer-grouped postpone splits for 
merge-on-read. */
+final class PostponeMergeInputSource extends 
AbstractNonCoordinatedSource<Split> {
+
+    private static final long serialVersionUID = 1L;
+
+    private final List<Split> inputs;
+
+    PostponeMergeInputSource(PostponeMergePlan plan) {
+        this.inputs = new ArrayList<>();
+
+        Map<Pair<BinaryRow, Integer>, List<DataSplit>> realBuckets = new 
LinkedHashMap<>();
+        for (DataSplit split : plan.realSplits()) {
+            realBuckets
+                    .computeIfAbsent(
+                            Pair.of(split.partition().copy(), split.bucket()),
+                            ignored -> new ArrayList<>())
+                    .add(split);
+        }
+        for (List<DataSplit> splits : realBuckets.values()) {
+            inputs.add(mergeRealSplits(splits));
+        }
+        inputs.addAll(plan.postponeSplits());
+    }
+
+    @Override
+    public Boundedness getBoundedness() {
+        return Boundedness.BOUNDED;
+    }
+
+    @Override
+    public SourceReader<Split, SimpleSourceSplit> 
createReader(SourceReaderContext readerContext) {
+        return new Reader();
+    }
+
+    private class Reader extends AbstractNonCoordinatedSourceReader<Split> {
+
+        private boolean emitted;
+
+        @Override
+        public InputStatus pollNext(ReaderOutput<Split> output) {
+            if (!emitted) {
+                for (Split input : inputs) {
+                    output.collect(input);
+                }
+                emitted = true;
+            }
+            return InputStatus.END_OF_INPUT;
+        }
+    }
+
+    private static DataSplit mergeRealSplits(List<DataSplit> splits) {
+        if (splits.size() == 1) {
+            return splits.get(0);
+        }
+
+        DataSplit first = splits.get(0);
+        List<DataFileMeta> dataFiles = new ArrayList<>();
+        List<DeletionFile> deletionFiles = new ArrayList<>();
+        boolean hasDeletionFiles = false;
+        boolean rawConvertible = true;
+        for (DataSplit split : splits) {
+            dataFiles.addAll(split.dataFiles());
+            if (split.deletionFiles().isPresent()) {
+                deletionFiles.addAll(split.deletionFiles().get());
+                hasDeletionFiles = true;
+            } else {
+                
deletionFiles.addAll(Collections.nCopies(split.dataFiles().size(), null));
+            }
+            rawConvertible &= split.rawConvertible();
+        }
+
+        DataSplit.Builder builder =
+                DataSplit.builder()
+                        .withSnapshot(first.snapshotId())
+                        .withPartition(first.partition())
+                        .withBucket(first.bucket())
+                        .withBucketPath(first.bucketPath())
+                        .withTotalBuckets(first.totalBuckets())
+                        .withDataFiles(dataFiles)
+                        .isStreaming(first.isStreaming())
+                        .rawConvertible(rawConvertible);
+        if (hasDeletionFiles) {
+            builder.withDataDeletionFiles(deletionFiles);
+        }
+        return builder.build();
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeOnRead.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeOnRead.java
new file mode 100644
index 0000000000..b0054e066d
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeOnRead.java
@@ -0,0 +1,160 @@
+/*
+ * 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.CoreOptions;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.flink.NestedProjectedRowData;
+import org.apache.paimon.flink.sink.FlinkStreamPartitioner;
+import org.apache.paimon.flink.utils.InternalTypeInfo;
+import org.apache.paimon.flink.utils.JavaTypeInfo;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.table.sink.ChannelComputer;
+import org.apache.paimon.table.source.PostponeMergePlan;
+import org.apache.paimon.table.source.PostponeMergeReadBuilder;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.SerializationUtils;
+
+import org.apache.flink.api.common.eventtime.WatermarkStrategy;
+import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.table.data.RowData;
+
+import javax.annotation.Nullable;
+
+/** Builds the Flink batch topology for postpone merge-on-read. */
+final class PostponeMergeOnRead {
+
+    static final int PARTITION = 0;
+    static final int BUCKET = 1;
+    static final int INPUT_KIND = 2;
+    static final int REAL_SPLIT = 3;
+    static final int KEY = 4;
+    static final int WRITER_LOCAL_ORDER = 5;
+    static final int ROW_KIND = 6;
+    static final int VALUE = 7;
+
+    static final byte REAL_SPLIT_KIND = 0;
+    static final byte POSTPONE_RECORD_KIND = 1;
+
+    static final RowType CARRIER_TYPE =
+            RowType.of(
+                    DataTypes.BYTES().notNull(),
+                    DataTypes.INT().notNull(),
+                    DataTypes.TINYINT().notNull(),
+                    DataTypes.BYTES(),
+                    DataTypes.BYTES(),
+                    DataTypes.BIGINT().notNull(),
+                    DataTypes.TINYINT().notNull(),
+                    DataTypes.BYTES());
+
+    private PostponeMergeOnRead() {}
+
+    static boolean configured(Table table) {
+        if (!(table instanceof FileStoreTable)) {
+            return false;
+        }
+        FileStoreTable fileStoreTable = (FileStoreTable) table;
+        return fileStoreTable.coreOptions().postponeMergeOnRead()
+                && fileStoreTable.bucketMode() == BucketMode.POSTPONE_MODE
+                && !fileStoreTable.primaryKeys().isEmpty();
+    }
+
+    static boolean usesCustomSource(Table table) {
+        return configured(table)
+                && ((FileStoreTable) table).coreOptions().startupMode()
+                        != CoreOptions.StartupMode.COMPACTED_FULL;
+    }
+
+    static DataStream<RowData> build(
+            StreamExecutionEnvironment env,
+            String sourceName,
+            FileStoreTable table,
+            PostponeMergeReadBuilder readBuilder,
+            PostponeMergePlan plan,
+            int parallelism,
+            TypeInformation<RowData> outputType,
+            @Nullable NestedProjectedRowData outerProject,
+            @Nullable Long limit,
+            boolean blobAsDescriptor) {
+        PostponeMergeInputSource inputSource = new 
PostponeMergeInputSource(plan);
+        DataStream<Split> inputs =
+                env.fromSource(
+                                new PaimonDataStreamSource<>(inputSource, 
table),
+                                WatermarkStrategy.noWatermarks(),
+                                sourceName + " - Postpone merge inputs",
+                                new JavaTypeInfo<>(Split.class))
+                        .forceNonParallel();
+
+        DataStream<Split> distributedInputs = 
FlinkStreamPartitioner.rebalance(inputs, parallelism);
+        SingleOutputStreamOperator<InternalRow> carriers =
+                distributedInputs
+                        .transform(
+                                "Postpone merge input",
+                                InternalTypeInfo.fromRowType(CARRIER_TYPE),
+                                new PostponeMergeInputOperator(
+                                        readBuilder,
+                                        plan.keyType(),
+                                        plan.mergeReadType(),
+                                        plan.bucketRouter()))
+                        .setParallelism(parallelism);
+
+        DataStream<InternalRow> partitionedCarriers =
+                FlinkStreamPartitioner.partition(
+                        carriers, new CarrierChannelComputer(), parallelism);
+        return partitionedCarriers
+                .transform(
+                        "Postpone merge on read",
+                        outputType,
+                        new PostponeMergeOperator(
+                                readBuilder,
+                                plan.resultReadType(),
+                                table.options(),
+                                outerProject,
+                                limit,
+                                blobAsDescriptor))
+                .setParallelism(parallelism);
+    }
+
+    private static class CarrierChannelComputer implements 
ChannelComputer<InternalRow> {
+
+        private static final long serialVersionUID = 1L;
+
+        private int numChannels;
+
+        @Override
+        public void setup(int numChannels) {
+            this.numChannels = numChannels;
+        }
+
+        @Override
+        public int channel(InternalRow carrier) {
+            BinaryRow partition =
+                    
SerializationUtils.deserializeBinaryRow(carrier.getBinary(PARTITION));
+            return ChannelComputer.select(partition, carrier.getInt(BUCKET), 
numChannels);
+        }
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeOperator.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeOperator.java
new file mode 100644
index 0000000000..5f4fc25098
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/source/PostponeMergeOperator.java
@@ -0,0 +1,279 @@
+/*
+ * 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.CoreOptions;
+import org.apache.paimon.KeyValue;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.disk.IOManager;
+import org.apache.paimon.flink.FlinkRowData;
+import org.apache.paimon.flink.FlinkRowDataWithBlob;
+import org.apache.paimon.flink.FlinkRowWrapper;
+import org.apache.paimon.flink.NestedProjectedRowData;
+import org.apache.paimon.reader.RecordReaderIterator;
+import org.apache.paimon.sort.BinaryExternalSortBuffer;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.PostponeMergeRead;
+import org.apache.paimon.table.source.PostponeMergeReadBuilder;
+import org.apache.paimon.table.source.SplitSerializer;
+import org.apache.paimon.types.RowKind;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.IteratorRecordReader;
+import org.apache.paimon.utils.MutableObjectIterator;
+import org.apache.paimon.utils.SerializationUtils;
+
+import org.apache.flink.streaming.api.operators.AbstractStreamOperator;
+import org.apache.flink.streaming.api.operators.BoundedOneInput;
+import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.table.data.RowData;
+
+import javax.annotation.Nullable;
+
+import java.util.Arrays;
+import java.util.Iterator;
+import java.util.Map;
+import java.util.NoSuchElementException;
+import java.util.Set;
+
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.BUCKET;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.CARRIER_TYPE;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.INPUT_KIND;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.KEY;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.PARTITION;
+import static 
org.apache.paimon.flink.source.PostponeMergeOnRead.POSTPONE_RECORD_KIND;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.REAL_SPLIT;
+import static 
org.apache.paimon.flink.source.PostponeMergeOnRead.REAL_SPLIT_KIND;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.ROW_KIND;
+import static org.apache.paimon.flink.source.PostponeMergeOnRead.VALUE;
+import static 
org.apache.paimon.flink.source.PostponeMergeOnRead.WRITER_LOCAL_ORDER;
+
+/** Spillably sorts routed carriers and merges each target bucket through 
Paimon Core. */
+final class PostponeMergeOperator extends AbstractStreamOperator<RowData>
+        implements OneInputStreamOperator<InternalRow, RowData>, 
BoundedOneInput {
+
+    private static final long serialVersionUID = 1L;
+
+    private final PostponeMergeReadBuilder readBuilder;
+    private final RowType resultReadType;
+    private final Map<String, String> tableOptions;
+    @Nullable private final NestedProjectedRowData outerProject;
+    @Nullable private final Long limit;
+    private final boolean blobAsDescriptor;
+
+    private transient IOManager ioManager;
+    private transient BinaryExternalSortBuffer buffer;
+    private transient PostponeMergeRead read;
+    private transient FlinkRowData flinkRow;
+    @Nullable private transient NestedProjectedRowData projectedRow;
+    @Nullable private transient RecordLimiter recordLimiter;
+
+    PostponeMergeOperator(
+            PostponeMergeReadBuilder readBuilder,
+            RowType resultReadType,
+            Map<String, String> tableOptions,
+            @Nullable NestedProjectedRowData outerProject,
+            @Nullable Long limit,
+            boolean blobAsDescriptor) {
+        this.readBuilder = readBuilder;
+        this.resultReadType = resultReadType;
+        this.tableOptions = tableOptions;
+        this.outerProject = outerProject;
+        this.limit = limit;
+        this.blobAsDescriptor = blobAsDescriptor;
+    }
+
+    @Override
+    public void open() throws Exception {
+        super.open();
+        ioManager =
+                IOManager.create(
+                        getContainingTask()
+                                .getEnvironment()
+                                .getIOManager()
+                                .getSpillingDirectoriesPaths());
+        CoreOptions options = CoreOptions.fromMap(tableOptions);
+        buffer =
+                BinaryExternalSortBuffer.create(
+                        ioManager,
+                        CARRIER_TYPE,
+                        new int[] {PARTITION, BUCKET, INPUT_KIND, 
WRITER_LOCAL_ORDER},
+                        options.sortSpillBufferSize(),
+                        options.pageSize(),
+                        options.localSortMaxNumFileHandles(),
+                        options.spillCompressOptions(),
+                        options.writeBufferSpillDiskSize());
+        read = readBuilder.newRead().withIOManager(ioManager);
+
+        Set<Integer> blobFields = 
FlinkRowWrapper.blobFieldIndexes(resultReadType);
+        flinkRow =
+                blobFields.isEmpty()
+                        ? new FlinkRowData(null)
+                        : new FlinkRowDataWithBlob(null, blobFields, 
blobAsDescriptor);
+        projectedRow = NestedProjectedRowData.copy(outerProject);
+        recordLimiter = RecordLimiter.create(limit);
+    }
+
+    @Override
+    public void processElement(StreamRecord<InternalRow> element) throws 
Exception {
+        buffer.write(element.getValue());
+    }
+
+    @Override
+    public void endInput() throws Exception {
+        if (buffer.isEmpty()) {
+            return;
+        }
+
+        SortedCarriers carriers = new SortedCarriers(buffer.sortedIterator());
+        while (carriers.hasNext() && !reachLimit()) {
+            mergeBucket(carriers);
+        }
+    }
+
+    private void mergeBucket(SortedCarriers carriers) throws Exception {
+        InternalRow first = carriers.peek();
+        BucketKey bucketKey = new BucketKey(first.getBinary(PARTITION), 
first.getInt(BUCKET));
+
+        DataSplit realSplit = null;
+        if (sameBucket(first, bucketKey) && first.getByte(INPUT_KIND) == 
REAL_SPLIT_KIND) {
+            realSplit =
+                    (DataSplit) 
SplitSerializer.deserialize(carriers.next().getBinary(REAL_SPLIT));
+        }
+
+        Iterator<KeyValue> postponeRecords =
+                new Iterator<KeyValue>() {
+                    @Override
+                    public boolean hasNext() {
+                        return carriers.hasNext()
+                                && sameBucket(carriers.peek(), bucketKey)
+                                && carriers.peek().getByte(INPUT_KIND) == 
POSTPONE_RECORD_KIND;
+                    }
+
+                    @Override
+                    public KeyValue next() {
+                        if (!hasNext()) {
+                            throw new NoSuchElementException();
+                        }
+                        InternalRow carrier = carriers.next();
+                        return new KeyValue()
+                                .replace(
+                                        
SerializationUtils.deserializeBinaryRow(
+                                                carrier.getBinary(KEY)),
+                                        
RowKind.fromByteValue(carrier.getByte(ROW_KIND)),
+                                        
SerializationUtils.deserializeBinaryRow(
+                                                carrier.getBinary(VALUE)));
+                    }
+                };
+
+        try (RecordReaderIterator<InternalRow> rows =
+                new RecordReaderIterator<>(
+                        read.createBucketMergeReader(
+                                realSplit, new 
IteratorRecordReader<>(postponeRecords)))) {
+            if (carriers.hasNext() && sameBucket(carriers.peek(), bucketKey)) {
+                throw new IllegalStateException(
+                        "Unexpected postpone merge carrier kind "
+                                + carriers.peek().getByte(INPUT_KIND)
+                                + " for one bucket.");
+            }
+            while (rows.hasNext() && !reachLimit()) {
+                flinkRow.replace(rows.next());
+                RowData result =
+                        projectedRow == null ? flinkRow : 
projectedRow.replaceRow(flinkRow);
+                output.collect(new StreamRecord<>(result));
+                if (recordLimiter != null) {
+                    recordLimiter.increment();
+                }
+            }
+        }
+    }
+
+    private boolean reachLimit() {
+        return recordLimiter != null && recordLimiter.reachLimit();
+    }
+
+    private static boolean sameBucket(InternalRow carrier, BucketKey 
bucketKey) {
+        return carrier.getInt(BUCKET) == bucketKey.bucket
+                && Arrays.equals(carrier.getBinary(PARTITION), 
bucketKey.partition);
+    }
+
+    @Override
+    public void close() throws Exception {
+        try {
+            if (buffer != null) {
+                buffer.clear();
+            }
+        } finally {
+            try {
+                if (ioManager != null) {
+                    ioManager.close();
+                }
+            } finally {
+                super.close();
+            }
+        }
+    }
+
+    private static class BucketKey {
+
+        private final byte[] partition;
+        private final int bucket;
+
+        private BucketKey(byte[] partition, int bucket) {
+            this.partition = partition;
+            this.bucket = bucket;
+        }
+    }
+
+    private static class SortedCarriers {
+
+        private final MutableObjectIterator<BinaryRow> iterator;
+        private final BinaryRow reuse = new 
BinaryRow(CARRIER_TYPE.getFieldCount());
+        @Nullable private BinaryRow next;
+
+        private SortedCarriers(MutableObjectIterator<BinaryRow> iterator) {
+            this.iterator = iterator;
+        }
+
+        private boolean hasNext() {
+            if (next == null) {
+                try {
+                    next = iterator.next(reuse);
+                } catch (Exception e) {
+                    throw new RuntimeException("Failed to read sorted postpone 
carriers.", e);
+                }
+            }
+            return next != null;
+        }
+
+        private BinaryRow peek() {
+            if (!hasNext()) {
+                throw new NoSuchElementException();
+            }
+            return next;
+        }
+
+        private BinaryRow next() {
+            BinaryRow result = peek();
+            next = null;
+            return result;
+        }
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PostponeBucketTableITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PostponeBucketTableITCase.java
index d9f8781da3..d27f6dd615 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PostponeBucketTableITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/PostponeBucketTableITCase.java
@@ -18,7 +18,15 @@
 
 package org.apache.paimon.flink;
 
+import org.apache.paimon.catalog.Catalog;
+import org.apache.paimon.catalog.CatalogContext;
+import org.apache.paimon.catalog.CatalogFactory;
+import org.apache.paimon.catalog.Identifier;
 import org.apache.paimon.flink.util.AbstractTestBase;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.index.DeletionVectorMeta;
+import org.apache.paimon.manifest.IndexManifestEntry;
+import org.apache.paimon.table.FileStoreTable;
 
 import org.apache.paimon.shade.guava30.com.google.common.collect.ImmutableList;
 
@@ -41,6 +49,7 @@ import java.util.concurrent.ThreadLocalRandom;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
 
+import static 
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX;
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatCode;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -180,6 +189,70 @@ public class PostponeBucketTableITCase extends 
AbstractTestBase {
         
assertThat(collect(tEnv.executeSql(query))).hasSameElementsAs(expected);
     }
 
+    @Test
+    public void testMergeOnRead() throws Exception {
+        String warehouse = getTempDirPath();
+        TableEnvironment tEnv =
+                tableEnvironmentBuilder()
+                        .batchMode()
+                        .parallelism(3)
+                        .setConf(TableConfigOptions.TABLE_DML_SYNC, true)
+                        .build();
+
+        tEnv.executeSql(
+                "CREATE CATALOG mycat WITH (\n"
+                        + "  'type' = 'paimon',\n"
+                        + "  'warehouse' = '"
+                        + warehouse
+                        + "'\n"
+                        + ")");
+        tEnv.executeSql("USE CATALOG mycat");
+        tEnv.executeSql(
+                "CREATE TABLE T (\n"
+                        + "  pt INT,\n"
+                        + "  k INT,\n"
+                        + "  v STRING,\n"
+                        + "  seq INT,\n"
+                        + "  PRIMARY KEY (pt, k) NOT ENFORCED\n"
+                        + ") PARTITIONED BY (pt) WITH (\n"
+                        + "  'bucket' = '-2',\n"
+                        + "  'sequence.field' = 'seq',\n"
+                        + "  'postpone.default-bucket-num' = '4',\n"
+                        + "  'source.split.target-size' = '1 B'\n"
+                        + ")");
+
+        // Create multiple real-bucket splits.
+        tEnv.executeSql("INSERT INTO T VALUES (0, 1, 'base-1', 10)").await();
+        tEnv.executeSql("INSERT INTO T VALUES (0, 2, 'base-2', 10)").await();
+        tEnv.executeSql(
+                        "INSERT INTO T /*+ 
OPTIONS('postpone.batch-write-fixed-bucket' = 'false') */ "
+                                + "VALUES (0, 1, 'new-1', 20), "
+                                + "(0, 1, 'newer-1', 20), "
+                                + "(0, 2, 'older-2', 5), (1, 3, 'new-3', 1)")
+                .await();
+
+        // MOR is opt-in; postpone records remain hidden by default.
+        assertThat(collect(tEnv.executeSql("SELECT * FROM T")))
+                .containsExactlyInAnyOrder("+I[0, 1, base-1, 10]", "+I[0, 2, 
base-2, 10]");
+
+        String hint = "/*+ OPTIONS('postpone.merge-on-read' = 'true') */";
+        assertThat(collect(tEnv.executeSql("SELECT * FROM T " + hint)))
+                .containsExactlyInAnyOrder(
+                        "+I[0, 1, newer-1, 20]", "+I[0, 2, base-2, 10]", 
"+I[1, 3, new-3, 1]");
+        assertThat(collect(tEnv.executeSql("SELECT v FROM T " + hint + " WHERE 
pt = 0 AND k = 2")))
+                .containsExactly("+I[base-2]");
+        assertThat(collect(tEnv.executeSql("SELECT * FROM T " + hint + " WHERE 
v = 'base-1'")))
+                .isEmpty();
+        assertThat(collect(tEnv.executeSql("SELECT COUNT(*) FROM T " + hint)))
+                .containsExactly("+I[3]");
+
+        // MOR remains correct after postpone files are compacted into real 
buckets.
+        tEnv.executeSql("CALL sys.compact(`table` => 'default.T')").await();
+        assertThat(collect(tEnv.executeSql("SELECT * FROM T " + hint)))
+                .containsExactlyInAnyOrder(
+                        "+I[0, 1, newer-1, 20]", "+I[0, 2, base-2, 10]", 
"+I[1, 3, new-3, 1]");
+    }
+
     @Test
     public void testOverwriteWithoutBatchWriteFixedBucket() throws Exception {
         String warehouse = getTempDirPath();
@@ -748,7 +821,8 @@ public class PostponeBucketTableITCase extends 
AbstractTestBase {
                         + ") WITH (\n"
                         + "  'bucket' = '-2',\n"
                         + "  'postpone.batch-write-fixed-bucket' = 'false',\n"
-                        + "  'deletion-vectors.enabled' = 'true'\n"
+                        + "  'deletion-vectors.enabled' = 'true',\n"
+                        + "  'deletion-vectors.merge-on-read' = 'true'\n"
                         + ")");
 
         tEnv.executeSql("INSERT INTO T VALUES (1, 10), (2, 20), (3, 30), (4, 
40)").await();
@@ -761,14 +835,41 @@ public class PostponeBucketTableITCase extends 
AbstractTestBase {
         assertThat(collect(tEnv.executeSql("SELECT * FROM T")))
                 .containsExactlyInAnyOrder(
                         "+I[1, 11]", "+I[2, 20]", "+I[3, 30]", "+I[4, 40]", 
"+I[5, 51]");
+        assertThat(deletionVectorCardinality(warehouse)).isPositive();
 
         tEnv.executeSql("INSERT INTO T VALUES (2, 52), (3, 32)").await();
+        assertThat(collect(tEnv.executeSql("SELECT * FROM T")))
+                .containsExactlyInAnyOrder(
+                        "+I[1, 11]", "+I[2, 20]", "+I[3, 30]", "+I[4, 40]", 
"+I[5, 51]");
+        String mergeOnReadHint =
+                "/*+ OPTIONS("
+                        + "'postpone.merge-on-read' = 'true', "
+                        + "'deletion-vectors.merge-on-read' = 'true') */";
+        assertThat(collect(tEnv.executeSql("SELECT * FROM T " + 
mergeOnReadHint)))
+                .containsExactlyInAnyOrder(
+                        "+I[1, 11]", "+I[2, 52]", "+I[3, 32]", "+I[4, 40]", 
"+I[5, 51]");
+
         tEnv.executeSql("CALL sys.compact(`table` => 'default.T')").await();
         assertThat(collect(tEnv.executeSql("SELECT * FROM T")))
                 .containsExactlyInAnyOrder(
                         "+I[1, 11]", "+I[2, 52]", "+I[3, 32]", "+I[4, 40]", 
"+I[5, 51]");
     }
 
+    private long deletionVectorCardinality(String warehouse) throws Exception {
+        try (Catalog catalog =
+                CatalogFactory.createCatalog(CatalogContext.create(new 
Path(warehouse)))) {
+            FileStoreTable table =
+                    (FileStoreTable) 
catalog.getTable(Identifier.create("default", "T"));
+            return table.store().newIndexFileHandler()
+                    .scan(table.latestSnapshot().get(), 
DELETION_VECTORS_INDEX).stream()
+                    .map(IndexManifestEntry::indexFile)
+                    .filter(index -> index.dvRanges() != null)
+                    .flatMap(index -> index.dvRanges().values().stream())
+                    .mapToLong(DeletionVectorMeta::cardinality)
+                    .sum();
+        }
+    }
+
     @Test
     public void testSameKeyPreserveOrder() throws Exception {
         String warehouse = getTempDirPath();
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 e9b782f660..296d825a29 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
@@ -45,6 +45,7 @@ import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.apache.flink.table.catalog.ObjectIdentifier;
 import org.apache.flink.table.connector.source.DynamicTableSource;
+import org.apache.flink.table.connector.source.LookupTableSource;
 import org.apache.flink.table.connector.source.ScanTableSource;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.types.DataType;
@@ -59,6 +60,7 @@ import java.util.Optional;
 
 import static org.apache.paimon.options.OptionsUtils.PAIMON_PREFIX;
 import static org.assertj.core.api.AssertionsForClassTypes.assertThat;
+import static org.assertj.core.api.AssertionsForClassTypes.assertThatThrownBy;
 
 /** Tests for {@link DataTableSource}. */
 class DataTableSourceTest {
@@ -99,6 +101,40 @@ class DataTableSourceTest {
         
assertThat(sourceStream2.getParallelism()).isEqualTo(sEnv2.getParallelism());
     }
 
+    @Test
+    void testInferPostponeMergeParallelism() throws Exception {
+        FileStoreTable fileStoreTable = createPostponeTable(false);
+        writePostponeData(fileStoreTable);
+
+        DataTableSource tableSource =
+                new DataTableSource(
+                        ObjectIdentifier.of("cat", "db", "table"), 
fileStoreTable, false, null);
+        PaimonDataStreamScanProvider runtimeProvider = 
runtimeProvider(tableSource);
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.createLocalEnvironment();
+        env.setParallelism(-1);
+
+        DataStream<RowData> sourceStream =
+                runtimeProvider.produceDataStream(s -> Optional.empty(), env);
+
+        assertThat(sourceStream.getParallelism()).isEqualTo(3);
+    }
+
+    @Test
+    void testPostponeMergeRejectsLookupAndDynamicFiltering() throws Exception {
+        FileStoreTable fileStoreTable = createPostponeTable(true);
+        DataTableSource tableSource =
+                new DataTableSource(
+                        ObjectIdentifier.of("cat", "db", "table"), 
fileStoreTable, false, null);
+
+        
assertThat(tableSource.listAcceptedFilterFields()).isEqualTo(Collections.emptyList());
+        assertThatThrownBy(() -> 
tableSource.applyDynamicFiltering(Collections.singletonList("pt")))
+                .isInstanceOf(IllegalStateException.class)
+                .hasMessageContaining("postpone merge-on-read");
+        assertThatThrownBy(() -> 
tableSource.getLookupRuntimeProvider(lookupContext()))
+                .isInstanceOf(UnsupportedOperationException.class)
+                .hasMessageContaining("not supported for lookup reads");
+    }
+
     @Test
     public void testInferStreamParallelism() throws Exception {
         FileStoreTable fileStoreTable = createTable(ImmutableMap.of("bucket", 
"-1"));
@@ -228,6 +264,29 @@ class DataTableSourceTest {
         return FileStoreTableFactory.create(fileIO, tablePath, tableSchema);
     }
 
+    private FileStoreTable createPostponeTable(boolean partitioned) throws 
Exception {
+        FileIO fileIO = LocalFileIO.create();
+        Path tablePath = new Path(path.toString());
+        SchemaManager schemaManager = new SchemaManager(fileIO, tablePath);
+        Schema.Builder schemaBuilder =
+                Schema.newBuilder()
+                        .column("pt", DataTypes.INT())
+                        .column("a", DataTypes.INT())
+                        .column("b", DataTypes.BIGINT())
+                        .primaryKey("pt", "a")
+                        .option("bucket", "-2")
+                        .option("postpone.merge-on-read", "true");
+        if (partitioned) {
+            schemaBuilder.partitionKeys("pt");
+        } else {
+            schemaBuilder
+                    .option("source.split.target-size", "1 B")
+                    .option("scan.infer-parallelism.max", "3");
+        }
+        TableSchema tableSchema = 
schemaManager.createTable(schemaBuilder.build());
+        return FileStoreTableFactory.create(fileIO, tablePath, tableSchema);
+    }
+
     private void writeData(FileStoreTable table) throws Exception {
         InnerTableWrite writer = table.newWrite("test");
         TableCommitImpl commit = table.newCommit("test");
@@ -236,4 +295,44 @@ class DataTableSourceTest {
         commit.close();
         writer.close();
     }
+
+    private void writePostponeData(FileStoreTable table) throws Exception {
+        InnerTableWrite writer = table.newWrite("test");
+        TableCommitImpl commit = table.newCommit("test");
+        writer.write(GenericRow.of(1, 1, 1L));
+        writer.write(GenericRow.of(1, 2, 2L));
+        writer.write(GenericRow.of(1, 3, 3L));
+        commit.commit(writer.prepareCommit());
+        commit.close();
+        writer.close();
+    }
+
+    private LookupTableSource.LookupContext lookupContext() {
+        return new LookupTableSource.LookupContext() {
+            @Override
+            public int[][] getKeys() {
+                return new int[][] {{1}};
+            }
+
+            @Override
+            public <T> TypeInformation<T> createTypeInformation(DataType 
dataType) {
+                throw new UnsupportedOperationException();
+            }
+
+            @Override
+            public <T> TypeInformation<T> createTypeInformation(LogicalType 
logicalType) {
+                throw new UnsupportedOperationException();
+            }
+
+            @Override
+            public DynamicTableSource.DataStructureConverter 
createDataStructureConverter(
+                    DataType dataType) {
+                throw new UnsupportedOperationException();
+            }
+
+            public boolean preferCustomShuffle() {
+                return false;
+            }
+        };
+    }
 }
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 9f6c46c2a7..4f8a097e47 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
@@ -41,6 +41,7 @@ import java.nio.file.Path;
 
 import static org.apache.paimon.flink.LogicalTypeConversion.toLogicalType;
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
@@ -152,6 +153,27 @@ public class FlinkSourceBuilderTest {
         
assertThat(transformation.getSource()).isInstanceOf(PaimonDataStreamSource.class);
     }
 
+    @Test
+    public void testPostponeMergeOnReadRejectsContinuousSource() throws 
Exception {
+        Identifier identifier = Identifier.create("default", 
"postpone_merge_on_read");
+        catalog.createTable(
+                identifier,
+                Schema.newBuilder()
+                        .column("a", DataTypes.INT())
+                        .primaryKey("a")
+                        .option("bucket", "-2")
+                        .option("postpone.merge-on-read", "true")
+                        .build(),
+                false);
+        Table table = catalog.getTable(identifier);
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+
+        assertThatThrownBy(
+                        () -> new 
FlinkSourceBuilder(table).env(env).sourceBounded(false).build())
+                .isInstanceOf(UnsupportedOperationException.class)
+                .hasMessageContaining("only supported for batch reads");
+    }
+
     @Test
     public void testMonitorSourceBuildSourceWrapsWithPaimonDataStreamSource() 
throws Exception {
         Table table = createTable("monitor_source", false, -1, true);
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala
index c9d19a73ec..e3f70b8172 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonBaseScan.scala
@@ -157,14 +157,14 @@ abstract class PaimonBaseScan(table: InnerTable)
   }
 
   override def estimateStatistics: Statistics = {
-    planPostponeMerge(SparkSession.active.sparkContext.defaultParallelism) 
match {
-      case Some(plan) =>
-        PaimonStatistics(
-          plan.corePlan.splits().asScala.toArray,
-          readTableRowType,
-          table.rowType(),
-          table.statistics())
-      case None => super.estimateStatistics
+    if (postponeMergeOnRead.enabled) {
+      val splits =
+        planPostponeMerge(SparkSession.active.sparkContext.defaultParallelism)
+          .map(_.corePlan.splits().asScala.toArray)
+          .getOrElse(Array.empty[Split])
+      PaimonStatistics(splits, readTableRowType, table.rowType(), 
table.statistics())
+    } else {
+      super.estimateStatistics
     }
   }
 
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonScanBuilder.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonScanBuilder.scala
index c17bc35ba8..4a19890f44 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonScanBuilder.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonScanBuilder.scala
@@ -97,7 +97,7 @@ class PaimonScanBuilder(val table: InnerTable)
 
   // Spark does not support push down aggregation for streaming scan.
   override def pushAggregation(aggregation: Aggregation): Boolean = {
-    if (PostponeMergeOnRead.createReadBuilder(table, 
pushedPartitionFilters).isDefined) {
+    if (PostponeMergeOnRead.usesCustomSource(table)) {
       return false
     }
 
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
index 5dcdb802ff..501874afe2 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PostponeMergeOnRead.scala
@@ -18,6 +18,7 @@
 
 package org.apache.paimon.spark
 
+import org.apache.paimon.CoreOptions
 import org.apache.paimon.partition.PartitionPredicate
 import org.apache.paimon.predicate.PredicateBuilder
 import org.apache.paimon.spark.PostponeMergeOnRead.MergePlan
@@ -45,9 +46,13 @@ final private[spark] class PostponeMergeOnRead(scan: 
PaimonBaseScan) {
 
   @transient private var mergePlan: MergePlan = _
 
-  def enabled: Boolean = mergeReadBuilder.isDefined
+  def enabled: Boolean = PostponeMergeOnRead.usesCustomSource(scan.table)
 
   def plan(defaultBucketNum: Int): Option[MergePlan] = synchronized {
+    if (!enabled) {
+      return None
+    }
+
     mergeReadBuilder.map {
       builder =>
         if (mergePlan == null) {
@@ -86,6 +91,15 @@ private[spark] object PostponeMergeOnRead {
     }
   }
 
+  private[spark] def usesCustomSource(table: Table): Boolean = {
+    table match {
+      case fileStoreTable: FileStoreTable =>
+        configured(fileStoreTable) &&
+        fileStoreTable.coreOptions().startupMode() != 
CoreOptions.StartupMode.COMPACTED_FULL
+      case _ => false
+    }
+  }
+
   private[spark] def createReadBuilder(
       table: Table,
       partitionFilters: Array[PartitionPredicate]): 
Option[PostponeMergeReadBuilder] = {
@@ -94,7 +108,7 @@ private[spark] object PostponeMergeOnRead {
         val partitionFilter =
           if (partitionFilters.isEmpty) null
           else PartitionPredicate.and(partitionFilters.toList.asJava)
-        val result = PostponeMergeReadBuilder.create(fileStoreTable, 
partitionFilter)
+        val result = 
PostponeMergeReadBuilder.createSnapshotBound(fileStoreTable, partitionFilter)
         if (result.isPresent) Some(result.get) else None
       case _ => None
     }
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PaimonStrategy.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PaimonStrategy.scala
index 03fbbb42c4..24b32671c9 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PaimonStrategy.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/execution/PaimonStrategy.scala
@@ -50,7 +50,7 @@ import 
org.apache.spark.sql.catalyst.plans.logical.{AddPartitions, CreateTableAs
 import org.apache.spark.sql.catalyst.plans.physical.{Distribution, 
UnspecifiedDistribution}
 import org.apache.spark.sql.catalyst.util.ArrayData
 import org.apache.spark.sql.connector.catalog.{Identifier, 
PaimonLookupCatalog, TableCatalog}
-import org.apache.spark.sql.execution.{FilterExec, GlobalLimitExec, 
PaimonDescribeTableExec, ProjectExec, SparkPlan, SparkStrategy, UnaryExecNode}
+import org.apache.spark.sql.execution.{FilterExec, GlobalLimitExec, 
LeafExecNode, PaimonDescribeTableExec, ProjectExec, SparkPlan, SparkStrategy, 
UnaryExecNode}
 import org.apache.spark.sql.execution.datasources.v2.{DataSourceV2Implicits, 
DataSourceV2Relation, DataSourceV2ScanRelation}
 import org.apache.spark.sql.execution.exchange.ShuffleExchangeLike
 import org.apache.spark.sql.execution.joins.{BroadcastHashJoinExec, 
BroadcastNestedLoopJoinExec}
@@ -72,7 +72,7 @@ case class PaimonStrategy(spark: SparkSession)
 
     case PhysicalOperation(projects, filters, relation: 
DataSourceV2ScanRelation) =>
       relation.scan match {
-        case scan: PaimonScan if PostponeMergeOnRead.enabled(scan.table) =>
+        case scan: PaimonScan if 
PostponeMergeOnRead.usesCustomSource(scan.table) =>
           scan.planPostponeMerge(spark.sparkContext.defaultParallelism) match {
             case Some(mergePlan) =>
               val inputScan = PostponeMergeInputScan(mergePlan)
@@ -96,7 +96,8 @@ case class PaimonStrategy(spark: SparkSession)
               val filtered =
                 filters.reduceLeftOption(And).map(FilterExec(_, 
merge)).getOrElse(merge)
               ProjectExec(projects, filtered) :: Nil
-            case None => Nil
+            case None =>
+              ProjectExec(projects, EmptyPostponeMergeExec(relation.output)) 
:: Nil
           }
         case _ => Nil
       }
@@ -328,6 +329,11 @@ case class PaimonStrategy(spark: SparkSession)
   }
 }
 
+private case class EmptyPostponeMergeExec(output: Seq[Attribute]) extends 
LeafExecNode {
+
+  override protected def doExecute(): RDD[InternalRow] = 
sparkContext.emptyRDD[InternalRow]
+}
+
 case class LateralVectorSearchExec(
     innerTable: InnerTable,
     columnName: String,
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
index 187ddf82de..805b747bb6 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
@@ -19,6 +19,7 @@
 package org.apache.paimon.spark.sql
 
 import org.apache.paimon.catalog.{Catalog, CatalogLoader, DelegateCatalog, 
Identifier}
+import 
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX
 import org.apache.paimon.fs.Path
 import org.apache.paimon.spark.{PaimonScan, PaimonSparkTestBase}
 import org.apache.paimon.spark.procedure.SparkPostponeCompactProcedure
@@ -26,9 +27,10 @@ import org.apache.paimon.table.{CatalogEnvironment, 
FileStoreTableFactory}
 
 import org.apache.spark.TaskContext
 import org.apache.spark.sql.Row
-import org.apache.spark.sql.execution.aggregate.BaseAggregateExec
 import org.apache.spark.sql.execution.datasources.v2.DataSourceV2ScanRelation
 
+import scala.collection.JavaConverters._
+
 class PostponeBucketTableTest extends PaimonSparkTestBase {
 
   test("Postpone bucket table: write with different bucket number") {
@@ -160,15 +162,12 @@ class PostponeBucketTableTest extends PaimonSparkTestBase 
{
       withSparkSQLConf("spark.paimon.postpone.merge-on-read" -> "true") {
         checkAnswer(sql("SELECT * FROM t ORDER BY k"), Seq(Row(1, "base-1"), 
Row(2, "base-2")))
         val plan = sql("SELECT * FROM 
t").queryExecution.executedPlan.toString()
-        assert(!plan.contains("PostponeMergeOnRead"), plan)
+        assert(plan.contains("PostponeMergeOnRead"), plan)
 
         val aggregate = sql("SELECT count(*) FROM t")
         checkAnswer(aggregate, Seq(Row(2L)))
-        assert(
-          aggregate.queryExecution.executedPlan.collect {
-            case _: BaseAggregateExec => true
-          }.isEmpty,
-          aggregate.queryExecution.executedPlan)
+        val aggregatePlan = aggregate.queryExecution.executedPlan.toString()
+        assert(aggregatePlan.contains("HashAggregate"), aggregatePlan)
       }
 
       sql("CREATE TABLE normal_t (k INT, v STRING) USING paimon")
@@ -285,6 +284,57 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
     }
   }
 
+  test("Postpone bucket table: Spark combines postpone and deletion-vector 
merge on read") {
+    withTable("t") {
+      sql("""
+            |CREATE TABLE t (
+            |  k INT,
+            |  v STRING
+            |) TBLPROPERTIES (
+            |  'primary-key' = 'k',
+            |  'bucket' = '-2',
+            |  'postpone.batch-write-fixed-bucket' = 'false',
+            |  'deletion-vectors.enabled' = 'true',
+            |  'deletion-vectors.merge-on-read' = 'true'
+            |)
+            |""".stripMargin)
+
+      sql("INSERT INTO t VALUES (1, 'base-1'), (2, 'base-2'), (3, 'base-3'), 
(4, 'base-4')")
+      sql("CALL sys.compact(table => 't')")
+      sql("INSERT INTO t VALUES (1, 'real-1'), (5, 'real-5')")
+      sql("CALL sys.compact(table => 't')")
+      assert(deletionVectorCardinality("t") > 0)
+
+      sql("INSERT INTO t VALUES (2, 'postpone-2'), (3, 'postpone-3'), (6, 
'postpone-6')")
+      sql("DELETE FROM t WHERE k = 4")
+
+      checkAnswer(
+        sql("SELECT * FROM t ORDER BY k"),
+        Seq(
+          Row(1, "real-1"),
+          Row(2, "base-2"),
+          Row(3, "base-3"),
+          Row(4, "base-4"),
+          Row(5, "real-5")))
+
+      withSparkSQLConf(
+        "spark.paimon.postpone.merge-on-read" -> "true",
+        "spark.paimon.deletion-vectors.merge-on-read" -> "true") {
+        val query = sql("SELECT * FROM t ORDER BY k")
+        checkAnswer(
+          query,
+          Seq(
+            Row(1, "real-1"),
+            Row(2, "postpone-2"),
+            Row(3, "postpone-3"),
+            Row(5, "real-5"),
+            Row(6, "postpone-6")))
+        checkAnswer(sql("SELECT count(*) FROM t"), Seq(Row(5L)))
+        
assert(query.queryExecution.executedPlan.toString.contains("PostponeMergeOnRead"))
+      }
+    }
+  }
+
   test("Postpone bucket table: Spark merge on read respects full scan 
protection") {
     withTable("t") {
       sql("""
@@ -406,7 +456,66 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
         val ordinaryPlan =
           sql("SELECT * FROM t WHERE pt = 
'b'").queryExecution.executedPlan.toString
         assert(mergePlan.contains("PostponeMergeOnRead"), mergePlan)
-        assert(!ordinaryPlan.contains("PostponeMergeOnRead"), ordinaryPlan)
+        assert(ordinaryPlan.contains("PostponeMergeOnRead"), ordinaryPlan)
+      }
+    }
+  }
+
+  test("Postpone bucket table: Spark merge on read pins the selected 
snapshot") {
+    withTable("real_t", "empty_t") {
+      sql("""
+            |CREATE TABLE real_t (
+            |  k INT,
+            |  v STRING
+            |) TBLPROPERTIES (
+            |  'primary-key' = 'k',
+            |  'bucket' = '-2',
+            |  'postpone.batch-write-fixed-bucket' = 'true'
+            |)
+            |""".stripMargin)
+      sql("INSERT INTO real_t VALUES (1, 'base')")
+
+      sql("""
+            |CREATE TABLE empty_t (
+            |  k INT,
+            |  v STRING
+            |) TBLPROPERTIES (
+            |  'primary-key' = 'k',
+            |  'bucket' = '-2',
+            |  'postpone.batch-write-fixed-bucket' = 'true'
+            |)
+            |""".stripMargin)
+
+      withSparkSQLConf("spark.paimon.postpone.merge-on-read" -> "true") {
+        val pinnedReal = sql("SELECT * FROM real_t ORDER BY k")
+        val pinnedRealScan = 
pinnedReal.queryExecution.optimizedPlan.collectFirst {
+          case relation: DataSourceV2ScanRelation if 
relation.scan.isInstanceOf[PaimonScan] =>
+            relation.scan.asInstanceOf[PaimonScan]
+        }.get
+        assert(
+          pinnedRealScan
+            .planPostponeMerge(spark.sparkContext.defaultParallelism)
+            .isDefined)
+
+        withSparkSQLConf("spark.paimon.postpone.batch-write-fixed-bucket" -> 
"false") {
+          sql("INSERT INTO real_t VALUES (1, 'new'), (2, 'added')")
+        }
+        checkAnswer(pinnedReal, Seq(Row(1, "base")))
+        checkAnswer(sql("SELECT * FROM real_t ORDER BY k"), Seq(Row(1, "new"), 
Row(2, "added")))
+
+        val pinnedEmpty = sql("SELECT * FROM empty_t")
+        val pinnedEmptyScan = 
pinnedEmpty.queryExecution.optimizedPlan.collectFirst {
+          case relation: DataSourceV2ScanRelation if 
relation.scan.isInstanceOf[PaimonScan] =>
+            relation.scan.asInstanceOf[PaimonScan]
+        }.get
+        assert(
+          pinnedEmptyScan
+            .planPostponeMerge(spark.sparkContext.defaultParallelism)
+            .isEmpty)
+
+        sql("INSERT INTO empty_t VALUES (1, 'committed-later')")
+        checkAnswer(pinnedEmpty, Seq.empty)
+        checkAnswer(sql("SELECT * FROM empty_t"), Seq(Row(1, 
"committed-later")))
       }
     }
   }
@@ -742,6 +851,19 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
       )
     }
   }
+
+  private def deletionVectorCardinality(tableName: String): Long = {
+    val table = loadTable(tableName)
+    table
+      .store()
+      .newIndexFileHandler()
+      .scan(table.latestSnapshot().get(), DELETION_VECTORS_INDEX)
+      .asScala
+      .flatMap(entry => Option(entry.indexFile().dvRanges()).toSeq)
+      .flatMap(_.values().asScala)
+      .flatMap(meta => Option(meta.cardinality()).map(_.longValue()))
+      .sum
+  }
 }
 
 object PostponeBucketTableTest {


Reply via email to