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 c721d2c532 [core] Support latest snapshot delta scan in batch (#9142)
c721d2c532 is described below

commit c721d2c532af1825ecb84e9134b7f504292f7ded
Author: wangwj <[email protected]>
AuthorDate: Thu Aug 13 22:11:17 2026 +0800

    [core] Support latest snapshot delta scan in batch (#9142)
---
 docs/generated/core_configuration.html             |  2 +-
 .../main/java/org/apache/paimon/CoreOptions.java   |  6 +++
 .../org/apache/paimon/schema/SchemaValidation.java | 22 +++++++++
 .../table/source/AbstractBatchTableScan.java       |  3 +-
 .../paimon/table/source/AbstractDataTableScan.java | 13 +++++
 .../table/source/PostponeMergeReadBuilder.java     |  1 +
 .../apache/paimon/schema/SchemaValidationTest.java | 31 ++++++++++++
 .../paimon/table/source/StartupModeTest.java       | 56 ++++++++++++++++++++++
 .../apache/paimon/table/source/TableScanTest.java  | 15 ++++++
 .../apache/paimon/flink/BatchFileStoreITCase.java  | 13 +++++
 .../org/apache/paimon/flink/FlinkCatalogTest.java  |  4 +-
 11 files changed, 163 insertions(+), 3 deletions(-)

diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index 0c576c3934..ab8443e4cb 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -1468,7 +1468,7 @@ For an internal format table in a REST catalog, it also 
makes the catalog own th
             <td><h5>scan.mode</h5></td>
             <td style="word-wrap: break-word;">default</td>
             <td><p>Enum</p></td>
-            <td>Specify the scanning behavior of the source.<br /><br 
/>Possible values:<ul><li>"default": Determines actual startup mode according 
to other table properties. If "scan.timestamp-millis" is set the actual startup 
mode will be "from-timestamp", and if "scan.snapshot-id" or "scan.tag-name" is 
set the actual startup mode will be "from-snapshot". Otherwise the actual 
startup mode will be "latest-full".</li><li>"latest-full": For streaming 
sources, produces the latest snapshot  [...]
+            <td>Specify the scanning behavior of the source.<br /><br 
/>Possible values:<ul><li>"default": Determines actual startup mode according 
to other table properties. If "scan.timestamp-millis" is set the actual startup 
mode will be "from-timestamp", and if "scan.snapshot-id" or "scan.tag-name" is 
set the actual startup mode will be "from-snapshot". Otherwise the actual 
startup mode will be "latest-full".</li><li>"latest-full": For streaming 
sources, produces the latest snapshot  [...]
         </tr>
         <tr>
             <td><h5>scan.plan-auto-tag-for-read.time-retained</h5></td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 12e02ee2a6..b30e015442 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -4836,6 +4836,12 @@ public class CoreOptions implements Serializable {
                         + "without producing a snapshot at the beginning. "
                         + "For batch sources, behaves the same as the 
\"latest-full\" startup mode."),
 
+        LATEST_DELTA(
+                "latest-delta",
+                "For batch sources, reads newly changed files from the latest 
snapshot. "
+                        + "This mode does not search backwards for an APPEND 
snapshot, so a latest "
+                        + "COMPACT or OVERWRITE snapshot produces no records. 
Streaming sources are not supported."),
+
         COMPACTED_FULL(
                 "compacted-full",
                 "For streaming sources, produces a snapshot after the latest 
compaction on the table "
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java 
b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java
index 1bd587d4ca..607f5d20fa 100644
--- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java
+++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java
@@ -76,16 +76,20 @@ import static org.apache.paimon.CoreOptions.FIELDS_PREFIX;
 import static org.apache.paimon.CoreOptions.FIELDS_SEPARATOR;
 import static org.apache.paimon.CoreOptions.FULL_COMPACTION_DELTA_COMMITS;
 import static org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN;
+import static org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN_SCAN_MODE;
+import static 
org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN_TAG_TO_SNAPSHOT;
 import static org.apache.paimon.CoreOptions.INCREMENTAL_BETWEEN_TIMESTAMP;
 import static org.apache.paimon.CoreOptions.INCREMENTAL_TO_AUTO_TAG;
 import static org.apache.paimon.CoreOptions.MAP_STORAGE_LAYOUT;
 import static org.apache.paimon.CoreOptions.PRIMARY_KEY;
+import static org.apache.paimon.CoreOptions.SCAN_CREATION_TIME_MILLIS;
 import static org.apache.paimon.CoreOptions.SCAN_FILE_CREATION_TIME_MILLIS;
 import static org.apache.paimon.CoreOptions.SCAN_MODE;
 import static org.apache.paimon.CoreOptions.SCAN_SNAPSHOT_ID;
 import static org.apache.paimon.CoreOptions.SCAN_TAG_NAME;
 import static org.apache.paimon.CoreOptions.SCAN_TIMESTAMP;
 import static org.apache.paimon.CoreOptions.SCAN_TIMESTAMP_MILLIS;
+import static org.apache.paimon.CoreOptions.SCAN_VERSION;
 import static org.apache.paimon.CoreOptions.SCAN_WATERMARK;
 import static org.apache.paimon.CoreOptions.SNAPSHOT_NUM_RETAINED_MAX;
 import static org.apache.paimon.CoreOptions.SNAPSHOT_NUM_RETAINED_MIN;
@@ -523,6 +527,24 @@ public class SchemaValidation {
                             INCREMENTAL_BETWEEN,
                             INCREMENTAL_TO_AUTO_TAG),
                     Collections.singletonList(SCAN_FILE_CREATION_TIME_MILLIS));
+        } else if (options.startupMode() == 
CoreOptions.StartupMode.LATEST_DELTA) {
+            for (ConfigOption<?> option :
+                    Arrays.asList(
+                            SCAN_TIMESTAMP_MILLIS,
+                            SCAN_FILE_CREATION_TIME_MILLIS,
+                            SCAN_CREATION_TIME_MILLIS,
+                            SCAN_TIMESTAMP,
+                            SCAN_SNAPSHOT_ID,
+                            SCAN_TAG_NAME,
+                            SCAN_WATERMARK,
+                            SCAN_VERSION,
+                            INCREMENTAL_BETWEEN_TIMESTAMP,
+                            INCREMENTAL_BETWEEN,
+                            INCREMENTAL_TO_AUTO_TAG,
+                            INCREMENTAL_BETWEEN_SCAN_MODE,
+                            INCREMENTAL_BETWEEN_TAG_TO_SNAPSHOT)) {
+                checkOptionNotExistInMode(options, option, 
options.startupMode());
+            }
         } else {
             checkOptionNotExistInMode(options, SCAN_TIMESTAMP_MILLIS, 
options.startupMode());
             checkOptionNotExistInMode(
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractBatchTableScan.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractBatchTableScan.java
index 9e4e77cf42..b52424b937 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractBatchTableScan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractBatchTableScan.java
@@ -77,7 +77,8 @@ public abstract class AbstractBatchTableScan extends 
AbstractDataTableScan {
             if (options.toConfiguration()
                             .get(CoreOptions.BATCH_SCAN_MODE)
                             .equals(CoreOptions.BatchScanMode.NONE)
-                    && options.startupMode() != 
CoreOptions.StartupMode.INCREMENTAL) {
+                    && options.startupMode() != 
CoreOptions.StartupMode.INCREMENTAL
+                    && options.startupMode() != 
CoreOptions.StartupMode.LATEST_DELTA) {
                 snapshotReader.withLevelFilter(level -> level > 
0).enableValueFilter();
             }
         }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
index 6695e5cefe..87257a744f 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
@@ -286,6 +286,19 @@ abstract class AbstractDataTableScan implements 
DataTableScan {
                 return isStreaming
                         ? new ContinuousLatestStartingScanner(snapshotManager)
                         : new FullStartingScanner(snapshotManager);
+            case LATEST_DELTA:
+                checkArgument(
+                        !isStreaming,
+                        "'latest-delta' scan mode is only supported for batch 
sources.");
+                Snapshot latestSnapshot = snapshotManager.latestSnapshot();
+                if (latestSnapshot == null) {
+                    return new EmptyResultStartingScanner(snapshotManager);
+                }
+                return IncrementalDeltaStartingScanner.betweenSnapshotIds(
+                        latestSnapshot.id() - 1,
+                        latestSnapshot.id(),
+                        snapshotManager,
+                        ScanMode.DELTA);
             case COMPACTED_FULL:
                 if (options.changelogProducer() == 
ChangelogProducer.FULL_COMPACTION
                         || 
options.toConfiguration().contains(FULL_COMPACTION_DELTA_COMMITS)) {
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 147d12e05a..09e04e1cca 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
@@ -313,6 +313,7 @@ public final class PostponeMergeReadBuilder implements 
Serializable {
         }
         CoreOptions.StartupMode startupMode = 
table.coreOptions().startupMode();
         if (startupMode == CoreOptions.StartupMode.INCREMENTAL
+                || startupMode == CoreOptions.StartupMode.LATEST_DELTA
                 || startupMode == 
CoreOptions.StartupMode.FROM_FILE_CREATION_TIME
                 || startupMode == 
CoreOptions.StartupMode.FROM_CREATION_TIMESTAMP) {
             throw new UnsupportedOperationException(
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java 
b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java
index a1801ab9c0..3ba69852f3 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java
@@ -133,6 +133,37 @@ class SchemaValidationTest {
                         "must set only one key in 
[scan.timestamp-millis,scan.timestamp] when you use from-timestamp for 
scan.mode");
     }
 
+    @Test
+    public void testLatestDeltaOnlyAcceptsScanMode() {
+        Map<String, String> options = new HashMap<>();
+        options.put(CoreOptions.SCAN_MODE.key(), 
CoreOptions.StartupMode.LATEST_DELTA.toString());
+        assertThatNoException().isThrownBy(() -> 
validateTableSchemaExec(options));
+
+        Map<String, String> incompatibleOptions = new HashMap<>();
+        incompatibleOptions.put(CoreOptions.SCAN_TIMESTAMP_MILLIS.key(), "1");
+        
incompatibleOptions.put(CoreOptions.SCAN_FILE_CREATION_TIME_MILLIS.key(), "1");
+        incompatibleOptions.put(CoreOptions.SCAN_CREATION_TIME_MILLIS.key(), 
"1");
+        incompatibleOptions.put(CoreOptions.SCAN_TIMESTAMP.key(), "2026-08-10 
00:00:00");
+        incompatibleOptions.put(CoreOptions.SCAN_SNAPSHOT_ID.key(), "1");
+        incompatibleOptions.put(CoreOptions.SCAN_TAG_NAME.key(), "tag1");
+        incompatibleOptions.put(CoreOptions.SCAN_WATERMARK.key(), "1");
+        incompatibleOptions.put(CoreOptions.SCAN_VERSION.key(), "1");
+        
incompatibleOptions.put(CoreOptions.INCREMENTAL_BETWEEN_TIMESTAMP.key(), "1,2");
+        incompatibleOptions.put(CoreOptions.INCREMENTAL_BETWEEN.key(), "1,2");
+        incompatibleOptions.put(CoreOptions.INCREMENTAL_TO_AUTO_TAG.key(), 
"tag1");
+        
incompatibleOptions.put(CoreOptions.INCREMENTAL_BETWEEN_SCAN_MODE.key(), 
"delta");
+        
incompatibleOptions.put(CoreOptions.INCREMENTAL_BETWEEN_TAG_TO_SNAPSHOT.key(), 
"true");
+
+        incompatibleOptions.forEach(
+                (key, value) -> {
+                    Map<String, String> invalidOptions = new 
HashMap<>(options);
+                    invalidOptions.put(key, value);
+                    assertThatThrownBy(() -> 
validateTableSchemaExec(invalidOptions))
+                            .hasMessageContaining(
+                                    key + " must be null when you use 
latest-delta for scan.mode");
+                });
+    }
+
     @Test
     public void testTargetFileRowNumMustBePositive() {
         Map<String, String> options = new HashMap<>();
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/source/StartupModeTest.java 
b/paimon-core/src/test/java/org/apache/paimon/table/source/StartupModeTest.java
index 2d3b60f3ae..73010cae63 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/source/StartupModeTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/source/StartupModeTest.java
@@ -19,7 +19,9 @@
 package org.apache.paimon.table.source;
 
 import org.apache.paimon.CoreOptions;
+import org.apache.paimon.Snapshot;
 import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.disk.IOManager;
 import org.apache.paimon.fs.FileIOFinder;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.options.Options;
@@ -88,6 +90,60 @@ public class StartupModeTest extends ScannerTestBase {
                 
.isEqualTo(snapshotReader.withSnapshot(4).withMode(ScanMode.ALL).read().splits());
     }
 
+    @Test
+    public void testStartFromLatestDelta() throws Exception {
+        initializeTable(StartupMode.LATEST_DELTA);
+        initializeTestData(); // initialize 3 commits
+
+        TableScan.Plan plan = table.newScan().plan();
+        assertThat(plan.splits())
+                
.isEqualTo(snapshotReader.withSnapshot(3).withMode(ScanMode.DELTA).read().splits());
+
+        // Do not search backwards for an APPEND snapshot when the latest 
snapshot is COMPACT.
+        write.compact(binaryRow(1), 0, true);
+        commit.commit(4, write.prepareCommit(true, 4));
+        assertThat(table.snapshotManager().latestSnapshot().id()).isEqualTo(4);
+        assertThat(table.snapshotManager().latestSnapshot().commitKind())
+                .isEqualTo(Snapshot.CommitKind.COMPACT);
+        assertThat(table.newScan().plan().splits()).isEmpty();
+
+        writeAndCommit(5, rowData(1, 10, 103L));
+        assertThat(table.newScan().plan().splits())
+                
.isEqualTo(snapshotReader.withSnapshot(5).withMode(ScanMode.DELTA).read().splits());
+    }
+
+    @Test
+    public void testStartFromLatestDeltaWithoutSnapshot() throws Exception {
+        initializeTable(StartupMode.LATEST_DELTA);
+
+        assertThat(table.newScan().plan().splits()).isEmpty();
+        assertThatThrownBy(() -> table.newStreamScan().plan())
+                .satisfies(
+                        anyCauseMatches(
+                                IllegalArgumentException.class,
+                                "'latest-delta' scan mode is only supported 
for batch sources."));
+    }
+
+    @Test
+    public void testStartFromLatestDeltaDoesNotSkipLevelZero() throws 
Exception {
+        Map<String, String> properties = new HashMap<>();
+        properties.put(
+                CoreOptions.MERGE_ENGINE.key(), 
CoreOptions.MergeEngine.FIRST_ROW.toString());
+        initializeTable(StartupMode.LATEST_DELTA, properties);
+        try (IOManager ioManager = 
IOManager.create(tempDir.resolve("latest-delta").toString());
+                StreamTableWrite levelZeroWrite =
+                        table.newWrite(commitUser).withIOManager(ioManager);
+                StreamTableCommit levelZeroCommit = 
table.newCommit(commitUser)) {
+            levelZeroWrite.write(rowData(1, 10, 100L));
+            levelZeroCommit.commit(1, levelZeroWrite.prepareCommit(false, 1));
+        }
+
+        TableScan.Plan plan = table.newScan().plan();
+        assertThat(plan.splits()).isNotEmpty();
+        assertThat(plan.splits())
+                
.isEqualTo(snapshotReader.withSnapshot(1).withMode(ScanMode.DELTA).read().splits());
+    }
+
     @Test
     public void testStartFromLatestFull() throws Exception {
         initializeTable(StartupMode.LATEST_FULL);
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 2f8f8959ed..a60893aa65 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
@@ -496,6 +496,21 @@ public class TableScanTest extends ScannerTestBase {
                 .deleteReadTag(readProtectionTag);
     }
 
+    @Test
+    public void testPostponeMergeRejectsLatestDelta() {
+        Map<String, String> dynamicOptions = new HashMap<>();
+        dynamicOptions.put(CoreOptions.BUCKET.key(), "-2");
+        dynamicOptions.put(CoreOptions.POSTPONE_MERGE_ON_READ.key(), "true");
+        dynamicOptions.put(CoreOptions.SCAN_MODE.key(), "latest-delta");
+        FileStoreTable latestDeltaTable = table.copy(dynamicOptions);
+
+        assertThatThrownBy(
+                        () -> 
PostponeMergeReadBuilder.createSnapshotBound(latestDeltaTable, null))
+                .isInstanceOf(UnsupportedOperationException.class)
+                .hasMessageContaining("requires a full snapshot scan")
+                .hasMessageContaining("latest-delta");
+    }
+
     @Test
     public void testPostponeMergePlanAndRead() throws Exception {
         StreamTableWrite realWrite = table.newWrite(commitUser);
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
index 6ce5f53387..6c9570dddc 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
@@ -1165,6 +1165,19 @@ public class BatchFileStoreITCase extends 
CatalogITCaseBase {
                 .isEmpty();
     }
 
+    @Test
+    public void testLatestDeltaScanMode() {
+        sql("CREATE TABLE latest_delta (id INT, v STRING)");
+        sql("INSERT INTO latest_delta VALUES (1, 'A'), (2, 'B')");
+        sql("INSERT INTO latest_delta VALUES (3, 'C'), (4, 'D')");
+
+        assertThat(
+                        sql(
+                                "SELECT * FROM latest_delta "
+                                        + "/*+ 
OPTIONS('scan.mode'='latest-delta') */"))
+                .containsExactlyInAnyOrder(Row.of(3, "C"), Row.of(4, "D"));
+    }
+
     @Test
     public void testIncrementScanMode() throws Exception {
         sql(
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkCatalogTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkCatalogTest.java
index 2bb9c580e5..cc0d2a8083 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkCatalogTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkCatalogTest.java
@@ -913,7 +913,9 @@ public class FlinkCatalogTest extends FlinkCatalogTestBase {
                 options.put("incremental-between", "2,5");
             }
 
-            if (isStreaming && mode == CoreOptions.StartupMode.INCREMENTAL) {
+            if (isStreaming
+                    && (mode == CoreOptions.StartupMode.INCREMENTAL
+                            || mode == CoreOptions.StartupMode.LATEST_DELTA)) {
                 continue;
             }
             allOptions.add(options);

Reply via email to