Skip to content
Merged
Show file tree
Hide file tree
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
41 changes: 41 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,47 @@ Based on the [liblance RFC](https://github.com/lance-format/lance/discussions/60
| [x] | Dataset metadata | `lance_dataset_version()`, `lance_dataset_count_rows()`, `lance_dataset_latest_version()` |
| [x] | Filter pushdown | `lance_scanner_set_substrait_filter()` accepts a serialized Substrait `ExtendedExpression`; `lance_scanner_additional_sql_filter()` adds SQL predicates with AND before scanning starts |

## Segment-scoped array label filters

Ordinary scans can use a `LabelList` index for array membership through the
existing SQL and Substrait filter interfaces. For example, on a `List<Utf8>`
column named `labels`:

```sql
array_contains(labels, 'red')
array_contains(labels, 'red') AND array_contains(labels, 'blue')
array_contains(labels, 'red') OR array_contains(labels, 'blue')
array_has_all(labels, ['red', 'blue'])
array_has_any(labels, ['red', 'blue'])
```

Configure `lance_scanner_set_fragment_ids` and
`lance_scanner_set_scalar_index_segment` with the selected LabelList segment.
Pass a SQL filter to `lance_scanner_new`, or attach a serialized Substrait
`ExtendedExpression` with `lance_scanner_set_substrait_filter`. The Substrait
schema must describe the list field and its element type; replacing that field
with an unsupported-type placeholder cannot express a label predicate. Use
Lance/DataFusion's `array_has` (the canonical name of `array_contains`),
`array_has_all`, or `array_has_any` functions with correctly typed arguments.
Substrait takes precedence over the primary SQL filter; use
`lance_scanner_additional_sql_filter` when an additional condition must be ANDed
with it.

AND/OR membership expressions can reuse the same selected LabelList segment.
Other columns remain residual filters, evaluated before LIMIT/OFFSET. This API
selects one physical segment; it does not intersect indices on different columns.
Incomplete segment coverage falls back to scanning the entire explicit fragment
scope. Check `scalar_segments_searched` and `scalar_segment_fallbacks` in the
statistics callback to distinguish index acceleration from filter execution.

Bindings preserve Lance's function semantics, not the semantics of similarly
named functions in another SQL engine. In particular, a NULL search value in
`array_contains` does not match NULL array elements. `array_has_all` tests set
containment, not an ordered contiguous subsequence. The pinned Lance version
also treats an empty all-label query as true even for NULL lists. Integrators
should initially push only non-NULL constant labels with matching element types
and retain conditions whose semantics have not been verified in the calling engine.

## Distance-bounded vector search

After configuring a single-vector nearest-neighbor query, use
Expand Down
14 changes: 11 additions & 3 deletions include/lance/lance.h
Original file line number Diff line number Diff line change
Expand Up @@ -2256,13 +2256,21 @@ int32_t lance_scanner_set_index_segments(
* fragment domains and separately include any unindexed data they wish to read.
* The segment metadata must identify one key field present in the schema.
*
* BTree/Bitmap/LabelList searches use a necessary AND-conjunct of the
* full scanner filter on the selected logical index and require an Exact result.
* BTree/Bitmap/LabelList searches evaluate a candidate expression on the selected
* logical index, including AND, OR, IN and NULL-aware NOT, and require an Exact
* result. LabelList supports Lance array_has/array_contains, array_has_all and
* array_has_any through SQL or Substrait filters with matching list element types.
* These functions retain Lance semantics, including for NULL search values.
* An AND may retain only its supported necessary conditions. OR needs
* candidates for both branches; NOT requires its complete indexed subtree.
* Expressions exceeding 128 nodes or depth 32 use the scoped fallback.
* Filters containing IS [NOT] TRUE/FALSE also use that fallback until the Lance
* planner dependency preserves their NULL semantics under negation.
* use_scalar_index=false skips segment search and uses the scoped fallback;
* snapshot UUID and fragment validation still applies.
* AtMost/AtLeast results fall back to a full filtered scan of fragment_ids.
* All predicates are reapplied during candidate reads; other scalar indices
* are disabled. Legacy storage, OR/NOT-only filters,
* are disabled. Legacy storage, expressions without safe scoped candidates,
* overlays, fragment reuse, unsupported index types / result domains
* and missing coverage use the same domain without an index. No filter also
* falls back. LIMIT/OFFSET apply after the complete scanner filter, never to the
Expand Down
210 changes: 179 additions & 31 deletions src/scalar_segment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,19 +8,19 @@ use std::collections::HashSet;
use std::sync::Arc;
use std::time::Instant;

use datafusion::common::tree_node::TreeNode;
use datafusion::logical_expr::Expr;
use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
use lance::Dataset;
use lance::dataset::scanner::{
ExecutionStatsCallback, ExecutionSummaryCounts, RowAddrMask, Scanner,
};
use lance::dataset::scanner::{ExecutionStatsCallback, ExecutionSummaryCounts, Scanner};
use lance::index::{DatasetIndexExt, DatasetIndexInternalExt};
use lance::io::exec::utils::IndexMetrics;
use lance_core::{Error, Result};
use lance_datafusion::planner::Planner;
use lance_datafusion::utils::MetricsExt;
use lance_index::IndexType;
use lance_index::scalar::SearchResult;
use lance_index::scalar::expression::{PlannerIndexExt, ScalarIndexExpr, ScalarIndexSearch};
use lance_index::scalar::expression::{PlannerIndexExt, ScalarIndexExpr, ScalarIndexLoader};
use lance_index::scalar::{MetricsCollector, ScalarIndex};
use uuid::Uuid;

pub(crate) struct PreparedScalarSegment {
Expand All @@ -35,15 +35,83 @@ fn invalid(message: impl Into<String>) -> Error {
Error::invalid_input_source(message.into().into())
}

// Only descend through AND: a leaf below OR or NOT need not contain all matches
// of the full expression. The original expression is always reapplied by reader.
fn driver<'a>(expr: &'a ScalarIndexExpr, index_name: &str) -> Option<&'a ScalarIndexSearch> {
// Bound both recursive planning and the concurrent searches in Lance's evaluator.
const MAX_EXPRESSION_NODES: usize = 128;
const MAX_EXPRESSION_DEPTH: usize = 32;

fn scoped_expression(
expr: &ScalarIndexExpr,
index_name: &str,
column: &str,
require_complete: bool,
depth: usize,
remaining: &mut usize,
) -> std::result::Result<Option<ScalarIndexExpr>, &'static str> {
if depth > MAX_EXPRESSION_DEPTH || *remaining == 0 {
return Err("expression_budget");
}
*remaining -= 1;
let recurse = |child, complete, remaining: &mut usize| {
scoped_expression(child, index_name, column, complete, depth + 1, remaining)
};
match expr {
ScalarIndexExpr::Query(search) if search.index_name == index_name => Some(search),
ScalarIndexExpr::Query(search) => {
if search.index_name != index_name {
Ok(None)
} else if search.column != column {
Err("field_path")
} else {
Ok(Some(expr.clone()))
}
}
ScalarIndexExpr::And(lhs, rhs) => {
driver(lhs, index_name).or_else(|| driver(rhs, index_name))
let lhs = recurse(lhs, require_complete, remaining)?;
let rhs = recurse(rhs, require_complete, remaining)?;
Ok(match (lhs, rhs) {
(Some(lhs), Some(rhs)) => Some(ScalarIndexExpr::And(Box::new(lhs), Box::new(rhs))),
(lhs, rhs) if !require_complete => lhs.or(rhs),
_ => None,
})
}
ScalarIndexExpr::Or(lhs, rhs) => {
let lhs = recurse(lhs, require_complete, remaining)?;
let rhs = recurse(rhs, require_complete, remaining)?;
// Both branches must contribute a superset of their matches. Dropping
// an unavailable OR branch would silently exclude valid rows.
Ok(lhs
.zip(rhs)
.map(|(lhs, rhs)| ScalarIndexExpr::Or(Box::new(lhs), Box::new(rhs))))
}
ScalarIndexExpr::Not(inner) => {
// Negating a pruned AND would turn a safe superset into an unsafe
// subset. Preserve the complete subtree and let Lance track NULLs.
Ok(recurse(inner, true, remaining)?.map(|inner| ScalarIndexExpr::Not(Box::new(inner))))

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Negated IS TRUE and IS FALSE silently omit matching NULL rows through this newly enabled NOT branch. The pinned Lance planner represents key IS TRUE as key = true, so it reports NULL rows as SQL NULL instead of FALSE. Negation then excludes those rows from the candidate mask. Reapplying the complete filter cannot recover rows already excluded. Both BTree and Bitmap segments fail, with and without stable row IDs; the base implementation instead takes the correct scoped fallback for these predicates.

Preserve the scoped fallback for affected predicates until the planner repair is available, and cover both negated boolean tests on nullable columns.

Reproducer

Add this test to tests/c_api_test.rs, using its existing fixture helpers:

#[test]
fn gate_scalar_segment_not_is_true_keeps_nulls() {
    let key = Arc::new(arrow_array::BooleanArray::from(vec![
        None, Some(true), Some(false), Some(false),
        None, Some(true), Some(false), Some(true),
        None, Some(false), Some(true), Some(false),
    ]));
    let (_tmp, uri, uuids) = create_scalar_segment_fixture_from_key(
        lance_index::IndexType::BTree, false, None, &[&[0, 1], &[2]], key,
    );
    let (ids, _) = scalar_segment_ids(
        &uri, &uuids[0], &[0, 1], "NOT (key IS TRUE)", None, 0,
    );
    assert_eq!(ids, vec![0, 2, 3, 4, 6]);
}

Run cargo test --locked --test c_api_test gate_scalar_segment_not_is_true_keeps_nulls. On this head, the assertion fails with actual IDs [2, 3, 6] instead of [0, 2, 3, 4, 6]; the missing IDs are the NULL rows. A scan with scalar indices disabled returns the expected IDs.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed in 78626b1. I reproduced the missing NULL rows, then added a conservative guard on the original filter expression before scalar planning. Filters containing IS [NOT] TRUE/FALSE now retain the complete scoped scan, with scalar_segment_fallback_boolean_truth_test=1 and no segment search. The guard stays in place until the dependency includes lance-format/lance#9568; ordinary NOT over equality remains accelerated. Regressions cover BTree/Bitmap with and without stable row IDs, both negated truth tests, positive truth tests, nested AND/OR, fragment boundaries, and LIMIT/OFFSET. Validation: 455 Rust tests and all 3 native C/C++ integration tests passed, plus formatting, cargo check, and Clippy.

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Fixed in 78626b1: the original filter is checked before scalar planning, so boolean truth tests retain the complete scoped scan and matching NULL rows. Verified BTree/Bitmap with and without stable row IDs, nested predicates, fragment scope and pagination; ordinary NOT (key = true) still uses segment search.

}
}
}

struct SegmentIndexLoader<'a> {
index: Arc<dyn ScalarIndex>,
index_name: &'a str,
column: &'a str,
}

