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);