diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexEvaluator.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexEvaluator.java index 49958cc3be11..a23bd2c76a85 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexEvaluator.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexEvaluator.java @@ -39,6 +39,7 @@ import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Collection; +import java.util.Collections; import java.util.Deque; import java.util.HashSet; import java.util.List; @@ -71,7 +72,15 @@ public Optional evaluate(@Nullable Predicate predicate) { if (predicate == null) { return Optional.empty(); } - return awaitGlobalIndexResult(visitAsync(predicate)); + return await(visitAsync(predicate)).map(Evaluation::result); + } + + /** Evaluate the predicate and return the fields whose supported indexes contributed. */ + public Optional evaluateWithContributingFields(@Nullable Predicate predicate) { + if (predicate == null) { + return Optional.empty(); + } + return await(visitAsync(predicate)); } public Optional evaluateTopN(TopN topN) { @@ -84,11 +93,10 @@ public Optional evaluateTopN(TopN topN) { return Optional.empty(); } checkArgument(readers.size() == 1, "TopN expects one aggregated global index reader."); - return awaitGlobalIndexResult(readers.iterator().next().visitTopN(topN)); + return await(readers.iterator().next().visitTopN(topN)); } - private Optional awaitGlobalIndexResult( - CompletableFuture> future) { + private T await(CompletableFuture future) { try { return future.get(); } catch (InterruptedException e) { @@ -105,14 +113,14 @@ private Optional awaitGlobalIndexResult( } } - private CompletableFuture> visitAsync(Predicate predicate) { + private CompletableFuture> visitAsync(Predicate predicate) { if (predicate instanceof LeafPredicate) { return visitLeafAsync((LeafPredicate) predicate); } return visitCompoundAsync((CompoundPredicate) predicate); } - private CompletableFuture> visitLeafAsync(LeafPredicate predicate) { + private CompletableFuture> visitLeafAsync(LeafPredicate predicate) { Optional fieldRefOptional = predicate.fieldRefOptional(); if (!fieldRefOptional.isPresent()) { return CompletableFuture.completedFuture(Optional.empty()); @@ -145,18 +153,20 @@ private CompletableFuture> visitLeafAsync(LeafPredic compoundResult = childResult; } if (compoundResult.get().results().isEmpty()) { - return compoundResult; + break; } } - return compoundResult; + return compoundResult.map( + result -> + new Evaluation(result, Collections.singleton(fieldId))); }); } - private CompletableFuture> visitCompoundAsync( + private CompletableFuture> visitCompoundAsync( CompoundPredicate predicate) { List children = pruneRedundantIsNotNullForAnd(flattenChildren(predicate), predicate); - List>> childFutures = + List>> childFutures = new ArrayList<>(children.size()); for (Predicate child : children) { childFutures.add(visitAsync(child)); @@ -165,40 +175,67 @@ private CompletableFuture> visitCompoundAsync( return CompletableFuture.allOf(childFutures.toArray(new CompletableFuture[0])) .thenApply( v -> { - List> results = new ArrayList<>(); - for (CompletableFuture> f : childFutures) { + List> results = new ArrayList<>(); + for (CompletableFuture> f : childFutures) { results.add(f.join()); } return combineResults(results, predicate); }); } - private Optional combineResults( - List> results, CompoundPredicate predicate) { + private Optional combineResults( + List> results, CompoundPredicate predicate) { + Set contributingFieldIds = new HashSet<>(); if (predicate.function() instanceof Or) { GlobalIndexResult compoundResult = GlobalIndexResult.createEmpty(); - for (Optional childResult : results) { - if (!childResult.isPresent()) { + for (Optional child : results) { + if (!child.isPresent()) { return Optional.empty(); } - compoundResult = compoundResult.or(childResult.get()); + compoundResult = compoundResult.or(child.get().result()); + contributingFieldIds.addAll(child.get().contributingFieldIds()); } - return Optional.of(compoundResult); + return Optional.of(new Evaluation(compoundResult, contributingFieldIds)); } else { Optional compoundResult = Optional.empty(); - for (Optional childResult : results) { - if (childResult.isPresent()) { + for (Optional child : results) { + if (child.isPresent()) { if (compoundResult.isPresent()) { - compoundResult = Optional.of(compoundResult.get().and(childResult.get())); + compoundResult = + Optional.of(compoundResult.get().and(child.get().result())); } else { - compoundResult = childResult; + compoundResult = Optional.of(child.get().result()); } + contributingFieldIds.addAll(child.get().contributingFieldIds()); } if (compoundResult.isPresent() && compoundResult.get().results().isEmpty()) { - return compoundResult; + break; } } - return compoundResult; + return compoundResult.map(result -> new Evaluation(result, contributingFieldIds)); + } + } + + /** + * Matches and fields whose supported index results contributed; discarded branches excluded. + */ + public static final class Evaluation { + + private final GlobalIndexResult result; + private final Set contributingFieldIds; + + private Evaluation(GlobalIndexResult result, Collection contributingFieldIds) { + this.result = result; + this.contributingFieldIds = + Collections.unmodifiableSet(new HashSet<>(contributingFieldIds)); + } + + public GlobalIndexResult result() { + return result; + } + + public Set contributingFieldIds() { + return contributingFieldIds; } } diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java index 42ccc570e271..126ff3682c14 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java @@ -225,6 +225,57 @@ void testOrReturnsEmptyWhenChildUnsupported() { evaluator.close(); } + @Test + void testAndTracksOnlyEvaluatedFields() { + executor = Executors.newFixedThreadPool(2); + RowType rowType = rowType(); + + GlobalIndexEvaluator evaluator = + new GlobalIndexEvaluator( + rowType, + fieldId -> + fieldId == 0 + ? Collections.singletonList(readerReturning(resultOf(42))) + : Collections.emptyList()); + PredicateBuilder builder = new PredicateBuilder(rowType); + Predicate predicate = PredicateBuilder.and(builder.equal(0, 42), builder.equal(1, 99)); + + Optional evaluation = + evaluator.evaluateWithContributingFields(predicate); + + assertThat(evaluation).isPresent(); + assertThat(evaluation.get().contributingFieldIds()).containsExactly(0); + assertBitmapContainsExactly(evaluation.get().result().results(), 42L); + evaluator.close(); + } + + @Test + void testDiscardedOrBranchDoesNotContributeFields() { + executor = Executors.newFixedThreadPool(2); + RowType rowType = rowType(); + + GlobalIndexEvaluator evaluator = + new GlobalIndexEvaluator( + rowType, + fieldId -> + fieldId == 0 || fieldId == 2 + ? Collections.singletonList(readerReturning(resultOf(42))) + : Collections.emptyList()); + PredicateBuilder builder = new PredicateBuilder(rowType); + Predicate predicate = + PredicateBuilder.and( + PredicateBuilder.or(builder.equal(0, 42), builder.equal(1, 99)), + builder.equal(2, 42)); + + Optional evaluation = + evaluator.evaluateWithContributingFields(predicate); + + assertThat(evaluation).isPresent(); + assertThat(evaluation.get().contributingFieldIds()).containsExactly(2); + assertBitmapContainsExactly(evaluation.get().result().results(), 42L); + evaluator.close(); + } + @Test void testAndWithEmptyResultShortCircuits() { executor = Executors.newFixedThreadPool(2); diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java index 059d66160887..8929ef4c7bb6 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java @@ -318,12 +318,17 @@ private Optional evalGlobalIndex() { try (DataEvolutionGlobalIndexScanner scanner = optionalScanner.get()) { long lookupStart = System.nanoTime(); - Optional result = scanner.scan(globalIndexFilter); + Optional result = + scanner.scanWithCoverage(globalIndexFilter); long lookupDuration = System.nanoTime() - lookupStart; if (result.isPresent()) { long coverageStart = System.nanoTime(); GlobalIndexResult finalResult = - result.get().or(scanner.unindexedRows(globalIndexFilter)); + result.get() + .result() + .or( + scanner.unindexedRowsForContributingFields( + result.get().contributingFieldIds())); long coverageDuration = System.nanoTime() - coverageStart; long totalDuration = System.nanoTime() - totalStart; LOG.info( 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 498130a11dd3..9d6acaeddeda 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 @@ -372,6 +372,10 @@ public Optional scan(Predicate predicate) { return globalIndexEvaluator.evaluate(predicate); } + public Optional scanWithCoverage(Predicate predicate) { + return globalIndexEvaluator.evaluateWithContributingFields(predicate); + } + public Optional scan(TopN topN) { if (!isSupportedTopN(topN)) { return Optional.empty(); @@ -394,6 +398,15 @@ public GlobalIndexResult unindexedRows(Predicate predicate) { return GlobalIndexResult.create(rows); } + public GlobalIndexResult unindexedRowsForContributingFields( + Collection contributingFieldIds) { + RoaringNavigableMap64 rows = new RoaringNavigableMap64(); + for (Range range : coverage.unindexedRanges(contributingFieldIds)) { + rows.addRange(range); + } + return GlobalIndexResult.create(rows); + } + public GlobalIndexResult unindexedRows(TopN topN) { String fieldName = topN.orders().get(0).field().name(); RoaringNavigableMap64 rows = new RoaringNavigableMap64(); 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 75fbb2b5b7db..b35a1eb45c4e 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 @@ -25,6 +25,7 @@ import org.apache.paimon.data.InternalVector; import org.apache.paimon.fs.FileIO; import org.apache.paimon.globalindex.DataEvolutionGlobalIndexScanner; +import org.apache.paimon.globalindex.GlobalIndexEvaluator; import org.apache.paimon.globalindex.GlobalIndexIOMeta; import org.apache.paimon.globalindex.GlobalIndexReader; import org.apache.paimon.globalindex.GlobalIndexResult; @@ -221,12 +222,14 @@ protected RoaringNavigableMap64 rawPreFilter(List splits) RoaringNavigableMap64 include = new RoaringNavigableMap64(); try (DataEvolutionGlobalIndexScanner scanner = optionalScanner.get()) { - Optional result = scanner.scan(filter); + Optional result = scanner.scanWithCoverage(filter); if (!result.isPresent()) { return null; } - include.or(result.get().results()); - include.or(scanner.unindexedRows(filter).results()); + include.or(result.get().result().results()); + include.or( + scanner.unindexedRowsForContributingFields(result.get().contributingFieldIds()) + .results()); } catch (IOException e) { throw new RuntimeException(e); } diff --git a/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java b/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java index 375cac3989ff..9d855701305b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java @@ -151,6 +151,32 @@ public void testBTreeGlobalIndexWithCoreScan() throws Exception { assertThat(readF1).containsExactly("a200", "a300", "a400", "a56789"); } + @Test + public void testFullSearchIgnoresUnindexedAndResidualForCoverage() throws Exception { + write(100L); + createIndex("f1"); + + FileStoreTable table = + tableWithSearchMode((FileStoreTable) catalog.getTable(identifier()), "full"); + PredicateBuilder builder = new PredicateBuilder(table.rowType()); + Predicate predicate = + PredicateBuilder.and( + builder.equal(1, BinaryString.fromString("a42")), + builder.equal(2, BinaryString.fromString("b42"))); + ReadBuilder readBuilder = table.newReadBuilder().withFilter(predicate); + + TableScan.Plan plan = readBuilder.newScan().plan(); + + assertThat(plan.splits()).allMatch(IndexedSplit.class::isInstance); + assertThat( + plan.splits().stream() + .map(IndexedSplit.class::cast) + .flatMap(split -> split.rowRanges().stream()) + .collect(Collectors.toList())) + .containsExactly(new Range(42, 42)); + assertThat(readF1(readBuilder, plan)).containsExactly("a42"); + } + @Test public void testBTreeGlobalIndexTopNCandidatesAcrossRanges() throws Exception { write(100L); diff --git a/paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py b/paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py index 192c2e99bb8d..dde7f24f570e 100644 --- a/paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py +++ b/paimon-python/pypaimon/globalindex/data_evolution_global_index_scanner.py @@ -18,9 +18,12 @@ """Scanner for shard-based global indexes on data-evolution tables.""" from concurrent.futures import ThreadPoolExecutor -from typing import Collection, Optional +from typing import Collection, List, Optional -from pypaimon.globalindex.global_index_evaluator import GlobalIndexEvaluator +from pypaimon.globalindex.global_index_evaluator import ( + GlobalIndexEvaluation, + GlobalIndexEvaluator, +) from pypaimon.globalindex.global_index_meta import GlobalIndexIOMeta from pypaimon.globalindex.global_index_reader import GlobalIndexReader, _map_future from pypaimon.globalindex.global_index_result import GlobalIndexResult @@ -33,6 +36,9 @@ from pypaimon.utils.range import Range +_SUPPORTED_SCALAR_INDEX_TYPES = frozenset(('btree', 'bitmap')) + + class DataEvolutionGlobalIndexScanner: """Scanner for shard-based global indexes.""" @@ -48,6 +54,7 @@ def __init__( snapshot=None, partition_filter=None, ): + index_files = _supported_scalar_index_files(index_files) self._options = options or CoreOptions(Options.from_none()) self._executor = ThreadPoolExecutor( max_workers=thread_num or 32 @@ -103,24 +110,33 @@ def _create_evaluator(self, fields, file_io, index_path, index_files): options = self._options def readers_function(field: DataField) -> Collection[GlobalIndexReader]: + groups = [] group = index_metas.get(field.id) if group is not None: - return _create_readers( - file_io, index_path, group.metas, field, executor, options) + groups.append(group) extra_groups = extra_index_metas.get(field.id) - if not extra_groups: + if extra_groups: + groups.extend( + extra_group + for extra_group in extra_groups + if extra_group not in groups + ) + if not groups: return [] + if len(groups) == 1: + return _create_readers( + file_io, index_path, groups[0].metas, field, executor, options) union_coverage = Range.sort_and_merge_overlap( [ range_key - for group in extra_groups + for group in groups for range_key in group.coverage_ranges ], True, ) readers = [] - for group in extra_groups: + for group in groups: pad_ranges = _exclude_ranges(union_coverage, group.coverage_ranges) readers.extend( _create_readers( @@ -145,6 +161,7 @@ def create(table, index_files=None, partition_filter=None, predicate=None, from pypaimon.index.index_file_handler import IndexFileHandler if index_files is not None: + index_files = _supported_scalar_index_files(index_files) if len(index_files) == 0: return None core_options = _core_options(table) @@ -172,6 +189,8 @@ def index_file_filter(entry): if partition_filter is not None: if not partition_filter.test(entry.partition): return False + if not is_supported_scalar_index(entry.index_file): + return False global_index_meta = entry.index_file.global_index_meta if global_index_meta is None: return False @@ -209,14 +228,32 @@ def scan(self, predicate: Optional[Predicate]) -> Optional[GlobalIndexResult]: """Scan the global index with the given predicate.""" return self._evaluator.evaluate(predicate) + def scan_with_coverage( + self, predicate: Optional[Predicate] + ) -> Optional[GlobalIndexEvaluation]: + return self._evaluator.evaluate_with_contributing_fields(predicate) + def unindexed_rows(self, predicate: Optional[Predicate], - search_mode=None) -> GlobalIndexResult: + search_mode=None, + contributing_field_ids=None) -> GlobalIndexResult: """Return coarse row ids not covered by global indexes.""" + return GlobalIndexResult.from_ranges(self.unindexed_ranges( + predicate, + search_mode=search_mode, + contributing_field_ids=contributing_field_ids, + )) + + def unindexed_ranges(self, predicate: Optional[Predicate], + search_mode=None, + contributing_field_ids=None) -> List[Range]: + """Return row ranges not covered by global indexes.""" if self._coverage is None: - return GlobalIndexResult.create_empty() - return GlobalIndexResult.from_ranges( - self._coverage.unindexed_ranges( - self._fields, predicate, search_mode=search_mode)) + return [] + if contributing_field_ids is not None: + return self._coverage.unindexed_ranges( + contributing_field_ids, search_mode=search_mode) + return self._coverage.unindexed_ranges( + self._fields, predicate, search_mode=search_mode) def close(self): """Close the scanner and release resources.""" @@ -330,6 +367,18 @@ def _resolve_snapshot(table, snapshot): return snapshot_manager.get_latest_snapshot() +def is_supported_scalar_index(index_file): + return ( + index_file.global_index_meta is not None + and index_file.index_type in _SUPPORTED_SCALAR_INDEX_TYPES + ) + + +def _supported_scalar_index_files(index_files): + return [index_file for index_file in index_files + if is_supported_scalar_index(index_file)] + + def _core_options(table): options = getattr(table, "options", None) if options is None: @@ -413,15 +462,5 @@ def _create_inner_readers( fallback_scan_max_size=core_options.bitmap_index_fallback_scan_max_size(), )] - from pypaimon.globalindex.full_text import ( - FULL_TEXT_IDENTIFIER, - NativeFullTextGlobalIndexReader, - ) - if index_type == FULL_TEXT_IDENTIFIER: - return [ - NativeFullTextGlobalIndexReader(file_io, index_path, [io_meta]) - for io_meta in io_metas - ] - raise ValueError( "Unsupported global-index type in scanner: '%s'" % index_type) diff --git a/paimon-python/pypaimon/globalindex/global_index_evaluator.py b/paimon-python/pypaimon/globalindex/global_index_evaluator.py index 9fdf64a607cd..b0c8a0ed2ca7 100644 --- a/paimon-python/pypaimon/globalindex/global_index_evaluator.py +++ b/paimon-python/pypaimon/globalindex/global_index_evaluator.py @@ -20,7 +20,7 @@ import threading from collections import deque from concurrent.futures import Future -from typing import Callable, Collection, Dict, List, Optional +from typing import Callable, Collection, Dict, FrozenSet, List, NamedTuple, Optional from pypaimon.globalindex.global_index_reader import GlobalIndexReader, FieldRef from pypaimon.globalindex.global_index_result import GlobalIndexResult @@ -28,6 +28,13 @@ from pypaimon.schema.data_types import DataField +class GlobalIndexEvaluation(NamedTuple): + """Matches and fields whose supported indexes contributed.""" + + result: GlobalIndexResult + contributing_field_ids: FrozenSet[int] + + class GlobalIndexEvaluator: """Predicate evaluator for filtering data using global indexes. @@ -51,8 +58,17 @@ def evaluate( ) -> Optional[GlobalIndexResult]: if predicate is None: return None - future = self._visit_async(predicate) - return future.result() + evaluation = self._visit_async(predicate).result() + return evaluation.result if evaluation is not None else None + + def evaluate_with_contributing_fields( + self, + predicate: Optional[Predicate] + ) -> Optional[GlobalIndexEvaluation]: + """Return matches and fields whose supported indexes contributed.""" + if predicate is None: + return None + return self._visit_async(predicate).result() def _visit_async(self, predicate) -> Future: if isinstance(predicate, Predicate) and predicate.method in ('and', 'or'): @@ -94,7 +110,8 @@ def on_done(_): if remaining[0] == 0: try: all_done.set_result( - self._combine_reader_results(reader_futures) + self._combine_reader_results( + reader_futures, field_id) ) except Exception as e: all_done.set_exception(e) @@ -105,8 +122,8 @@ def on_done(_): return all_done def _combine_reader_results( - self, reader_futures: List[Future] - ) -> Optional[GlobalIndexResult]: + self, reader_futures: List[Future], field_id: int, + ) -> Optional[GlobalIndexEvaluation]: compound_result: Optional[GlobalIndexResult] = None for f in reader_futures: child_result = f.result() @@ -117,8 +134,10 @@ def _combine_reader_results( else: compound_result = child_result if compound_result.is_empty(): - return compound_result - return compound_result + break + if compound_result is None: + return None + return GlobalIndexEvaluation(compound_result, frozenset([field_id])) def _visit_compound_async(self, predicate: Predicate) -> Future: children = self._flatten_children(predicate.method, predicate.literals) @@ -150,26 +169,34 @@ def on_done(_): return all_done def _combine_results( - self, results: List[Optional[GlobalIndexResult]], method: str - ) -> Optional[GlobalIndexResult]: + self, results: List[Optional[GlobalIndexEvaluation]], method: str + ) -> Optional[GlobalIndexEvaluation]: + contributing_field_ids = set() if method == 'or': compound_result = GlobalIndexResult.create_empty() - for child_result in results: - if child_result is None: + for child in results: + if child is None: return None - compound_result = compound_result.or_(child_result) - return compound_result + compound_result = compound_result.or_(child.result) + contributing_field_ids.update(child.contributing_field_ids) + return GlobalIndexEvaluation(compound_result, + frozenset(contributing_field_ids)) else: compound_result: Optional[GlobalIndexResult] = None - for child_result in results: - if child_result is not None: + for child in results: + if child is not None: if compound_result is not None: - compound_result = compound_result.and_(child_result) + compound_result = compound_result.and_(child.result) else: - compound_result = child_result + compound_result = child.result + contributing_field_ids.update( + child.contributing_field_ids) if compound_result is not None and compound_result.is_empty(): - return compound_result - return compound_result + break + if compound_result is None: + return None + return GlobalIndexEvaluation(compound_result, + frozenset(contributing_field_ids)) def _flatten_children(self, method: str, children) -> list: result = [] diff --git a/paimon-python/pypaimon/read/scanner/file_scanner.py b/paimon-python/pypaimon/read/scanner/file_scanner.py index 3283a51bf905..e92205405de2 100755 --- a/paimon-python/pypaimon/read/scanner/file_scanner.py +++ b/paimon-python/pypaimon/read/scanner/file_scanner.py @@ -18,12 +18,13 @@ import logging import os import time -from typing import Callable, Dict, List, Optional, Set, Tuple +from typing import Callable, Dict, List, NamedTuple, Optional, Set, Tuple logger = logging.getLogger(__name__) from pypaimon.common.predicate import Predicate from pypaimon.globalindex import ScoredGlobalIndexResult +from pypaimon.globalindex.global_index_result import GlobalIndexResult from pypaimon.manifest.index_manifest_file import IndexManifestFile from pypaimon.manifest.manifest_file_manager import ManifestFileManager from pypaimon.manifest.manifest_list_manager import ManifestListManager @@ -56,11 +57,16 @@ from pypaimon.table.bucket_mode import BucketMode from pypaimon.table.special_fields import SpecialFields from pypaimon.table.source.deletion_file import DeletionFile +from pypaimon.utils.range import Range + + +class _GlobalIndexPlanningResult(NamedTuple): + indexed_result: GlobalIndexResult + unindexed_ranges: List[Range] def _row_ranges_from_predicate(predicate: Optional[Predicate]) -> Optional[List]: from pypaimon.table.special_fields import SpecialFields - from pypaimon.utils.range import Range if predicate is None: return None @@ -119,8 +125,6 @@ def _build_early_row_range_filter(row_ranges): if row_ranges is None or not row_ranges: return None - from pypaimon.utils.range import Range - def _filter(record): file_dict = record.get('_FILE') if file_dict is None: @@ -156,8 +160,6 @@ def _filter_manifest_files_by_row_ranges( Returns: Filtered list of manifest files """ - from pypaimon.utils.range import Range - filtered_files = [] for manifest in manifest_files: min_row_id = manifest.min_row_id @@ -270,6 +272,7 @@ def __init__( self.data_evolution = options.data_evolution_enabled() self.deletion_vectors_enabled = options.deletion_vectors_enabled() self._global_index_result = None + self._row_ranges = None self._scanned_snapshot = None self._scanned_snapshot_id = None # Opt-in scan-plan tracking. Stays ``None`` for the read hot path; @@ -430,7 +433,7 @@ def _apply_primary_key_sorted_indexes(self, splits): return list(PrimaryKeySortedIndexResult(evaluated).splits) def _create_data_evolution_split_generator(self): - row_ranges = None + row_ranges = getattr(self, '_row_ranges', None) score_getter = None # Fetch snapshot once and share with global index evaluation to avoid # a duplicate /snapshot REST round-trip (#7513). @@ -438,15 +441,36 @@ def _create_data_evolution_split_generator(self): self._scanned_snapshot = snapshot self._scanned_snapshot_id = snapshot.id if snapshot else None - global_index_result = self._global_index_result if self._global_index_result is not None \ - else self._eval_global_index(snapshot) - if global_index_result is not None: - row_ranges = global_index_result.results().to_range_list() - if isinstance(global_index_result, ScoredGlobalIndexResult): - score_getter = global_index_result.score_getter() + if row_ranges is None: + global_index_plan = self._global_index_result \ + if self._global_index_result is not None \ + else self._eval_global_index(snapshot) + if global_index_plan is not None: + if isinstance(global_index_plan, _GlobalIndexPlanningResult): + global_index_result = global_index_plan.indexed_result + row_ranges = Range.sort_and_merge_overlap( + global_index_result.results().to_range_list() + + global_index_plan.unindexed_ranges, + True, + ) + else: + global_index_result = global_index_plan + row_ranges = global_index_result.results().to_range_list() + if isinstance(global_index_result, ScoredGlobalIndexResult): + score_getter = global_index_result.score_getter() if row_ranges is None and self.predicate is not None: row_ranges = _row_ranges_from_predicate(self.predicate) + if row_ranges is not None and not row_ranges: + return [], DataEvolutionSplitGenerator( + self.table, + self.target_split_size, + self.open_file_cost, + {}, + row_ranges, + score_getter, + ) + # Filter manifest files by row ranges if available if row_ranges is not None: manifest_files = _filter_manifest_files_by_row_ranges(manifest_files, row_ranges) @@ -463,7 +487,7 @@ def _create_data_evolution_split_generator(self): self.open_file_cost, self._deletion_files_map(entries), row_ranges, - score_getter + score_getter, ) def plan_files(self) -> List[ManifestEntry]: @@ -495,12 +519,19 @@ def _eval_global_index(self, snapshot=None): if scanner is None: return None with scanner: - result = scanner.scan(self.predicate) - if result is None: + evaluation = scanner.scan_with_coverage(self.predicate) + if evaluation is None: return None scalar_mode = self.table.options.scalar_index_search_mode() - return result.or_( - scanner.unindexed_rows(self.predicate, search_mode=scalar_mode)) + return _GlobalIndexPlanningResult( + evaluation.result, + scanner.unindexed_ranges( + self.predicate, + search_mode=scalar_mode, + contributing_field_ids=( + evaluation.contributing_field_ids), + ), + ) except Exception: return None @@ -587,9 +618,23 @@ def with_slice(self, start_pos: int, end_pos: int) -> 'FileScanner': return self def with_global_index_result(self, result) -> 'FileScanner': + if self._row_ranges is not None: + raise ValueError( + "with_global_index_result and with_row_ranges are mutually exclusive") self._global_index_result = result return self + def with_row_ranges(self, row_ranges) -> 'FileScanner': + if not self.data_evolution: + raise ValueError("Row ranges are only supported for data evolution tables") + if row_ranges is None: + raise ValueError("row_ranges cannot be None") + if self._global_index_result is not None: + raise ValueError( + "with_row_ranges and with_global_index_result are mutually exclusive") + self._row_ranges = Range.sort_and_merge_overlap(list(row_ranges), True) + return self + def scan_with_stats(self) -> Tuple[Plan, ScanStats]: """Run one scan pass while recording :class:`ScanStats` counters. @@ -620,6 +665,8 @@ def _validate_chunk_shuffle_compat(self) -> None: raise ValueError("chunk_shuffle cannot combine with limit") if self._global_index_result is not None: raise ValueError("chunk_shuffle cannot combine with global index") + if self._row_ranges is not None: + raise ValueError("chunk_shuffle cannot combine with row ranges") # Only partition predicates are allowed: row-level / column-level # predicates would silently shrink each chunk's effective row count, # breaking the chunk_size contract DataLoader callers expect. diff --git a/paimon-python/pypaimon/read/table_scan.py b/paimon-python/pypaimon/read/table_scan.py index 410cae9a4e63..af1a939f2fb0 100755 --- a/paimon-python/pypaimon/read/table_scan.py +++ b/paimon-python/pypaimon/read/table_scan.py @@ -95,8 +95,9 @@ def _native_plan_supported(self) -> bool: def _native_plan_supported_impl(self) -> bool: """Fall back to the Python scanner for scans native can't carry: - shard/slice, chunk-shuffle, global-index, first-row merge-engine (Rust - drops L0), deletion vectors, postpone bucket (drops synthetic buckets), + shard/slice, chunk-shuffle, global-index/row-ranges, first-row + merge-engine (Rust drops L0), deletion vectors, postpone bucket + (drops synthetic buckets), a primary-key table whose trimmed PK is empty (PK equals the partition key; native may mark splits raw-convertible and skip merge), dynamic bucket / cross-partition PK tables (unconfirmed Rust parity), a stale @@ -113,6 +114,7 @@ def _native_plan_supported_impl(self) -> bool: or getattr(fs, 'start_pos_of_this_subtask', None) is not None or getattr(fs, 'chunk_shuffle', None) is not None or getattr(fs, '_global_index_result', None) is not None + or getattr(fs, '_row_ranges', None) is not None or getattr(fs, 'deletion_vectors_enabled', False) or getattr(fs, 'only_read_real_buckets', False)): return False @@ -363,6 +365,10 @@ def with_global_index_result(self, result) -> 'TableScan': self.file_scanner.with_global_index_result(result) return self + def with_row_ranges(self, row_ranges) -> 'TableScan': + self.file_scanner.with_row_ranges(row_ranges) + return self + def with_chunk_shuffle(self, seed: int, chunk_size: int) -> 'TableScan': self.file_scanner.with_chunk_shuffle(seed, chunk_size) return self diff --git a/paimon-python/pypaimon/table/source/full_text_read.py b/paimon-python/pypaimon/table/source/full_text_read.py index 46372c5edb13..036a6d9d3d69 100644 --- a/paimon-python/pypaimon/table/source/full_text_read.py +++ b/paimon-python/pypaimon/table/source/full_text_read.py @@ -197,8 +197,7 @@ def _read_raw_rows(self, raw_row_ranges): projection = [self._text_columns[0].name, SpecialFields.ROW_ID.name] read_builder = read_builder.with_projection(projection) - plan = read_builder.new_scan().with_global_index_result( - GlobalIndexResult.from_ranges(raw_row_ranges)).plan() + plan = read_builder.new_scan().with_row_ranges(raw_row_ranges).plan() return read_builder.new_read().to_arrow(plan.splits()) def _build_raw_index(self, row_ids, texts, row_range_start): diff --git a/paimon-python/pypaimon/table/source/vector_search_read.py b/paimon-python/pypaimon/table/source/vector_search_read.py index 199dc42c93ea..46c45f01f3d0 100644 --- a/paimon-python/pypaimon/table/source/vector_search_read.py +++ b/paimon-python/pypaimon/table/source/vector_search_read.py @@ -167,8 +167,8 @@ def _pre_filter(self, splits, snapshot=None): def _raw_pre_filter(self, splits, snapshot=None): if self._filter is None: return None - raw_rows = _bitmap_of_ranges(_raw_row_ranges(splits)) - if raw_rows.is_empty(): + raw_row_ranges = _raw_row_ranges(splits) + if not raw_row_ranges: return None seen = set() @@ -192,17 +192,21 @@ def _raw_pre_filter(self, splits, snapshot=None): if scanner is None: return None try: - result = scanner.scan(self._filter) - if result is None: + evaluation = scanner.scan_with_coverage(self._filter) + if evaluation is None: return None - include = result.results() - include = RoaringBitmap64.or_( - include, - scanner.unindexed_rows( + include_ranges = evaluation.result.results().to_range_list() + include_ranges.extend( + scanner.unindexed_ranges( self._filter, search_mode=self._table.options.scalar_index_search_mode(), - ).results()) - return RoaringBitmap64.and_(include, raw_rows) + contributing_field_ids=( + evaluation.contributing_field_ids), + )) + return Range.and_( + raw_row_ranges, + Range.sort_and_merge_overlap(include_ranges, True), + ) finally: scanner.close() @@ -322,8 +326,7 @@ def _read_raw_arrow(self, raw_row_ranges, include_filter, snapshot=None): read_builder = read_builder.with_filter(self._filter) read_builder = read_builder.with_projection( self._raw_search_projection(include_filter)) - plan = read_builder.new_scan().with_global_index_result( - GlobalIndexResult.from_ranges(raw_row_ranges)).plan() + plan = read_builder.new_scan().with_row_ranges(raw_row_ranges).plan() return read_builder.new_read().to_arrow(plan.splits()) def _score_raw_vectors(self, candidates, raw_vectors, query_vector, metric, top_k): @@ -639,7 +642,7 @@ def _filtered_raw_row_ranges(raw_row_ranges, pre_filter): return raw_row_ranges return Range.and_( raw_row_ranges, - Range.sort_and_merge_overlap(pre_filter.to_range_list(), True), + Range.sort_and_merge_overlap(pre_filter, True), ) @@ -678,13 +681,6 @@ def _bitmap_of_range(row_range): return bitmap -def _bitmap_of_ranges(ranges): - bitmap = RoaringBitmap64() - for row_range in ranges: - bitmap.add_range(row_range.from_, row_range.to) - return bitmap - - def _to_vector_list(value): if hasattr(value, "to_list"): return value.to_list() diff --git a/paimon-python/pypaimon/table/source/vector_search_scan.py b/paimon-python/pypaimon/table/source/vector_search_scan.py index 83e381d16b1b..8c94390b511f 100644 --- a/paimon-python/pypaimon/table/source/vector_search_scan.py +++ b/paimon-python/pypaimon/table/source/vector_search_scan.py @@ -22,6 +22,9 @@ from pypaimon.common.options.core_options import GlobalIndexSearchMode from pypaimon.globalindex.data_evolution_global_index_coverage import DataEvolutionGlobalIndexCoverage +from pypaimon.globalindex.data_evolution_global_index_scanner import ( + is_supported_scalar_index, +) from pypaimon.table.source.vector_search_split import ( IndexVectorSearchSplit, RawVectorSearchSplit, @@ -121,6 +124,8 @@ def index_file_filter(entry): field_id = global_index_meta.index_field_id if vector_column.id == field_id: return True + if not is_supported_scalar_index(entry.index_file): + return False for filter_field_id in filter_field_ids: if contains_field(global_index_meta, filter_field_id): return True @@ -154,7 +159,8 @@ def index_file_filter(entry): for index_file in all_index_files: meta = index_file.global_index_meta assert meta is not None - if meta.index_field_id == vector_column.id: + if (meta.index_field_id == vector_column.id + or not is_supported_scalar_index(index_file)): continue scalar_range = Range(meta.row_range_start, meta.row_range_end) if range_key.overlaps(scalar_range): @@ -182,6 +188,7 @@ def index_file_filter(entry): f for f in all_index_files if f.global_index_meta is not None and f.global_index_meta.index_field_id != vector_column.id + and is_supported_scalar_index(f) ] if self._filter is not None: scalar_unindexed_ranges = DataEvolutionGlobalIndexCoverage( @@ -231,7 +238,9 @@ def _scalar_index_files_for_ranges(all_index_files, row_ranges, vector_field_id) scalar_files = [] for index_file in all_index_files: meta = index_file.global_index_meta - if meta is None or meta.index_field_id == vector_field_id: + if (meta is None + or meta.index_field_id == vector_field_id + or not is_supported_scalar_index(index_file)): continue if _has_intersection(row_ranges, Range(meta.row_range_start, meta.row_range_end)): scalar_files.append(index_file) diff --git a/paimon-python/pypaimon/tests/global_index_evaluator_test.py b/paimon-python/pypaimon/tests/global_index_evaluator_test.py index 7f3b49388a3a..d979275ec002 100644 --- a/paimon-python/pypaimon/tests/global_index_evaluator_test.py +++ b/paimon-python/pypaimon/tests/global_index_evaluator_test.py @@ -173,6 +173,64 @@ def readers_fn(field): self.assertIsNone(result) evaluator.close() + def test_and_tracks_only_evaluated_fields(self): + fields = _make_fields() + indexed = GlobalIndexResult.from_range(Range(42, 42)) + + evaluator = GlobalIndexEvaluator( + fields, + lambda field: [StubGlobalIndexReader(indexed)] + if field.id == 0 else [], + ) + predicate = Predicate( + method='and', index=None, field=None, + literals=[ + Predicate(method='equal', index=0, field='a', literals=[42]), + Predicate(method='equal', index=1, field='b', literals=[99]), + ], + ) + + evaluation = evaluator.evaluate_with_contributing_fields(predicate) + + self.assertIsNotNone(evaluation) + self.assertEqual(frozenset([0]), evaluation.contributing_field_ids) + self.assertEqual([Range(42, 42)], + evaluation.result.results().to_range_list()) + evaluator.close() + + def test_discarded_or_branch_does_not_contribute_fields(self): + fields = _make_fields() + indexed = GlobalIndexResult.from_range(Range(42, 42)) + + evaluator = GlobalIndexEvaluator( + fields, + lambda field: [StubGlobalIndexReader(indexed)] + if field.id in (0, 2) else [], + ) + predicate = Predicate( + method='and', index=None, field=None, + literals=[ + Predicate( + method='or', index=None, field=None, + literals=[ + Predicate(method='equal', index=0, field='a', + literals=[42]), + Predicate(method='equal', index=1, field='b', + literals=[99]), + ], + ), + Predicate(method='equal', index=2, field='c', literals=[42]), + ], + ) + + evaluation = evaluator.evaluate_with_contributing_fields(predicate) + + self.assertIsNotNone(evaluation) + self.assertEqual(frozenset([2]), evaluation.contributing_field_ids) + self.assertEqual([Range(42, 42)], + evaluation.result.results().to_range_list()) + evaluator.close() + def test_and_with_disjoint_results(self): fields = _make_fields() result_a = GlobalIndexResult.from_range(Range(1, 3)) diff --git a/paimon-python/pypaimon/tests/global_index_scalar_search_mode_test.py b/paimon-python/pypaimon/tests/global_index_scalar_search_mode_test.py index 1985af3287a9..84d3e9468380 100644 --- a/paimon-python/pypaimon/tests/global_index_scalar_search_mode_test.py +++ b/paimon-python/pypaimon/tests/global_index_scalar_search_mode_test.py @@ -26,6 +26,7 @@ from pypaimon.globalindex.data_evolution_global_index_scanner import ( DataEvolutionGlobalIndexScanner, ) +from pypaimon.utils.range import Range def _ranges(result): @@ -47,6 +48,14 @@ def _coverage(options): table, snapshot, None, [SimpleNamespace(global_index_meta=meta)]) +def _scanner(coverage): + scanner = DataEvolutionGlobalIndexScanner.__new__( + DataEvolutionGlobalIndexScanner) + scanner._coverage = coverage + scanner._fields = [1] + return scanner + + class ScalarGlobalIndexSearchModeTest(unittest.TestCase): def test_default_values(self): @@ -93,15 +102,13 @@ def test_coverage_honours_search_mode_override(self): def test_scanner_applies_passed_scalar_mode(self): coverage = _coverage(CoreOptions(Options.from_none())) - scanner = SimpleNamespace(_coverage=coverage, _fields=[1]) - result = DataEvolutionGlobalIndexScanner.unindexed_rows( - scanner, None, search_mode=GlobalIndexSearchMode.FULL) - self.assertEqual([(100, 199)], _ranges(result)) + result = _scanner(coverage).unindexed_ranges( + None, search_mode=GlobalIndexSearchMode.FULL) + self.assertEqual([Range(100, 199)], result) def test_scanner_default_is_scalar_mode(self): coverage = _coverage(CoreOptions(Options.from_none())) - scanner = SimpleNamespace(_coverage=coverage, _fields=[1]) - result = DataEvolutionGlobalIndexScanner.unindexed_rows(scanner, None) + result = _scanner(coverage).unindexed_rows(None) self.assertEqual([], _ranges(result)) diff --git a/paimon-python/pypaimon/tests/global_index_test.py b/paimon-python/pypaimon/tests/global_index_test.py index 7959ffecfc7f..07adeeb25d99 100644 --- a/paimon-python/pypaimon/tests/global_index_test.py +++ b/paimon-python/pypaimon/tests/global_index_test.py @@ -26,6 +26,7 @@ from pypaimon.common.predicate import Predicate from pypaimon.common.predicate_builder import PredicateBuilder from pypaimon.globalindex.global_index_meta import GlobalIndexMeta +from pypaimon.globalindex.global_index_evaluator import GlobalIndexEvaluation from pypaimon.globalindex.global_index_result import GlobalIndexResult from pypaimon.index.index_file_meta import IndexFileMeta from pypaimon.index.index_file_handler import IndexFileHandler @@ -35,6 +36,7 @@ BatchModeMixin, DataEvolutionTestBase, ) +from pypaimon.utils.roaring_bitmap import RoaringBitmap64 from pypaimon.utils.range import Range @@ -217,8 +219,11 @@ def test_detail_mode_uses_table_data_ranges(self): class GlobalIndexScalarFallbackTest(unittest.TestCase): - def test_eval_global_index_merges_unindexed_rows_when_index_scan_succeeds(self): - from pypaimon.read.scanner.file_scanner import FileScanner + def test_eval_global_index_keeps_unindexed_ranges_out_of_bitmap(self): + from pypaimon.read.scanner.file_scanner import ( + FileScanner, + _GlobalIndexPlanningResult, + ) class _Options: def global_index_enabled(self): @@ -237,24 +242,63 @@ class _Table: scanner.table = _Table() index_result = GlobalIndexResult.from_range(Range(1, 1)) - unindexed = GlobalIndexResult.from_range(Range(5, 6)) + unindexed = [Range(5, 6)] fake_scanner = unittest.mock.MagicMock() - fake_scanner.scan.return_value = index_result - fake_scanner.unindexed_rows.return_value = unindexed + fake_scanner.scan_with_coverage.return_value = GlobalIndexEvaluation( + index_result, frozenset([0])) + fake_scanner.unindexed_ranges.return_value = unindexed fake_scanner.__enter__.return_value = fake_scanner fake_scanner.__exit__.return_value = None with unittest.mock.patch( "pypaimon.globalindex.data_evolution_global_index_scanner.DataEvolutionGlobalIndexScanner.create", - return_value=fake_scanner): + return_value=fake_scanner), unittest.mock.patch.object( + GlobalIndexResult, + "from_ranges", + side_effect=AssertionError("fallback ranges entered bitmap")): result = scanner._eval_global_index(snapshot=object()) + self.assertIsInstance(result, _GlobalIndexPlanningResult) + self.assertIs(index_result, result.indexed_result) + self.assertEqual(unindexed, result.unindexed_ranges) + fake_scanner.unindexed_ranges.assert_called_once_with( + predicate, + search_mode=GlobalIndexSearchMode.FULL, + contributing_field_ids=frozenset([0]), + ) + + def test_split_planning_merges_indexed_and_unindexed_ranges(self): + from pypaimon.read.scanner.file_scanner import ( + FileScanner, + _GlobalIndexPlanningResult, + ) + + scanner = FileScanner.__new__(FileScanner) + scanner.manifest_scanner = unittest.mock.MagicMock( + return_value=([], unittest.mock.Mock(id=3))) + scanner._global_index_result = None + scanner._eval_global_index = unittest.mock.MagicMock( + return_value=_GlobalIndexPlanningResult( + GlobalIndexResult.from_range(Range(1, 1)), + [Range(10, 10 ** 12)], + )) + scanner.predicate = Predicate( + method="equal", index=0, field="id", literals=[1]) + scanner.read_manifest_entries = unittest.mock.MagicMock(return_value=[]) + scanner.table = unittest.mock.Mock() + scanner.target_split_size = 1 + scanner.open_file_cost = 1 + scanner._deletion_files_map = unittest.mock.MagicMock(return_value={}) + + with unittest.mock.patch( + "pypaimon.read.scanner.file_scanner.DataEvolutionSplitGenerator" + ) as split_generator: + scanner._create_data_evolution_split_generator() + self.assertEqual( - [Range(1, 1), Range(5, 6)], - result.results().to_range_list(), + [Range(1, 1), Range(10, 10 ** 12)], + split_generator.call_args[0][4], ) - fake_scanner.unindexed_rows.assert_called_once_with( - predicate, search_mode=GlobalIndexSearchMode.FULL) def test_eval_global_index_keeps_none_as_full_scan(self): from pypaimon.read.scanner.file_scanner import FileScanner @@ -273,7 +317,7 @@ class _Table: scanner.table = _Table() fake_scanner = unittest.mock.MagicMock() - fake_scanner.scan.return_value = None + fake_scanner.scan_with_coverage.return_value = None fake_scanner.__enter__.return_value = fake_scanner fake_scanner.__exit__.return_value = None @@ -295,6 +339,27 @@ class PlanSnapshotFetchRegressionTest( 'bucket': '-1', } + @pytest.mark.python_plan + def test_plan_accepts_row_ranges_without_bitmap(self): + table = self._create_table() + self._write_arrow(table, pa.table( + {'id': [1, 2, 3], 'name': ['a', 'b', 'c'], + 'age': [10, 20, 30], 'city': ['x', 'y', 'z']}, + schema=self.pa_schema)) + + read_builder = table.new_read_builder() + ranges = [Range(0, 10 ** 12)] + with patch.object( + RoaringBitmap64, + 'to_range_list', + side_effect=AssertionError('row ranges entered a bitmap')): + plan = read_builder.new_scan().with_row_ranges(ranges).plan() + + result = read_builder.new_read().to_arrow(plan.splits()) + self.assertEqual([1, 2, 3], sorted(result.column('id').to_pylist())) + self.assertEqual( + [], read_builder.new_scan().with_row_ranges([]).plan().splits()) + @pytest.mark.python_plan def test_plan_fetches_latest_snapshot_only_once(self): table = self._create_table() diff --git a/paimon-python/pypaimon/tests/native_plan_test.py b/paimon-python/pypaimon/tests/native_plan_test.py index b590bb4f90b3..3abecc22f829 100644 --- a/paimon-python/pypaimon/tests/native_plan_test.py +++ b/paimon-python/pypaimon/tests/native_plan_test.py @@ -64,6 +64,7 @@ def _scan(native_enabled, file_scanner): file_scanner.start_pos_of_this_subtask = None # no slice file_scanner.chunk_shuffle = None # no chunk-shuffle file_scanner._global_index_result = None # no global-index result + file_scanner._row_ranges = None # no explicit row ranges file_scanner.deletion_vectors_enabled = False # no deletion vectors file_scanner.data_evolution = False # no data evolution file_scanner.only_read_real_buckets = False # not postpone bucket @@ -174,8 +175,8 @@ def test_plan_forwards_filter_limit_partition_and_time_travel(self): ) def test_plan_falls_back_when_scan_is_not_plain(self): - # Native planning does not carry shard/slice, global-index, or - # incremental scans -> must fall back to the file scanner. + # Native planning does not carry shard/slice, global-index, row ranges, + # or incremental scans -> must fall back to the file scanner. def check(setup): fs = Mock(partition_key_predicate=None) sentinel = object() @@ -191,6 +192,7 @@ def check(setup): check(lambda s, fs: setattr(fs, 'start_pos_of_this_subtask', 0)) check(lambda s, fs: setattr(fs, 'chunk_shuffle', (1, 100))) check(lambda s, fs: setattr(fs, '_global_index_result', object())) + check(lambda s, fs: setattr(fs, '_row_ranges', [object()])) check(lambda s, fs: setattr(fs, 'deletion_vectors_enabled', True)) check(lambda s, fs: setattr(fs, 'only_read_real_buckets', True)) check(lambda s, fs: (setattr(s.table, 'is_primary_key_table', True), diff --git a/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py index 735a0edab039..7bf7c9e9270e 100644 --- a/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py +++ b/paimon-python/pypaimon/tests/scanner/chunk_shuffle_split_generator_test.py @@ -1123,6 +1123,14 @@ def test_workers_union_equals_full_table(self): self.assertEqual(actual.column('id').to_pylist(), list(range(200))) self.assertEqual(actual.column('payload').to_pylist(), self._payloads(range(200))) + def test_row_ranges_with_chunk_shuffle_rejected(self): + table, _ = self._create_de_table('cs_de_row_ranges') + scan = table.new_read_builder().new_scan() \ + .with_row_ranges([Range(0, 0)]) \ + .with_chunk_shuffle(seed=1, chunk_size=10) + with self.assertRaisesRegex(ValueError, "row ranges"): + scan.plan() + def test_deterministic_plan_across_calls(self): table, pa_schema = self._create_de_table('cs_de_determinism') for c in range(3): diff --git a/paimon-python/pypaimon/tests/vector_search_filter_test.py b/paimon-python/pypaimon/tests/vector_search_filter_test.py index 8cd2d7e6a27e..e3bbccc911b5 100644 --- a/paimon-python/pypaimon/tests/vector_search_filter_test.py +++ b/paimon-python/pypaimon/tests/vector_search_filter_test.py @@ -36,6 +36,7 @@ from pypaimon.common.predicate_builder import PredicateBuilder from pypaimon.globalindex.btree.btree_index_meta import BTreeIndexMeta from pypaimon.globalindex.global_index_meta import GlobalIndexIOMeta, GlobalIndexMeta +from pypaimon.globalindex.global_index_evaluator import GlobalIndexEvaluation from pypaimon.globalindex.global_index_reader import _completed_future from pypaimon.globalindex.global_index_result import GlobalIndexResult from pypaimon.globalindex.vector_search import VectorSearch @@ -169,7 +170,7 @@ def _bitmap(*row_ids): def _install_raw_vector_read_builder(table, vector_column_name, row_id_to_vector, calls=None): - """Install a fake raw read builder which honors GlobalIndexResult ranges.""" + """Install a fake raw read builder which honors row ranges.""" import pyarrow as pa calls = calls if calls is not None else {} @@ -185,10 +186,9 @@ class _Scan: def __init__(self): self._row_ids = [] - def with_global_index_result(self, result): - ranges = result.results().to_range_list() + def with_row_ranges(self, ranges): calls["raw_read_count"] = calls.get("raw_read_count", 0) + 1 - calls["global_index_ranges"] = ranges + calls["global_index_ranges"] = list(ranges) self._row_ids = [ row_id for row_id in sorted(row_id_to_vector) @@ -234,7 +234,7 @@ def new_read(self): def _install_raw_full_text_read_builder(table, text_column_name, row_id_to_text, calls=None): - """Install a fake raw read builder which honors GlobalIndexResult ranges.""" + """Install a fake raw read builder which honors row ranges.""" import pyarrow as pa calls = calls if calls is not None else {} @@ -250,9 +250,8 @@ class _Scan: def __init__(self): self._row_ids = [] - def with_global_index_result(self, result): - ranges = result.results().to_range_list() - calls["global_index_ranges"] = ranges + def with_row_ranges(self, ranges): + calls["global_index_ranges"] = list(ranges) self._row_ids = [ row_id for row_id in sorted(row_id_to_text) @@ -993,6 +992,50 @@ def test_scan_attaches_overlapping_scalar_index_files(self): (splits_sorted[1].row_range_start, splits_sorted[1].row_range_end)) + def test_unsupported_scalar_coverage_still_plans_raw_split(self): + from pypaimon.table.source.vector_search_split import ( + IndexVectorSearchSplit, + RawVectorSearchSplit, + ) + + entries = [ + _entry(None, field_id=1, index_type="lumina-vector-ann", + file_name="vec.index", row_range_start=0, row_range_end=9), + _entry(None, field_id=0, index_type="full-text", + file_name="id-ft.index", row_range_start=0, row_range_end=9), + ] + table = _StubTable( + fields=[self.id_field, self.embedding_field], entries=entries) + table.options = CoreOptions(Options({ + "scalar-index.search-mode": "full", + "vector-index.search-mode": "full", + })) + self._scan_patch.stop() + self._travel_patch.stop() + _patch_snapshot( + self, entries, types.SimpleNamespace(id=1, next_row_id=10)) + + predicate = Predicate( + method="equal", index=0, field="id", literals=[5]) + splits = ( + VectorSearchBuilderImpl(table) + .with_vector_column("embedding") + .with_query_vector([1.0, 0.0, 0.0, 0.0]) + .with_limit(3) + .with_filter(predicate) + .new_vector_search_scan() + .scan() + .splits() + ) + + index = [s for s in splits if isinstance(s, IndexVectorSearchSplit)] + raw = [s for s in splits if isinstance(s, RawVectorSearchSplit)] + self.assertEqual(1, len(index)) + self.assertEqual([], index[0].scalar_index_files) + self.assertEqual(1, len(raw)) + self.assertEqual([Range(0, 9)], raw[0].row_ranges) + self.assertEqual([], raw[0].scalar_index_files) + def test_read_threads_prefilter_bitmap_as_include_row_ids(self): """preFilter bitmap from scanner.scan(filter) must reach each split's VectorSearch, offset-rebased to local coords by OffsetGlobalIndexReader. @@ -1588,8 +1631,9 @@ def test_raw_vector_pre_filter_uses_scalar_search_mode(self): "scalar-index.search-mode": "detail", })) scanner = mock.MagicMock() - scanner.scan.return_value = GlobalIndexResult.create_empty() - scanner.unindexed_rows.return_value = GlobalIndexResult.create_empty() + scanner.scan_with_coverage.return_value = GlobalIndexEvaluation( + GlobalIndexResult.create_empty(), frozenset([0])) + scanner.unindexed_ranges.return_value = [] reader = DataEvolutionVectorRead( table, limit=3, @@ -1602,11 +1646,79 @@ def test_raw_vector_pre_filter_uses_scalar_search_mode(self): "pypaimon.globalindex.data_evolution_global_index_scanner." "DataEvolutionGlobalIndexScanner.create", return_value=scanner): - reader._raw_pre_filter([ + result = reader._raw_pre_filter([ RawVectorSearchSplit([Range(0, 9)], [scalar_file])]) - scanner.unindexed_rows.assert_called_once_with( - predicate, search_mode=GlobalIndexSearchMode.DETAIL) + self.assertEqual([], result) + scanner.unindexed_ranges.assert_called_once_with( + predicate, + search_mode=GlobalIndexSearchMode.DETAIL, + contributing_field_ids=frozenset([0]), + ) + + def test_raw_vector_pre_filter_keeps_full_fallback_as_ranges(self): + from pypaimon.table.source.vector_search_read import DataEvolutionVectorRead + from pypaimon.table.source.vector_search_split import RawVectorSearchSplit + + predicate = Predicate(method="equal", index=0, field="id", literals=[5]) + scalar_file = self.entries[2].index_file + table = _StubTable(fields=[self.id_field, self.embedding_field], entries=[]) + table.options = CoreOptions(Options({ + "scalar-index.search-mode": "full", + })) + scanner = mock.MagicMock() + scanner.scan_with_coverage.return_value = GlobalIndexEvaluation( + GlobalIndexResult.from_range(Range(5, 5)), frozenset([0])) + scanner.unindexed_ranges.return_value = [Range(10, 10 ** 12)] + scanner.unindexed_rows.side_effect = AssertionError( + "FULL fallback must not enter a bitmap") + reader = DataEvolutionVectorRead( + table, + limit=3, + vector_column=self.embedding_field, + query_vector=[1.0, 0.0, 0.0, 0.0], + filter_=predicate, + ) + + with mock.patch( + "pypaimon.globalindex.data_evolution_global_index_scanner." + "DataEvolutionGlobalIndexScanner.create", + return_value=scanner): + result = reader._raw_pre_filter([ + RawVectorSearchSplit([Range(0, 20)], [scalar_file])]) + + self.assertEqual([Range(5, 5), Range(10, 20)], result) + scanner.unindexed_ranges.assert_called_once_with( + predicate, + search_mode=GlobalIndexSearchMode.FULL, + contributing_field_ids=frozenset([0]), + ) + + def test_raw_vector_read_passes_ranges_without_bitmap(self): + from pypaimon.table.source.vector_search_read import DataEvolutionVectorRead + + table = _StubTable(fields=[self.embedding_field], entries=[]) + calls = _install_raw_vector_read_builder( + table, + "embedding", + {10: [1.0, 0.0, 0.0, 0.0]}, + ) + ranges = [Range(10, 10 ** 12)] + reader = DataEvolutionVectorRead( + table, + limit=1, + vector_column=self.embedding_field, + query_vector=[1.0, 0.0, 0.0, 0.0], + ) + + with mock.patch.object( + GlobalIndexResult, + "from_ranges", + side_effect=AssertionError("row ranges entered a bitmap")): + result = reader._read_raw_arrow(ranges, include_filter=True) + + self.assertEqual(ranges, calls["global_index_ranges"]) + self.assertEqual(1, result.num_rows) def test_scan_threads_builder_options_to_raw_split_index_type(self): from pypaimon.table.source.vector_search_split import RawVectorSearchSplit @@ -1729,23 +1841,29 @@ def close(self_inner): has_nulls=False) with mock.patch( - "pypaimon.globalindex.btree.lazy_filtered_btree_reader.BTreeIndexReader", - _StubBTreeReader): + "pypaimon.globalindex.data_evolution_global_index_scanner." + "_exclude_ranges", + side_effect=AssertionError("single group must not compute padding")): with mock.patch( - "pypaimon.globalindex.sorted_file_global_index_reader.SortedIndexFileMeta.deserialize", - return_value=wide_meta): - scanner = DataEvolutionGlobalIndexScanner( - fields=table.fields, - file_io=table.file_io, - index_path="/unused", - index_files=[shard_a, shard_b], - ) - try: - result = scanner.scan( - Predicate(method="equal", index=0, field="id", - literals=[7])) - finally: - scanner.close() + "pypaimon.globalindex.btree.lazy_filtered_btree_reader." + "BTreeIndexReader", + _StubBTreeReader): + with mock.patch( + "pypaimon.globalindex.sorted_file_global_index_reader." + "SortedIndexFileMeta.deserialize", + return_value=wide_meta): + scanner = DataEvolutionGlobalIndexScanner( + fields=table.fields, + file_io=table.file_io, + index_path="/unused", + index_files=[shard_a, shard_b], + ) + try: + result = scanner.scan( + Predicate(method="equal", index=0, field="id", + literals=[7])) + finally: + scanner.close() self.assertIsNotNone(result) hits = sorted(list(result.results())) @@ -1753,6 +1871,134 @@ def close(self_inner): # Must not be empty despite shard_a being empty (no short-circuit). self.assertEqual([7], hits) + def test_primary_and_extra_field_indexes_share_coverage(self): + from pypaimon.globalindex.global_index_reader import GlobalIndexReader + from pypaimon.globalindex.data_evolution_global_index_scanner import ( + DataEvolutionGlobalIndexScanner, + ) + + fields = [_field(0, "a"), _field(1, "b"), _field(2, "c")] + primary = _entry(None, field_id=2, index_type="btree", + file_name="c-primary.index", + row_range_start=0, row_range_end=4).index_file + extra = _entry(None, field_id=0, index_type="btree", + file_name="a-c.index", + row_range_start=5, row_range_end=9).index_file + extra.global_index_meta.extra_field_ids = [2] + table = _StubTable(fields=fields, entries=[]) + table.options = CoreOptions(Options({ + "scalar-index.search-mode": "full", + })) + + class _StubReader(GlobalIndexReader): + def __init__(self_inner, file_name): + self_inner._file_name = file_name + + def visit_equal(self_inner, field_ref, literal): + bitmap = RoaringBitmap64() + bitmap.add(1 if self_inner._file_name == "c-primary.index" else 2) + return _completed_future(GlobalIndexResult.create(bitmap)) + + 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.data_evolution_global_index_scanner." + "_create_inner_readers", + side_effect=_stub_create_inner_readers): + scanner = DataEvolutionGlobalIndexScanner( + fields=fields, + file_io=object(), + index_path="/unused", + index_files=[primary, extra], + options=table.options, + table=table, + snapshot=types.SimpleNamespace(next_row_id=10), + ) + try: + evaluation = scanner.scan_with_coverage( + Predicate(method="equal", index=2, field="c", + literals=[42])) + fallback = scanner.unindexed_rows( + None, + search_mode=GlobalIndexSearchMode.FULL, + contributing_field_ids=( + evaluation.contributing_field_ids), + ) + finally: + scanner.close() + + result = evaluation.result.or_(fallback) + self.assertTrue(fallback.results().is_empty()) + self.assertEqual([1, 7], sorted(result.results())) + + def test_unsupported_extra_field_index_does_not_poison_primary(self): + from pypaimon.globalindex.global_index_reader import GlobalIndexReader + from pypaimon.globalindex.data_evolution_global_index_scanner import ( + DataEvolutionGlobalIndexScanner, + ) + + fields = [_field(0, "a"), _field(1, "b"), _field(2, "c")] + primary = _entry(None, field_id=2, index_type="btree", + file_name="c-primary.index", + row_range_start=0, row_range_end=4).index_file + unsupported = _entry(None, field_id=0, index_type="es-index", + file_name="a-c.index", + row_range_start=5, row_range_end=9).index_file + unsupported.global_index_meta.extra_field_ids = [2] + table = _StubTable(fields=fields, entries=[]) + table.options = CoreOptions(Options({ + "scalar-index.search-mode": "full", + })) + + class _StubReader(GlobalIndexReader): + def visit_equal(self_inner, field_ref, literal): + bitmap = RoaringBitmap64() + bitmap.add(1) + return _completed_future(GlobalIndexResult.create(bitmap)) + + def close(self_inner): + pass + + observed_types = [] + + def _stub_create_inner_readers( + index_type, file_io, index_path, field, io_metas, + executor=None, options=None): + observed_types.append(index_type) + return [_StubReader()] + + with mock.patch( + "pypaimon.globalindex.data_evolution_global_index_scanner." + "_create_inner_readers", + side_effect=_stub_create_inner_readers): + scanner = DataEvolutionGlobalIndexScanner.create( + table, + index_files=[primary, unsupported], + snapshot=types.SimpleNamespace(next_row_id=10), + ) + try: + evaluation = scanner.scan_with_coverage( + Predicate(method="equal", index=2, field="c", + literals=[42])) + fallback = scanner.unindexed_ranges( + None, + search_mode=GlobalIndexSearchMode.FULL, + contributing_field_ids=( + evaluation.contributing_field_ids), + ) + finally: + scanner.close() + + self.assertEqual(["btree"], observed_types) + self.assertEqual([1], sorted(evaluation.result.results())) + self.assertEqual([Range(5, 9)], fallback) + def test_extra_field_groups_are_padded_before_and(self): from pypaimon.globalindex.global_index_reader import GlobalIndexReader from pypaimon.globalindex.data_evolution_global_index_scanner import ( @@ -1933,70 +2179,66 @@ def _stub_create_inner_readers( self.assertIsNone(result) - def test_native_fulltext_index_is_dispatched_by_scanner(self): - """Non-btree scalar global indexes (full-text, etc.) must be - instantiated by DataEvolutionGlobalIndexScanner — previously only 'btree' was - handled and everything else was silently dropped, making text-column - pre-filter a no-op.""" - from pypaimon.globalindex.global_index_result import GlobalIndexResult + def test_full_text_index_is_not_scalar_coverage(self): + from pypaimon.globalindex.global_index_reader import GlobalIndexReader from pypaimon.globalindex.data_evolution_global_index_scanner import ( DataEvolutionGlobalIndexScanner, ) - name_field = _field(0, "name", "STRING") - emb_field = _field(1, "embedding", "FLOAT") - full_text_shard = _entry( - None, field_id=0, index_type="full-text", - file_name="name-ft.index", - row_range_start=0, row_range_end=9, - external_path="oss://bucket/name-ft.index").index_file - table = _StubTable(fields=[name_field, emb_field], entries=[]) - - captured_ctor_args = [] - visit_calls = [] - - from pypaimon.globalindex.global_index_reader import _completed_future as _cf - - class _StubFullTextReader: - def __init__(self_inner, file_io, index_path, io_metas): - captured_ctor_args.append( - (file_io, index_path, list(io_metas))) + field = _field(0, "name", "STRING") + btree = _entry(None, field_id=0, index_type="btree", + file_name="name-btree.index", + row_range_start=0, row_range_end=4).index_file + full_text = _entry(None, field_id=0, index_type="full-text", + file_name="name-ft.index", + row_range_start=5, row_range_end=9).index_file + table = _StubTable(fields=[field], entries=[]) + table.options = CoreOptions(Options({ + "scalar-index.search-mode": "full", + })) + class _StubReader(GlobalIndexReader): def visit_equal(self_inner, field_ref, literal): - visit_calls.append(("equal", literal)) bm = RoaringBitmap64() - bm.add(4) - return _cf(GlobalIndexResult.create(bm)) + bm.add(1) + return _completed_future(GlobalIndexResult.create(bm)) def close(self_inner): pass + observed_types = [] + + def _stub_create_inner_readers( + index_type, file_io, index_path, field, io_metas, + executor=None, options=None): + observed_types.append(index_type) + return [_StubReader()] + with mock.patch( - "pypaimon.globalindex.full_text.NativeFullTextGlobalIndexReader", - _StubFullTextReader): - scanner = DataEvolutionGlobalIndexScanner( - fields=table.fields, - file_io=table.file_io, - index_path="/unused", - index_files=[full_text_shard], + "pypaimon.globalindex.data_evolution_global_index_scanner." + "_create_inner_readers", + side_effect=_stub_create_inner_readers): + scanner = DataEvolutionGlobalIndexScanner.create( + table, + index_files=[btree, full_text], + snapshot=types.SimpleNamespace(next_row_id=10), ) try: - result = scanner.scan( + evaluation = scanner.scan_with_coverage( Predicate(method="equal", index=0, field="name", literals=["x"])) + fallback = scanner.unindexed_ranges( + None, + search_mode=GlobalIndexSearchMode.FULL, + contributing_field_ids=( + evaluation.contributing_field_ids), + ) finally: scanner.close() - # Native full-text reader was instantiated (it would NOT be before this fix). - self.assertEqual(1, len(captured_ctor_args)) - _, _, io_metas = captured_ctor_args[0] - self.assertEqual("oss://bucket/name-ft.index", - io_metas[0].external_path) - # visit_equal was dispatched all the way through evaluator → union → - # offset → stub native full-text reader. - self.assertEqual([("equal", "x")], visit_calls) - # Row id 4 is inside [0,9] so offset rebase is a no-op. - self.assertEqual([4], sorted(list(result.results()))) + self.assertEqual(["btree"], observed_types) + self.assertEqual([1], sorted(evaluation.result.results())) + self.assertEqual([Range(5, 9)], fallback) def test_like_predicate_is_dispatched_to_reader(self): """Evaluator must dispatch ``like`` to reader.visit_like — otherwise @@ -2725,8 +2967,8 @@ def splits(self_inner): return ["split"] class _Scan: - def with_global_index_result(self_inner, result): - calls["global_index_ranges"] = result.results().to_range_list() + def with_row_ranges(self_inner, ranges): + calls["global_index_ranges"] = list(ranges) return self_inner def plan(self_inner):