Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 47 additions & 27 deletions src/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ use lance_core::datatypes::BlobHandling;
use lance_index::scalar::FullTextSearchQuery;
use lance_index::vector::ApproxMode;
use lance_io::stream::RecordBatchStream;
use lance_table::format::IndexMetadata;
use lance_table::format::{Fragment, IndexMetadata};
use uuid::Uuid;

use crate::async_dispatcher::{self, LanceCallback};
Expand Down Expand Up @@ -421,7 +421,10 @@ impl LanceScanner {
if self.include_deleted_rows {
scanner.include_deleted_rows();
}
self.apply_fragment_filter(&mut scanner)?;
let apply_fragment_filter_after_nearest = self.nearest.is_some() && !self.prefilter;
if !apply_fragment_filter_after_nearest {
self.apply_fragment_filter(&mut scanner)?;
}
if self.index_segments.is_some() && self.nearest.is_none() {
return Err(lance_core::Error::invalid_input_source(
"index_segments requires nearest() to be configured".into(),
Expand All @@ -437,8 +440,8 @@ impl LanceScanner {
"fragment_ids cannot be combined with an FTS query context; split the query by FTS index segment UUID instead".into(),
));
}
// nearest() checks the current prefilter setting before accepting a
// fragment-scoped search. Enable it before installing the query.
// nearest() checks this setting at configuration time. In postfilter mode, defer
// applying fragment IDs until after nearest() is installed.
if self.prefilter {
scanner.prefilter(true);
}
Expand Down Expand Up @@ -495,6 +498,9 @@ impl LanceScanner {
scanner.with_index_segments(segments.clone())?;
}
}
if apply_fragment_filter_after_nearest {
self.apply_fragment_filter(&mut scanner)?;
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] with_fragments does not become a post-filter based on call order

Scanner::with_fragments only stores self.fragments; the builder does not preserve whether it was called before or after nearest(). During plan creation, that fragment scope still restricts relevant index segments and flat-search inputs before Top-K. Moving this call below nearest() therefore bypasses ensure_not_fragment_scan(), but does not implement post-ranking filtering.

For example, if the global Top-K rows are all in fragment 0 while fragment 1 is selected, a real post-filter should return no rows, whereas this plan can return the Top-K rows from fragment 1.

Please implement the fragment predicate as an explicit post-ranking execution node, or clarify that fragments define the search domain and keep the prefilter semantics. Add a two-fragment indexed regression test that distinguishes these outcomes.

if let Some(fts) = &self.fts_query {
scanner.full_text_search(fts.clone())?;
}
Expand All @@ -505,6 +511,10 @@ impl LanceScanner {
Some(PreparedFtsExecution {
context: Arc::clone(context),
segments,
scope_prefilter_to_fts_segments: self.prefilter
&& (self.filter.is_some()
|| self.substrait_filter.is_some()
|| !self.additional_sql_filters.is_empty()),
batch_size: self.batch_size,
scan_statistics_callback: self.scan_statistics_callback.clone(),
})
Expand Down Expand Up @@ -558,6 +568,7 @@ impl LanceScanner {
struct PreparedFtsExecution {
context: Arc<FtsQueryContextInner>,
segments: Vec<IndexMetadata>,
scope_prefilter_to_fts_segments: bool,
batch_size: Option<usize>,
scan_statistics_callback: Option<ExecutionStatsCallback>,
}
Expand Down Expand Up @@ -596,11 +607,21 @@ impl PreparedScanner {
let Some(distributed_fts) = self.distributed_fts else {
return self.scanner.try_into_stream().await;
};
let plan = self.scanner.create_plan().await?;
let selected_segments_have_current_fragments = segments_have_current_fragments(
let selected_fragments = selected_current_fts_fragments(
&distributed_fts.context.dataset,
&distributed_fts.segments,
)?;
let selected_segments_have_current_fragments = !selected_fragments.is_empty();
let mut scanner = self.scanner;
if distributed_fts.scope_prefilter_to_fts_segments
&& selected_segments_have_current_fragments
{
// The scanner is already split by the selected FTS segment(s). Applying the same
// fragment scope before plan creation lets Lance restrict scalar-index segment loads
// for the TVF prefilter without changing unfiltered FTS scans.
Comment on lines 514 to +621

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Add regression coverage for both new query-planning branches

This PR introduces two behavior changes, but the only test diff updates a helper call in the existing unfiltered prepared-FTS plan test.

Please add:

  1. An indexed, two-fragment vector test for fragment_ids + prefilter=false whose expected result distinguishes global Top-K post-filtering from fragment-scoped Top-K.
  2. A prepared-FTS test with at least two segments, one selected segment, a SQL filter, and prefilter=true. Assert returned IDs and global scores, and use scan statistics or plan assertions to verify that only the selected segment’s fragments are loaded for the prefilter.

This is required because both branches affect result-selection boundaries, not just execution cost.

scanner.with_fragments(selected_fragments);
}
let plan = scanner.create_plan().await?;
let (plan, rewritten) = rewrite_prepared_fts_plan(
plan,
&distributed_fts.segments,
Expand Down Expand Up @@ -661,37 +682,34 @@ fn select_fts_segments(
Ok(selected)
}

fn segments_have_current_fragments(
fn selected_current_fts_fragments(
dataset: &lance::Dataset,
segments: &[IndexMetadata],
) -> Result<bool> {
let current_fragment_ids = dataset
.get_fragments()
.into_iter()
.map(|fragment| {
u32::try_from(fragment.id()).map_err(|_| {
lance_core::Error::internal(format!(
"current fragment id {} exceeds the validated u32 FTS coverage range",
fragment.id()
))
})
})
.collect::<Result<std::collections::HashSet<_>>>()?;
) -> Result<Vec<Fragment>> {
let mut selected_fragment_ids = std::collections::HashSet::new();
for segment in segments {
let fragment_bitmap = segment.fragment_bitmap.as_ref().ok_or_else(|| {
lance_core::Error::internal(format!(
"prepared FTS segment {} lost its validated fragment coverage",
segment.uuid
))
})?;
if fragment_bitmap
.iter()
.any(|fragment_id| current_fragment_ids.contains(&fragment_id))
{
return Ok(true);
selected_fragment_ids.extend(fragment_bitmap.iter());
}

let mut selected_fragments = Vec::new();
for fragment in dataset.get_fragments() {
let fragment_id = u32::try_from(fragment.id()).map_err(|_| {
lance_core::Error::internal(format!(
"current fragment id {} exceeds the validated u32 FTS coverage range",
fragment.id()
))
})?;
if selected_fragment_ids.contains(&fragment_id) {
selected_fragments.push(fragment.metadata().clone());
}
}
Ok(false)
Ok(selected_fragments)
}

#[derive(Default)]
Expand Down Expand Up @@ -3357,7 +3375,9 @@ mod tests {
);

let has_current_fragments =
segments_have_current_fragments(&distributed.context.dataset, &segments).unwrap();
!selected_current_fts_fragments(&distributed.context.dataset, &segments)
.unwrap()
.is_empty();
let (rewritten, counts) = rewrite_prepared_fts_plan(
plan,
&segments,
Expand Down
Loading