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 96458b5daf [python] Fix extra-field global index padding (#8385)
96458b5daf is described below

commit 96458b5daf14a10d017f733b1ee1606437e8444f
Author: QuakeWang <[email protected]>
AuthorDate: Tue Jun 30 17:21:10 2026 +0800

    [python] Fix extra-field global index padding (#8385)
    
    PyPaimon `GlobalIndexScanner` flattened multi-column global index
    metadata by field id. Extra-field scans therefore lost the primary-index
    group boundary used by Java, and padding was applied as an independent
    reader.
    
    That caused two correctness issues:
    - Groups with different or non-contiguous coverage could drop valid hits
    from ranges not indexed by a shorter group.
    - Padding could convert a real reader's `None` result, meaning the
    predicate cannot be evaluated by the index, into concrete row-id hits.
    
    This PR groups index files by primary field, computes extra-field
    padding from union coverage minus each group's indexed coverage, and
    applies padding only after the real group reader returns a non-`None`
    result.
---
 .../pypaimon/globalindex/global_index_scanner.py   | 145 +++++++++++++++--
 .../pypaimon/tests/vector_search_filter_test.py    | 180 +++++++++++++++++++++
 2 files changed, 311 insertions(+), 14 deletions(-)

diff --git a/paimon-python/pypaimon/globalindex/global_index_scanner.py 
b/paimon-python/pypaimon/globalindex/global_index_scanner.py
index 7f2470d382..2755082550 100644
--- a/paimon-python/pypaimon/globalindex/global_index_scanner.py
+++ b/paimon-python/pypaimon/globalindex/global_index_scanner.py
@@ -22,7 +22,7 @@ from typing import Collection, Optional
 
 from pypaimon.globalindex.global_index_evaluator import GlobalIndexEvaluator
 from pypaimon.globalindex.global_index_meta import GlobalIndexIOMeta
-from pypaimon.globalindex.global_index_reader import GlobalIndexReader
+from pypaimon.globalindex.global_index_reader import GlobalIndexReader, 
_map_future
 from pypaimon.globalindex.global_index_result import GlobalIndexResult
 from pypaimon.common.options.core_options import CoreOptions
 from pypaimon.common.options.options import Options
@@ -64,13 +64,15 @@ class GlobalIndexScanner:
 
     def _create_evaluator(self, fields, file_io, index_path, index_files):
         index_metas = {}
+        extra_index_metas = {}
         for index_file in index_files:
             global_index_meta = index_file.global_index_meta
             if global_index_meta is None:
                 continue
 
             index_type = index_file.index_type
-            field_ids = [global_index_meta.index_field_id]
+            index_field_id = global_index_meta.index_field_id
+            field_ids = [index_field_id]
             if global_index_meta.extra_field_ids is not None:
                 field_ids.extend(global_index_meta.extra_field_ids)
 
@@ -83,21 +85,48 @@ class GlobalIndexScanner:
             range_key = Range(
                 global_index_meta.row_range_start,
                 global_index_meta.row_range_end)
-            for field_id in field_ids:
-                if field_id not in index_metas:
-                    index_metas[field_id] = {}
-                if index_type not in index_metas[field_id]:
-                    index_metas[field_id][index_type] = {}
-                if range_key not in index_metas[field_id][index_type]:
-                    index_metas[field_id][index_type][range_key] = []
-                index_metas[field_id][index_type][range_key].append(io_meta)
+            group = index_metas.get(index_field_id)
+            if group is None:
+                group = _IndexMetaFileGroup(index_field_id, field_ids)
+                index_metas[index_field_id] = group
+                for extra_field_id in field_ids[1:]:
+                    extra_index_metas.setdefault(extra_field_id, 
[]).append(group)
+            elif group.field_ids != tuple(field_ids):
+                raise ValueError(
+                    "Primary field %s owns multiple indexes with different "
+                    "columns %s and %s; a primary column can own at most one "
+                    "index." % (index_field_id, list(group.field_ids), 
field_ids)
+                )
+            group.add_file(index_type, range_key, io_meta)
 
         executor = self._executor
         options = self._options
 
         def readers_function(field: DataField) -> 
Collection[GlobalIndexReader]:
-            return _create_readers(
-                file_io, index_path, index_metas.get(field.id), field, 
executor, options)
+            group = index_metas.get(field.id)
+            if group is not None:
+                return _create_readers(
+                    file_io, index_path, group.metas, field, executor, options)
+
+            extra_groups = extra_index_metas.get(field.id)
+            if not extra_groups:
+                return []
+            union_coverage = Range.sort_and_merge_overlap(
+                [
+                    range_key
+                    for group in extra_groups
+                    for range_key in group.coverage_ranges
+                ],
+                True,
+            )
+            readers = []
+            for group in extra_groups:
+                pad_ranges = _exclude_ranges(union_coverage, 
group.coverage_ranges)
+                readers.extend(
+                    _create_readers(
+                        file_io, index_path, group.metas, field, executor,
+                        options, pad_ranges=pad_ranges))
+            return readers
 
         return GlobalIndexEvaluator(fields, readers_function)
 
@@ -199,6 +228,80 @@ class GlobalIndexScanner:
         self.close()
 
 
+class _IndexMetaFileGroup:
+    def __init__(self, index_field_id, field_ids):
+        self.index_field_id = index_field_id
+        self.field_ids = tuple(field_ids)
+        self.metas = {}
+        self.coverage_ranges = []
+
+    def add_file(self, index_type, range_key, io_meta):
+        self.coverage_ranges.append(range_key)
+        self.metas.setdefault(index_type, {}).setdefault(range_key, 
[]).append(io_meta)
+
+
+class _PaddingGlobalIndexReader(GlobalIndexReader):
+    def __init__(self, wrapped, padding):
+        self._wrapped = wrapped
+        self._padding = padding
+
+    def _pad(self, future):
+        return _map_future(
+            future,
+            lambda result: None if result is None else 
result.or_(self._padding))
+
+    def visit_equal(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_equal(field_ref, literal))
+
+    def visit_not_equal(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_not_equal(field_ref, literal))
+
+    def visit_less_than(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_less_than(field_ref, literal))
+
+    def visit_less_or_equal(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_less_or_equal(field_ref, literal))
+
+    def visit_greater_than(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_greater_than(field_ref, literal))
+
+    def visit_greater_or_equal(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_greater_or_equal(field_ref, 
literal))
+
+    def visit_is_null(self, field_ref):
+        return self._pad(self._wrapped.visit_is_null(field_ref))
+
+    def visit_is_not_null(self, field_ref):
+        return self._pad(self._wrapped.visit_is_not_null(field_ref))
+
+    def visit_in(self, field_ref, literals):
+        return self._pad(self._wrapped.visit_in(field_ref, literals))
+
+    def visit_not_in(self, field_ref, literals):
+        return self._pad(self._wrapped.visit_not_in(field_ref, literals))
+
+    def visit_starts_with(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_starts_with(field_ref, literal))
+
+    def visit_ends_with(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_ends_with(field_ref, literal))
+
+    def visit_contains(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_contains(field_ref, literal))
+
+    def visit_like(self, field_ref, literal):
+        return self._pad(self._wrapped.visit_like(field_ref, literal))
+
+    def visit_between(self, field_ref, from_v, to_v):
+        return self._pad(self._wrapped.visit_between(field_ref, from_v, to_v))
+
+    def visit_not_between(self, field_ref, from_v, to_v):
+        return self._pad(self._wrapped.visit_not_between(field_ref, from_v, 
to_v))
+
+    def close(self):
+        self._wrapped.close()
+
+
 def _resolve_snapshot(table, snapshot):
     if snapshot is not None:
         return snapshot
@@ -232,7 +335,17 @@ def _core_options(table):
     return options
 
 
-def _create_readers(file_io, index_path, index_type_metas, field, 
executor=None, options=None):
+def _exclude_ranges(base_ranges, excluded_ranges):
+    excluded_ranges = Range.sort_and_merge_overlap(excluded_ranges, True)
+    result = []
+    for base_range in Range.sort_and_merge_overlap(base_ranges, True):
+        result.extend(base_range.exclude(excluded_ranges))
+    return Range.sort_and_merge_overlap(result, True)
+
+
+def _create_readers(
+        file_io, index_path, index_type_metas, field, executor=None,
+        options=None, pad_ranges=None):
     """Create readers for a specific field, dispatched by index_type.
 
     Unknown indexTypes raise — a silent skip would make
@@ -260,7 +373,11 @@ def _create_readers(file_io, index_path, index_type_metas, 
field, executor=None,
                     OffsetGlobalIndexReader(
                         inner, range_key.from_, range_key.to))
         if offset_readers:
-            readers.append(UnionGlobalIndexReader(offset_readers))
+            reader = UnionGlobalIndexReader(offset_readers)
+            if pad_ranges:
+                padding = GlobalIndexResult.from_ranges(pad_ranges)
+                reader = _PaddingGlobalIndexReader(reader, padding)
+            readers.append(reader)
     return readers
 
 
diff --git a/paimon-python/pypaimon/tests/vector_search_filter_test.py 
b/paimon-python/pypaimon/tests/vector_search_filter_test.py
index 58f2fdc26e..19b9bca2d4 100644
--- a/paimon-python/pypaimon/tests/vector_search_filter_test.py
+++ b/paimon-python/pypaimon/tests/vector_search_filter_test.py
@@ -1805,6 +1805,186 @@ class 
VectorSearchMultiShardScalarTest(unittest.TestCase):
         # Must not be empty despite shard_a being empty (no short-circuit).
         self.assertEqual([7], hits)
 
+    def test_extra_field_groups_are_padded_before_and(self):
+        from pypaimon.globalindex.global_index_reader import GlobalIndexReader
+        from pypaimon.globalindex.global_index_scanner import (
+            GlobalIndexScanner,
+        )
+
+        a_field = _field(0, "a")
+        b_field = _field(1, "b")
+        c_field = _field(2, "c")
+
+        short = _entry(None, field_id=0, index_type="btree",
+                       file_name="a-c.index",
+                       row_range_start=0, row_range_end=4).index_file
+        short.global_index_meta.extra_field_ids = [2]
+        long = _entry(None, field_id=1, index_type="btree",
+                      file_name="b-c.index",
+                      row_range_start=0, row_range_end=9).index_file
+        long.global_index_meta.extra_field_ids = [2]
+
+        class _StubReader(GlobalIndexReader):
+            def __init__(self_inner, file_name):
+                self_inner._file_name = file_name
+
+            def visit_equal(self_inner, field_ref, literal):
+                bm = RoaringBitmap64()
+                if self_inner._file_name == "a-c.index":
+                    for row_id in [1, 3, 4]:
+                        bm.add(row_id)
+                else:
+                    for row_id in [1, 3, 7, 8]:
+                        bm.add(row_id)
+                return _completed_future(GlobalIndexResult.create(bm))
+
+            def close(self_inner):
+                pass
+
+        def _stub_create_inner_readers(
+                index_type, file_io, index_path, field, io_metas,
+                executor=None, options=None):
+            return [_StubReader(io_meta.file_name) for io_meta in io_metas]
+
+        with mock.patch(
+                
"pypaimon.globalindex.global_index_scanner._create_inner_readers",
+                side_effect=_stub_create_inner_readers):
+            scanner = GlobalIndexScanner(
+                fields=[a_field, b_field, c_field],
+                file_io=object(),
+                index_path="/unused",
+                index_files=[short, long],
+            )
+            try:
+                result = scanner.scan(
+                    Predicate(method="equal", index=2, field="c",
+                              literals=[42]))
+            finally:
+                scanner.close()
+
+        self.assertIsNotNone(result)
+        # The short group is all-hit padded for rows 5..9 before AND-ing with
+        # the long group. Row 4 is filtered out by the long group, while tail
+        # rows 7 and 8 survive because the short group has not indexed them.
+        self.assertEqual([1, 3, 7, 8], sorted(list(result.results())))
+
+    def test_extra_field_groups_pad_missing_coverage_before_and(self):
+        from pypaimon.globalindex.global_index_reader import GlobalIndexReader
+        from pypaimon.globalindex.global_index_scanner import (
+            GlobalIndexScanner,
+        )
+
+        a_field = _field(0, "a")
+        b_field = _field(1, "b")
+        c_field = _field(2, "c")
+
+        a_early = _entry(None, field_id=0, index_type="btree",
+                         file_name="a-c-early.index",
+                         row_range_start=2, row_range_end=3).index_file
+        a_early.global_index_meta.extra_field_ids = [2]
+        a_late = _entry(None, field_id=0, index_type="btree",
+                        file_name="a-c-late.index",
+                        row_range_start=7, row_range_end=9).index_file
+        a_late.global_index_meta.extra_field_ids = [2]
+        b_full = _entry(None, field_id=1, index_type="btree",
+                        file_name="b-c-full.index",
+                        row_range_start=0, row_range_end=9).index_file
+        b_full.global_index_meta.extra_field_ids = [2]
+
+        class _StubReader(GlobalIndexReader):
+            def __init__(self_inner, file_name):
+                self_inner._file_name = file_name
+
+            def visit_equal(self_inner, field_ref, literal):
+                bm = RoaringBitmap64()
+                if self_inner._file_name == "a-c-early.index":
+                    bm.add(0)  # global 2 after offset
+                elif self_inner._file_name == "a-c-late.index":
+                    for row_id in [0, 1]:  # global 7, 8 after offset
+                        bm.add(row_id)
+                else:
+                    for row_id in [1, 2, 5, 7, 8]:
+                        bm.add(row_id)
+                return _completed_future(GlobalIndexResult.create(bm))
+
+            def close(self_inner):
+                pass
+
+        def _stub_create_inner_readers(
+                index_type, file_io, index_path, field, io_metas,
+                executor=None, options=None):
+            return [_StubReader(io_meta.file_name) for io_meta in io_metas]
+
+        with mock.patch(
+                
"pypaimon.globalindex.global_index_scanner._create_inner_readers",
+                side_effect=_stub_create_inner_readers):
+            scanner = GlobalIndexScanner(
+                fields=[a_field, b_field, c_field],
+                file_io=object(),
+                index_path="/unused",
+                index_files=[a_early, a_late, b_full],
+            )
+            try:
+                result = scanner.scan(
+                    Predicate(method="equal", index=2, field="c",
+                              literals=[42]))
+            finally:
+                scanner.close()
+
+        self.assertIsNotNone(result)
+        # The first group has not indexed [0,1] and [4,6], so these ranges are
+        # neutral under AND. Its indexed ranges still filter normally.
+        self.assertEqual([1, 2, 5, 7, 8], sorted(list(result.results())))
+
+    def test_extra_field_padding_does_not_convert_none_to_hits(self):
+        from pypaimon.globalindex.global_index_reader import GlobalIndexReader
+        from pypaimon.globalindex.global_index_scanner import (
+            GlobalIndexScanner,
+        )
+
+        a_field = _field(0, "a", "STRING")
+        b_field = _field(1, "b", "STRING")
+        c_field = _field(2, "c", "STRING")
+
+        short = _entry(None, field_id=0, index_type="btree",
+                       file_name="a-c.index",
+                       row_range_start=0, row_range_end=4).index_file
+        short.global_index_meta.extra_field_ids = [2]
+        long = _entry(None, field_id=1, index_type="btree",
+                      file_name="b-c.index",
+                      row_range_start=0, row_range_end=9).index_file
+        long.global_index_meta.extra_field_ids = [2]
+
+        class _StubReader(GlobalIndexReader):
+            def visit_contains(self_inner, field_ref, literal):
+                return _completed_future(None)
+
+            def close(self_inner):
+                pass
+
+        def _stub_create_inner_readers(
+                index_type, file_io, index_path, field, io_metas,
+                executor=None, options=None):
+            return [_StubReader() for _ in io_metas]
+
+        with mock.patch(
+                
"pypaimon.globalindex.global_index_scanner._create_inner_readers",
+                side_effect=_stub_create_inner_readers):
+            scanner = GlobalIndexScanner(
+                fields=[a_field, b_field, c_field],
+                file_io=object(),
+                index_path="/unused",
+                index_files=[short, long],
+            )
+            try:
+                result = scanner.scan(
+                    Predicate(method="contains", index=2, field="c",
+                              literals=["x"]))
+            finally:
+                scanner.close()
+
+        self.assertIsNone(result)
+
     def test_tantivy_fulltext_index_is_dispatched_by_scanner(self):
         """Non-btree scalar global indexes (tantivy-fulltext, etc.) must be
         instantiated by GlobalIndexScanner — previously only 'btree' was

Reply via email to