diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java index 67a73da98695..7b77a32b493e 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java @@ -150,6 +150,21 @@ public static Optional> buildIndexStream( PartitionPredicate partitionPredicate, Options userOptions) throws Exception { + Long lastSafeSnapshotId = table.snapshotManager().latestSnapshotId(); + table = + table.copy( + new HashMap() { + { + put( + CoreOptions.COMMIT_LAST_SAFE_SNAPSHOT.key(), + Long.toString( + lastSafeSnapshotId == null + ? 0 + : lastSafeSnapshotId)); + put(CoreOptions.COMMIT_STRICT_MODE_ENABLED.key(), "false"); + } + }); + List> allStreams = new ArrayList<>(); for (String indexColumn : indexColumns) { SortedGlobalIndexScanner indexScanner =