From fd8b8c57bc91f13d4b222db4053698335d73b0d4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Wed, 23 Sep 2026 15:40:28 +0800 Subject: [PATCH 1/2] Safe snapshot for sorted index build --- .../globalindex/SortedIndexTopoBuilder.java | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) 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..3dd8ffde4c65 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 @@ -118,6 +118,22 @@ public static boolean buildIndex( 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"); + } + + }); + + Optional> written = buildIndexStream( env, From 51249c6f43218c7b5b6e71d515aab98f133a637e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Wed, 23 Sep 2026 19:40:37 +0800 Subject: [PATCH 2/2] Fix minus --- .../globalindex/SortedIndexTopoBuilder.java | 31 +++++++++---------- 1 file changed, 15 insertions(+), 16 deletions(-) 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 3dd8ffde4c65..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 @@ -118,22 +118,6 @@ public static boolean buildIndex( 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"); - } - - }); - - Optional> written = buildIndexStream( env, @@ -166,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 =