#[async_trait::async_trait]
impl ScalarIndexLoader for SegmentIndexLoader<'_> {
async fn load_index(
&self,
column: &str,
index_name: &str,
_metrics: &dyn MetricsCollector,
) -> Result<Arc<dyn ScalarIndex>> {
if column != self.column || index_name != self.index_name {
return Err(invalid(
"scalar expression references an index outside the selected segment",
));
}
_ => None,
// Reuse the one physical segment; a dataset loader would search every
// segment of the logical index and violate the caller's restriction.
Ok(self.index.clone())
}
}

Expand Down Expand Up @@ -193,6 +261,18 @@ impl PreparedScalarSegment {
let Some(filter) = reader.get_expr_filter()? else {
return Ok(Some("no_filter"));
};
// The pinned Lance planner lowers IS TRUE/FALSE to nullable equality.
// Under negation this loses matching NULL rows before the full recheck.
// Inspect the original expression before planning erases that distinction;
// remove this fallback only after adopting lance-format/lance#9568.
if filter.exists(|expr| {
Ok(matches!(
expr,
Expr::IsTrue(_) | Expr::IsFalse(_) | Expr::IsNotTrue(_) | Expr::IsNotFalse(_)
))
})? {
return Ok(Some("boolean_truth_test"));
}
let stored_schema: arrow_schema::Schema = self.dataset.schema().into();
// get_expr_filter validates against the scanner's full filterable
// schema, including metadata columns absent from the stored schema.
Expand All @@ -208,19 +288,25 @@ impl PreparedScalarSegment {
let planner = Planner::new(Arc::new(stored_schema));
let index_info = self.dataset.scalar_index_info().await?;
let filter_plan = planner.create_filter_plan(filter, &index_info, true)?;
let Some(search) = filter_plan
.index_query
.as_ref()
.and_then(|expr| driver(expr, &index_meta.name))
else {
let Some(expr) = filter_plan.index_query.as_ref() else {
return Ok(Some("no_driver"));
};
if search.column != field.name {
return Ok(Some("field_path"));
}
let mut remaining = MAX_EXPRESSION_NODES;
let expr = match scoped_expression(
expr,
&index_meta.name,
&field.name,
false,
0,
&mut remaining,
) {
Ok(Some(expr)) => expr,
Ok(None) => return Ok(Some("no_driver")),
Err(reason) => return Ok(Some(reason)),
};
let index = self
.dataset
.open_scalar_index(&search.column, &self.segment_uuid, metrics)
.open_scalar_index(&field.name, &self.segment_uuid, metrics)
.await?;
// These implementations can return exact candidates. Keep the runtime
// Exact check below: a type alone is not a guarantee for every query.
Expand All @@ -235,26 +321,88 @@ impl PreparedScalarSegment {
return Ok(Some("row_id_domain"));
}
let started = Instant::now();
let result = index.search(search.query.as_ref(), metrics).await?;
let loader = SegmentIndexLoader {
index,
index_name: &index_meta.name,
column: &field.name,
};
let result = expr.evaluate(&loader, metrics).await?;
stats.all_times.insert(
"scalar_segment_search_time".into(),
started.elapsed().as_nanos().min(usize::MAX as u128) as usize,
);
stats
.all_counts
.insert("scalar_segments_searched".into(), 1);
let SearchResult::Exact(rows) = result else {
if !result.is_exact() {
return Ok(Some("inexact_result"));
};
stats.all_counts.insert(
"scalar_segment_candidate_rows".into(),
rows.len().unwrap_or(0) as usize,
);
}
if let Some(rows) = result.upper.max_len() {
stats
.all_counts
.insert("scalar_segment_candidate_rows".into(), rows as usize);
} else {
// A complement mask has no finite cardinality without a row universe.
// Do not report zero candidates for a successful NOT search.
stats
.all_counts
.insert("scalar_segment_candidate_rows_unknown".into(), 1);
}
// Do not truncate candidates at LIMIT. The reader evaluates the complete
// filter before applying its existing limit/offset operators.
// The raw selected bitmap can overlap NULL rows; the full filter removes
// those as well. The metric above counts semantic TRUE rows, not mask size.
reader.with_row_addr_prefilter(RowAddrMask::from_allowed(rows.selected_rows().clone()));
// filter before applying its existing limit/offset operators. Its explicit
// fragment domain also bounds complement masks produced by NOT.
reader.with_row_addr_prefilter(result.upper);
Ok(None)
}
}

#[cfg(test)]
mod tests {
use super::*;
use lance_index::scalar::SargableQuery;
use lance_index::scalar::expression::ScalarIndexSearch;

fn leaf(index_name: &str) -> ScalarIndexExpr {
ScalarIndexExpr::Query(ScalarIndexSearch {
column: "key".into(),
index_name: index_name.into(),
index_type: "BTree".into(),
query: Arc::new(SargableQuery::IsNull()),
needs_recheck: false,
fragment_bitmap: None,
})
}

fn select(
expr: &ScalarIndexExpr,
) -> std::result::Result<Option<ScalarIndexExpr>, &'static str> {
let mut remaining = MAX_EXPRESSION_NODES;
scoped_expression(expr, "selected", "key", false, 0, &mut remaining)
}

#[test]
fn never_negate_a_partial_candidate_expression() {
let partial = ScalarIndexExpr::And(Box::new(leaf("selected")), Box::new(leaf("other")));
assert!(select(&partial).unwrap().is_some());
let negated = ScalarIndexExpr::Not(Box::new(partial));
assert!(select(&negated).unwrap().is_none());
let alternative = ScalarIndexExpr::Or(Box::new(leaf("selected")), Box::new(negated));
assert!(select(&alternative).unwrap().is_none());
}

#[test]
fn expression_budget_bounds_depth_and_concurrent_searches() {
let mut deep = leaf("selected");
for _ in 0..=MAX_EXPRESSION_DEPTH {
deep = ScalarIndexExpr::Not(Box::new(deep));
}
assert_eq!(select(&deep).unwrap_err(), "expression_budget");
let mut wide = leaf("selected");
for _ in 0..6 {
wide = ScalarIndexExpr::Or(Box::new(wide.clone()), Box::new(wide));
}
assert!(select(&wide).unwrap().is_some());
wide = ScalarIndexExpr::Or(Box::new(wide.clone()), Box::new(wide));
assert_eq!(select(&wide).unwrap_err(), "expression_budget");
}
}
Loading
Loading