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 555749f923 [core][flink] Pin DV vector-search read path to the plan 
snapshot (#8946)
555749f923 is described below

commit 555749f923ee597da07915e0d5fac990f4fbba3f
Author: XiaoHongbo <[email protected]>
AuthorDate: Fri Jul 31 10:13:44 2026 +0800

    [core][flink] Pin DV vector-search read path to the plan snapshot (#8946)
---
 .../DataEvolutionGlobalIndexScanner.java           |  5 ++--
 .../source/AbstractDataEvolutionVectorRead.java    | 13 +++++++---
 .../table/source/DataEvolutionBatchVectorRead.java |  1 +
 .../table/source/DataEvolutionVectorRead.java      |  1 +
 .../table/source/DataEvolutionVectorScan.java      | 22 +++++++++++++++-
 .../table/source/GlobalIndexLiveRowFilter.java     |  8 ++++--
 .../org/apache/paimon/table/source/VectorScan.java | 10 ++++++++
 .../table/source/VectorSearchBuilderImpl.java      |  3 ++-
 .../table/source/VectorSearchBuilderTest.java      | 29 ++++++++++++++++++++++
 .../vectorsearch/FlinkDataEvolutionVectorRead.java |  2 ++
 10 files changed, 85 insertions(+), 9 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
 
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
index 215e1b31c5..7dbdb06853 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java
@@ -195,11 +195,12 @@ public class DataEvolutionGlobalIndexScanner implements 
Closeable {
 
     public static Optional<DataEvolutionGlobalIndexScanner> create(
             FileStoreTable table, Collection<IndexFileMeta> indexFiles) {
-        return create(table, null, indexFiles);
+        return create(table, null, null, indexFiles);
     }
 
     public static Optional<DataEvolutionGlobalIndexScanner> create(
             FileStoreTable table,
+            @Nullable Snapshot pinnedSnapshot,
             @Nullable PartitionPredicate partitionFilter,
             Collection<IndexFileMeta> indexFiles) {
         List<IndexFileMeta> globalIndexFiles = globalIndexFiles(indexFiles);
@@ -209,7 +210,7 @@ public class DataEvolutionGlobalIndexScanner implements 
Closeable {
         return Optional.of(
                 new DataEvolutionGlobalIndexScanner(
                         table,
-                        tryTravelOrLatest(table),
+                        pinnedSnapshot != null ? pinnedSnapshot : 
tryTravelOrLatest(table),
                         partitionFilter,
                         table.coreOptions().toConfiguration(),
                         table.rowType(),
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java
index 2545a14777..1f505cc128 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java
@@ -18,6 +18,7 @@
 
 package org.apache.paimon.table.source;
 
+import org.apache.paimon.Snapshot;
 import org.apache.paimon.data.InternalArray;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.data.InternalVector;
@@ -82,6 +83,9 @@ public abstract class AbstractDataEvolutionVectorRead 
implements Serializable {
     protected final DataField vectorColumn;
     protected final Map<String, String> options;
 
+    /** Snapshot the plan was built against; pins live-row filtering to it. */
+    @Nullable protected transient Snapshot planSnapshot;
+
     private static final Comparator<long[]> WEAKEST_SCORE_FIRST =
             Comparator.<long[]>comparingDouble(a -> Float.intBitsToFloat((int) 
a[1]))
                     .thenComparing((a, b) -> Long.compare(b[0], a[0]));
@@ -133,7 +137,8 @@ public abstract class AbstractDataEvolutionVectorRead 
implements Serializable {
         }
 
         RoaringNavigableMap64 liveRows =
-                GlobalIndexLiveRowFilter.liveRows(table, partitionFilter, 
indexedRowRanges);
+                GlobalIndexLiveRowFilter.liveRows(
+                        table, planSnapshot, partitionFilter, 
indexedRowRanges);
         RoaringNavigableMap64 matchedRows = scalarMatchedRows(splits);
 
         List<RoaringNavigableMap64> includeRowIds = new 
ArrayList<>(splits.size());
@@ -173,7 +178,8 @@ public abstract class AbstractDataEvolutionVectorRead 
implements Serializable {
         }
 
         Optional<DataEvolutionGlobalIndexScanner> optionalScanner =
-                DataEvolutionGlobalIndexScanner.create(table, partitionFilter, 
scalarIndexFiles);
+                DataEvolutionGlobalIndexScanner.create(
+                        table, planSnapshot, partitionFilter, 
scalarIndexFiles);
         if (!optionalScanner.isPresent()) {
             return new RoaringNavigableMap64();
         }
@@ -206,7 +212,8 @@ public abstract class AbstractDataEvolutionVectorRead 
implements Serializable {
             scalarIndexFiles.addAll(split.scalarIndexFiles());
         }
         Optional<DataEvolutionGlobalIndexScanner> optionalScanner =
-                DataEvolutionGlobalIndexScanner.create(table, partitionFilter, 
scalarIndexFiles);
+                DataEvolutionGlobalIndexScanner.create(
+                        table, planSnapshot, partitionFilter, 
scalarIndexFiles);
         if (!optionalScanner.isPresent()) {
             return null;
         }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionBatchVectorRead.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionBatchVectorRead.java
index 9df3e2df0b..6c83a62b23 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionBatchVectorRead.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionBatchVectorRead.java
@@ -62,6 +62,7 @@ public class DataEvolutionBatchVectorRead extends 
AbstractDataEvolutionVectorRea
 
     @Override
     public List<GlobalIndexResult> readBatch(VectorScan.Plan plan) {
+        this.planSnapshot = plan.snapshot();
         return readBatch(plan.splits());
     }
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorRead.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorRead.java
index 967aecd6bb..9179dd7512 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorRead.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorRead.java
@@ -61,6 +61,7 @@ public class DataEvolutionVectorRead extends 
AbstractDataEvolutionVectorRead imp
 
     @Override
     public GlobalIndexResult read(VectorScan.Plan plan) {
+        this.planSnapshot = plan.snapshot();
         return readSplits(plan.splits());
     }
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
index e9a26020d5..2064b5614e 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionVectorScan.java
@@ -55,6 +55,7 @@ public class DataEvolutionVectorScan implements VectorScan {
     @Nullable private final Predicate filter;
     private final DataField vectorColumn;
     private final Map<String, String> options;
+    @Nullable private final Snapshot pinnedSnapshot;
 
     public DataEvolutionVectorScan(
             FileStoreTable table,
@@ -62,10 +63,21 @@ public class DataEvolutionVectorScan implements VectorScan {
             @Nullable Predicate filter,
             DataField vectorColumn,
             @Nullable Map<String, String> options) {
+        this(table, partitionFilter, filter, vectorColumn, options, null);
+    }
+
+    public DataEvolutionVectorScan(
+            FileStoreTable table,
+            @Nullable PartitionPredicate partitionFilter,
+            @Nullable Predicate filter,
+            DataField vectorColumn,
+            @Nullable Map<String, String> options,
+            @Nullable Snapshot pinnedSnapshot) {
         this.table = table;
         this.partitionFilter = partitionFilter;
         this.filter = filter;
         this.vectorColumn = vectorColumn;
+        this.pinnedSnapshot = pinnedSnapshot;
         this.options =
                 options == null
                         ? Collections.emptyMap()
@@ -77,7 +89,9 @@ public class DataEvolutionVectorScan implements VectorScan {
         Objects.requireNonNull(vectorColumn, "Vector column must be set");
 
         Set<Integer> filterFieldIds = collectFieldIds(table.rowType(), filter);
-        @Nullable Snapshot snapshot = TimeTravelUtil.tryTravelOrLatest(table);
+        @Nullable
+        Snapshot snapshot =
+                pinnedSnapshot != null ? pinnedSnapshot : 
TimeTravelUtil.tryTravelOrLatest(table);
         IndexFileHandler indexFileHandler = 
table.store().newIndexFileHandler();
         Filter<IndexManifestEntry> indexFileFilter =
                 entry -> {
@@ -181,11 +195,17 @@ public class DataEvolutionVectorScan implements 
VectorScan {
                             vectorIndexType));
         }
 
+        @Nullable Snapshot planSnapshot = snapshot;
         return new Plan() {
             @Override
             public List<VectorSearchSplit> splits() {
                 return splits;
             }
+
+            @Override
+            public Snapshot snapshot() {
+                return planSnapshot;
+            }
         };
     }
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/GlobalIndexLiveRowFilter.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/GlobalIndexLiveRowFilter.java
index ebd9c5af3f..d80ea2b25d 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/GlobalIndexLiveRowFilter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/GlobalIndexLiveRowFilter.java
@@ -40,19 +40,23 @@ class GlobalIndexLiveRowFilter {
     @Nullable
     static RoaringNavigableMap64 liveRows(
             @Nullable FileStoreTable table, @Nullable PartitionPredicate 
partitionFilter) {
-        return liveRows(table, partitionFilter, null);
+        return liveRows(table, null, partitionFilter, null);
     }
 
     @Nullable
     static RoaringNavigableMap64 liveRows(
             @Nullable FileStoreTable table,
+            @Nullable Snapshot pinnedSnapshot,
             @Nullable PartitionPredicate partitionFilter,
             @Nullable List<Range> rowRanges) {
         if (table == null || !table.coreOptions().deletionVectorsEnabled()) {
             return null;
         }
 
-        @Nullable Snapshot snapshot = TimeTravelUtil.tryTravelOrLatest(table);
+        // Pin to the index scan's snapshot so row-ids and live rows agree.
+        @Nullable
+        Snapshot snapshot =
+                pinnedSnapshot != null ? pinnedSnapshot : 
TimeTravelUtil.tryTravelOrLatest(table);
         if (snapshot == null) {
             return null;
         }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/VectorScan.java 
b/paimon-core/src/main/java/org/apache/paimon/table/source/VectorScan.java
index 80f667cc0b..8b1da8c1ca 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/source/VectorScan.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/source/VectorScan.java
@@ -18,6 +18,10 @@
 
 package org.apache.paimon.table.source;
 
+import org.apache.paimon.Snapshot;
+
+import javax.annotation.Nullable;
+
 import java.util.List;
 
 /** Vector scan to pre-filter and scan index files. */
@@ -28,5 +32,11 @@ public interface VectorScan {
     /** Plan of vector scan. */
     interface Plan {
         List<VectorSearchSplit> splits();
+
+        /** Snapshot the plan was built against; the read pins live-row 
filtering to it. */
+        @Nullable
+        default Snapshot snapshot() {
+            return null;
+        }
     }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/VectorSearchBuilderImpl.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/VectorSearchBuilderImpl.java
index fb0c972365..d5c2526a72 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/VectorSearchBuilderImpl.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/VectorSearchBuilderImpl.java
@@ -137,7 +137,8 @@ public class VectorSearchBuilderImpl implements 
VectorSearchBuilder {
                     filter,
                     pinnedSnapshot);
         }
-        return new DataEvolutionVectorScan(table, partitionFilter, filter, 
vectorColumn, options);
+        return new DataEvolutionVectorScan(
+                table, partitionFilter, filter, vectorColumn, options, 
pinnedSnapshot);
     }
 
     @Override
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
index 3cb9d18f52..3713e0a38b 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java
@@ -253,6 +253,35 @@ public class VectorSearchBuilderTest extends TableTestBase 
{
         assertThat(readIds(table, result)).containsExactly(2, 3);
     }
 
+    @Test
+    public void testVectorSearchPinsLiveRowFilterToPlanSnapshot() throws 
Exception {
+        catalog.createTable(
+                identifier("vector_search_pinned_snapshot"),
+                vectorSchemaBuilder(VECTOR_FIELD_NAME)
+                        .option(CoreOptions.DELETION_VECTORS_ENABLED.key(), 
"true")
+                        .build(),
+                false);
+        FileStoreTable table = 
getTable(identifier("vector_search_pinned_snapshot"));
+
+        float[][] vectors = {{0.0f, 0.0f}, {1.0f, 0.0f}, {2.0f, 0.0f}, {3.0f, 
0.0f}};
+        writeVectors(table, vectors);
+        buildAndCommitIndex(table, vectors);
+
+        VectorSearchBuilder builder =
+                table.newVectorSearchBuilder()
+                        .withVector(new float[] {0.0f, 0.0f})
+                        .withLimit(4)
+                        .withVectorColumn(VECTOR_FIELD_NAME);
+        VectorScan.Plan plan = builder.newVectorScan().scan();
+
+        // Delete row 0 after planning. The read stays pinned to the plan's
+        // snapshot, so row 0 (live when planned) is still returned.
+        commitDeletionVectors(table, 0L);
+
+        GlobalIndexResult result = builder.newVectorRead().read(plan);
+        assertThat(result.results()).contains(0L);
+    }
+
     @Test
     public void testVectorLiveRowPlanningSkipsUnindexedDeletionVectors() 
throws Exception {
         catalog.createTable(
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/vectorsearch/FlinkDataEvolutionVectorRead.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/vectorsearch/FlinkDataEvolutionVectorRead.java
index e1ea1e0d6a..0e2b9a9df2 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/vectorsearch/FlinkDataEvolutionVectorRead.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/vectorsearch/FlinkDataEvolutionVectorRead.java
@@ -87,6 +87,8 @@ public class FlinkDataEvolutionVectorRead extends 
DataEvolutionVectorRead {
 
     @Override
     public GlobalIndexResult read(VectorScan.Plan plan) {
+        // Pin the driver-side live-row / scalar pre-filter to the plan's 
snapshot.
+        this.planSnapshot = plan.snapshot();
         List<IndexVectorSearchSplit> indexSplits = new ArrayList<>();
         List<RawVectorSearchSplit> rawSplits = new ArrayList<>();
         splitSearchSplits(plan.splits(), indexSplits, rawSplits);

Reply via email to