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 87e20cc5c5 [core][python] Fix global-index live-row filter re-adding 
DV-deleted rows across overlapping files (#8939)
87e20cc5c5 is described below

commit 87e20cc5c5dbb32dd572c918796d2887421854ef
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 30 19:46:13 2026 +0800

    [core][python] Fix global-index live-row filter re-adding DV-deleted rows 
across overlapping files (#8939)
---
 .../table/source/GlobalIndexLiveRowFilter.java     | 24 +++++++--
 .../table/source/VectorSearchBuilderTest.java      | 58 ++++++++++++++++++++
 .../table/source/global_index_live_row_filter.py   | 38 ++++++++-----
 .../pypaimon/tests/vector_search_filter_test.py    | 63 ++++++++++++++++++++++
 paimon-python/pypaimon/utils/roaring_bitmap.py     |  4 ++
 5 files changed, 170 insertions(+), 17 deletions(-)

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 6539629251..ebd9c5af3f 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
@@ -66,10 +66,19 @@ class GlobalIndexLiveRowFilter {
             snapshotReader.withRowRanges(rowRanges);
         }
 
+        List<Split> splits = snapshotReader.read().splits();
         RoaringNavigableMap64 liveRows = new RoaringNavigableMap64();
-        for (Split split : snapshotReader.read().splits()) {
+        // Phase 1: union every file's row-id range.
+        for (Split split : splits) {
             if (split instanceof DataSplit) {
-                addLiveRows(table, liveRows, (DataSplit) split);
+                addRowRanges(liveRows, (DataSplit) split);
+            }
+        }
+        // Phase 2: subtract each DV. Ranges are all unioned first (no re-add),
+        // and peak memory stays at one DV at a time.
+        for (Split split : splits) {
+            if (split instanceof DataSplit) {
+                subtractDeletedRows(table, liveRows, (DataSplit) split);
             }
         }
         return liveRows;
@@ -89,7 +98,15 @@ class GlobalIndexLiveRowFilter {
         return includeRows.getLongCardinality() == range.count() ? null : 
includeRows;
     }
 
-    private static void addLiveRows(
+    private static void addRowRanges(RoaringNavigableMap64 liveRows, DataSplit 
split) {
+        for (DataFileMeta file : split.dataFiles()) {
+            if (file.firstRowId() != null) {
+                liveRows.addRange(file.nonNullRowIdRange());
+            }
+        }
+    }
+
+    private static void subtractDeletedRows(
             FileStoreTable table, RoaringNavigableMap64 liveRows, DataSplit 
split) {
         List<DataFileMeta> files = split.dataFiles();
         List<DeletionFile> deletionFiles = split.deletionFiles().orElse(null);
@@ -100,7 +117,6 @@ class GlobalIndexLiveRowFilter {
                 continue;
             }
             long firstRowId = file.nonNullFirstRowId();
-            liveRows.addRange(file.nonNullRowIdRange());
 
             Optional<DeletionVector> deletionVector;
             try {
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 85528388c9..3cb9d18f52 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
@@ -36,6 +36,7 @@ import 
org.apache.paimon.globalindex.testvector.TestVectorGlobalIndexer;
 import org.apache.paimon.globalindex.testvector.TestVectorGlobalIndexerFactory;
 import org.apache.paimon.index.IndexFileMeta;
 import org.apache.paimon.io.CompactIncrement;
+import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.io.DataIncrement;
 import org.apache.paimon.options.Options;
 import org.apache.paimon.partition.PartitionPredicate;
@@ -221,6 +222,37 @@ public class VectorSearchBuilderTest extends TableTestBase 
{
         assertThat(readIds(table, result)).containsExactly(2, 3);
     }
 
+    @Test
+    public void 
testVectorSearchExcludesDeletedRowsAcrossOverlappingPartialColumnFiles()
+            throws Exception {
+        catalog.createTable(
+                identifier("vector_search_overlapping_partial_column"),
+                vectorSchemaBuilder(VECTOR_FIELD_NAME)
+                        .option(CoreOptions.DELETION_VECTORS_ENABLED.key(), 
"true")
+                        .build(),
+                false);
+        FileStoreTable table = 
getTable(identifier("vector_search_overlapping_partial_column"));
+
+        float[][] vectors = {{0.0f, 0.0f}, {1.0f, 0.0f}, {2.0f, 0.0f}, {3.0f, 
0.0f}};
+        writeVectors(table, vectors);
+        // A partial-column file over the same [0,3] row-id range that carries 
no
+        // deletion vector; its range must not re-add the deleted rows.
+        writeOverlappingIdColumn(table, vectors.length);
+        buildAndCommitIndex(table, vectors);
+        commitDeletionVectors(table, 0L, 1L);
+
+        GlobalIndexResult result =
+                table.newVectorSearchBuilder()
+                        .withVector(new float[] {0.0f, 0.0f})
+                        .withLimit(4)
+                        .withVectorColumn(VECTOR_FIELD_NAME)
+                        .executeLocal();
+
+        assertThat(result.results()).contains(2L, 3L);
+        assertThat(result.results()).doesNotContain(0L, 1L);
+        assertThat(readIds(table, result)).containsExactly(2, 3);
+    }
+
     @Test
     public void testVectorLiveRowPlanningSkipsUnindexedDeletionVectors() 
throws Exception {
         catalog.createTable(
@@ -1543,6 +1575,32 @@ public class VectorSearchBuilderTest extends 
TableTestBase {
 
     // ====================== Helper methods ======================
 
+    private void writeOverlappingIdColumn(FileStoreTable table, int count) 
throws Exception {
+        long firstRowId = table.snapshotManager().latestSnapshot().nextRowId() 
- count;
+        RowType idType = 
table.rowType().project(Collections.singletonList("id"));
+        BatchWriteBuilder builder = table.newBatchWriteBuilder();
+        try (BatchTableWrite write = builder.newWrite().withWriteType(idType)) 
{
+            for (int i = 0; i < count; i++) {
+                write.write(GenericRow.of(i));
+            }
+            BatchTableCommit commit = builder.newCommit();
+            List<CommitMessage> messages = write.prepareCommit();
+            assignFirstRowId(messages, firstRowId);
+            commit.commit(messages);
+        }
+    }
+
+    private void assignFirstRowId(List<CommitMessage> messages, long 
firstRowId) {
+        for (CommitMessage message : messages) {
+            CommitMessageImpl impl = (CommitMessageImpl) message;
+            List<DataFileMeta> newFiles = new 
ArrayList<>(impl.newFilesIncrement().newFiles());
+            impl.newFilesIncrement().newFiles().clear();
+            for (DataFileMeta file : newFiles) {
+                
impl.newFilesIncrement().newFiles().add(file.assignFirstRowId(firstRowId));
+            }
+        }
+    }
+
     private void writeVectors(FileStoreTable table, float[][] vectors) throws 
Exception {
         BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
         try (BatchTableWrite write = writeBuilder.newWrite();
diff --git 
a/paimon-python/pypaimon/table/source/global_index_live_row_filter.py 
b/paimon-python/pypaimon/table/source/global_index_live_row_filter.py
index d0581ab7a3..4e7d2bfcdd 100644
--- a/paimon-python/pypaimon/table/source/global_index_live_row_filter.py
+++ b/paimon-python/pypaimon/table/source/global_index_live_row_filter.py
@@ -43,12 +43,20 @@ def live_rows(table, partition_filter=None) -> 
Optional[RoaringBitmap64]:
     if partition_filter is not None:
         read_builder = read_builder.with_partition_filter(partition_filter)
 
+    data_splits = [
+        (split.split if isinstance(split, QueryAuthSplit) else split)
+        for split in read_builder.new_scan().plan().splits()
+    ]
+    data_splits = [s for s in data_splits if isinstance(s, DataSplit)]
+
     rows = RoaringBitmap64()
-    scan = read_builder.new_scan()
-    for split in scan.plan().splits():
-        inner = split.split if isinstance(split, QueryAuthSplit) else split
-        if isinstance(inner, DataSplit):
-            rows = _add_live_rows(table, rows, inner)
+    # Phase 1: union every file's row-id range.
+    for split in data_splits:
+        _add_row_ranges(rows, split)
+    # Phase 2: subtract each DV. Ranges are all unioned first (no re-add), and
+    # peak memory stays at one DV at a time.
+    for split in data_splits:
+        _subtract_deleted_rows(table, rows, split)
     return rows
 
 
@@ -64,14 +72,19 @@ def for_range(live_row_ids: Optional[RoaringBitmap64],
     return None if include.cardinality() == row_range.count() else include
 
 
-def _add_live_rows(table, rows: RoaringBitmap64, split: DataSplit) -> 
RoaringBitmap64:
+def _add_row_ranges(rows: RoaringBitmap64, split: DataSplit) -> None:
+    for data_file in split.files:
+        row_id_range = data_file.row_id_range()
+        if row_id_range is not None:
+            rows.add_range(row_id_range.from_, row_id_range.to)
+
+
+def _subtract_deleted_rows(table, rows: RoaringBitmap64, split: DataSplit) -> 
None:
     deletion_files = split.data_deletion_files or []
     for i, data_file in enumerate(split.files):
-        row_id_range = data_file.row_id_range()
-        if row_id_range is None:
+        if data_file.row_id_range() is None:
             continue
 
-        rows.add_range(row_id_range.from_, row_id_range.to)
         deletion_file = deletion_files[i] if i < len(deletion_files) else None
         if deletion_file is None or deletion_file.cardinality == 0:
             continue
@@ -80,9 +93,8 @@ def _add_live_rows(table, rows: RoaringBitmap64, split: 
DataSplit) -> RoaringBit
         if deletion_vector.is_empty():
             continue
 
-        deleted_rows = RoaringBitmap64()
+        deleted = RoaringBitmap64()
         first_row_id = data_file.first_row_id
         for position in deletion_vector.bit_map():
-            deleted_rows.add(first_row_id + position)
-        rows = RoaringBitmap64.remove_all(rows, deleted_rows)
-    return rows
+            deleted.add(first_row_id + position)
+        rows.remove_all_inplace(deleted)
diff --git a/paimon-python/pypaimon/tests/vector_search_filter_test.py 
b/paimon-python/pypaimon/tests/vector_search_filter_test.py
index d677ba2a81..9b861a626c 100644
--- a/paimon-python/pypaimon/tests/vector_search_filter_test.py
+++ b/paimon-python/pypaimon/tests/vector_search_filter_test.py
@@ -422,6 +422,69 @@ class GlobalIndexLiveRowFilterTest(unittest.TestCase):
         self.assertTrue(calls["new_scan"])
         self.assertEqual([10, 12, 14], rows.to_list())
 
+    def 
test_live_rows_does_not_readd_deleted_rows_across_overlapping_files(self):
+        from pypaimon.read.split import DataSplit
+        from pypaimon.table.source import global_index_live_row_filter
+        from pypaimon.table.source.deletion_file import DeletionFile
+
+        class _File:
+            first_row_id = 10
+            row_count = 5
+
+            def row_id_range(self_inner):
+                return Range(10, 14)
+
+        # Anchor file carries a DV; the overlapping partial-column sibling does
+        # not. Its range must not re-add the rows the anchor's DV removed.
+        deletion_file = DeletionFile("dv", 0, 1, cardinality=2)
+        split = DataSplit(
+            files=[_File(), _File()],
+            partition=None,
+            bucket=0,
+            data_deletion_files=[deletion_file],
+        )
+
+        class _Plan:
+            def splits(self_inner):
+                return [split]
+
+        class _Scan:
+            def plan(self_inner):
+                return _Plan()
+
+        class _Builder:
+            def with_partition_filter(self_inner, predicate):
+                return self_inner
+
+            def new_scan(self_inner):
+                return _Scan()
+
+        class _Options:
+            def deletion_vectors_enabled(self_inner, default=False):
+                return True
+
+        class _Table:
+            options = _Options()
+            file_io = object()
+
+            def new_read_builder(self_inner):
+                return _Builder()
+
+        class _DeletionVector:
+            def is_empty(self_inner):
+                return False
+
+            def bit_map(self_inner):
+                return [1, 3]
+
+        with mock.patch(
+                "pypaimon.table.source.global_index_live_row_filter."
+                "DeletionVector.read",
+                return_value=_DeletionVector()):
+            rows = global_index_live_row_filter.live_rows(_Table())
+
+        self.assertEqual([10, 12, 14], rows.to_list())
+
 
 class VectorReaderFactoryTest(unittest.TestCase):
     """Vector reader factory compatibility."""
diff --git a/paimon-python/pypaimon/utils/roaring_bitmap.py 
b/paimon-python/pypaimon/utils/roaring_bitmap.py
index 4327c3d5f3..41b0a6e4d8 100644
--- a/paimon-python/pypaimon/utils/roaring_bitmap.py
+++ b/paimon-python/pypaimon/utils/roaring_bitmap.py
@@ -123,6 +123,10 @@ class RoaringBitmap64:
         result._data = a._data - b._data
         return result
 
+    def remove_all_inplace(self, other: 'RoaringBitmap64') -> None:
+        """Remove all values contained in ``other`` from this bitmap, in 
place."""
+        self._data -= other._data
+
     def serialize(self) -> bytes:
         """Serialize the bitmap to bytes."""
         return self._data.serialize()

Reply via email to