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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ Based on the [liblance RFC](https://github.com/lance-format/lance/discussions/60
|--------|-----------|-------------|
| [x] | Async scan | Callback-based `lance_scanner_scan_async()` for non-blocking scans |
| [x] | Dataset metadata | `lance_dataset_version()`, `lance_dataset_count_rows()`, `lance_dataset_latest_version()` |
| [x] | Substrait filter pushdown | `lance_scanner_set_substrait_filter()` accepts a serialized Substrait `ExtendedExpression` (preferred over SQL strings for query engines) |
| [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 |

## Building

Expand Down
16 changes: 16 additions & 0 deletions include/lance/lance.h
Original file line number Diff line number Diff line change
Expand Up @@ -906,6 +906,22 @@ int32_t lance_scanner_set_substrait_filter(
size_t len
);

/**
* Add an SQL filter that is combined with the selected primary filter using
* AND. The primary filter is the Substrait filter when set, otherwise it is
* the SQL filter passed to `lance_scanner_new`. Multiple additional SQL
* filters are also combined using AND.
*
* Must be called before the scan starts. The filter string is copied.
*
* @param filter Non-NULL, non-empty SQL filter expression
* @return 0 on success, -1 on error
*/
int32_t lance_scanner_additional_sql_filter(
LanceScanner* scanner,
const char* filter
);

/** Type of a dynamically named scan metric. */
typedef enum {
LANCE_SCAN_METRIC_COUNT = 0,
Expand Down
7 changes: 7 additions & 0 deletions include/lance/lance.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -1185,6 +1185,13 @@ class Scanner {
return substrait_filter(bytes.data(), bytes.size());
}

/// Add an SQL filter that is combined with the selected primary filter using AND.
Scanner& additional_sql_filter(const std::string& filter) {
if (lance_scanner_additional_sql_filter(handle_.get(), filter.c_str()) != 0)
check_error();
return *this;
}

/// Register a non-null callback for scan statistics after successful full exhaustion.
/// The registration applies to every stream derived from this scanner, including
/// concurrent streams and streams created after an earlier callback returns. The
Expand Down
83 changes: 77 additions & 6 deletions src/scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ pub struct LanceScanner {
columns: Option<Vec<String>>,
filter: Option<String>,
substrait_filter: Option<Vec<u8>>,
additional_sql_filters: Vec<String>,
limit: Option<i64>,
offset: Option<i64>,
batch_size: Option<usize>,
Expand Down Expand Up @@ -124,6 +125,7 @@ impl LanceScanner {
columns: None,
filter: None,
substrait_filter: None,
additional_sql_filters: Vec::new(),
limit: None,
offset: None,
batch_size: None,
Expand Down Expand Up @@ -177,6 +179,36 @@ impl LanceScanner {
Ok(())
}

fn apply_filter(&self, scanner: &mut lance::dataset::scanner::Scanner) -> Result<()> {
if let Some(substrait) = &self.substrait_filter {
scanner.filter_substrait(substrait)?;
} else if let Some(sql) = &self.filter {
scanner.filter(sql)?;
}

if self.additional_sql_filters.is_empty() {
return Ok(());
}

// Let Lance resolve every SQL expression against the scanner's full
// filterable schema. Besides stored columns, this includes metadata
// columns and query-generated columns such as _distance and _score.
let mut combined = scanner.get_expr_filter()?;
for sql in &self.additional_sql_filters {
let mut additional_scanner = scanner.clone();
additional_scanner.filter(sql)?;
let additional = additional_scanner
.get_expr_filter()?
.expect("additional SQL filter exists");
combined = Some(match combined {
Some(existing) => existing.and(additional),
None => additional,
});
}
scanner.filter_expr(combined.expect("additional SQL filter exists"));
Ok(())
}

/// Build the underlying Scanner and open a stream.
fn materialize_stream(&mut self) -> Result<()> {
let prepared_scanner = self.build_scanner()?;
Expand All @@ -193,12 +225,6 @@ impl LanceScanner {
if let Some(cols) = &self.columns {
scanner.project(cols)?;
}
// Substrait filter takes precedence over SQL filter when both are set.
if let Some(bytes) = &self.substrait_filter {
scanner.filter_substrait(bytes)?;
} else if let Some(filter) = &self.filter {
scanner.filter(filter)?;
}
if self.limit.is_some() || self.offset.is_some() {
scanner.limit(self.limit, self.offset)?;
}
Expand Down Expand Up @@ -270,6 +296,7 @@ impl LanceScanner {
} else {
None
};
self.apply_filter(&mut scanner)?;
if let Some(callback) = &self.scan_statistics_callback {
scanner.scan_stats_callback(callback.clone());
}
Expand Down Expand Up @@ -782,6 +809,50 @@ unsafe fn scanner_set_substrait_filter_inner(
Ok(0)
}

/// Add an SQL filter that is combined with the scanner's selected primary filter using AND.
///
/// The primary filter is the Substrait filter when one is set, otherwise it is the SQL filter
/// passed to `lance_scanner_new`. Multiple additional SQL filters are also combined using AND.
/// This must be called before the scan starts. The string is copied into the scanner.
///
/// Returns 0 on success, -1 on error (check `lance_last_error_*`).
#[unsafe(no_mangle)]
pub unsafe extern "C" fn lance_scanner_additional_sql_filter(
scanner: *mut LanceScanner,
filter: *const c_char,
) -> i32 {
scanner_poison_check!(scanner, -1);
scanner_ffi_try!(scanner, unsafe {
scanner_additional_sql_filter_inner(scanner, filter)
})
}

unsafe fn scanner_additional_sql_filter_inner(
scanner: *mut LanceScanner,
filter: *const c_char,
) -> Result<i32> {
if scanner.is_null() {
return Err(lance_core::Error::invalid_input_source(
"scanner is NULL".into(),
));
}
let filter = unsafe { helpers::parse_c_string(filter)? }
.ok_or_else(|| lance_core::Error::invalid_input_source("filter must not be NULL".into()))?;
if filter.is_empty() {
return Err(lance_core::Error::invalid_input_source(
"additional SQL filter must be non-empty".into(),
));
}
let scanner = unsafe { &mut *scanner };
if scanner.scan_started.load(Ordering::Acquire) {
return Err(lance_core::Error::invalid_input_source(
"additional SQL filter must be set before the scan starts".into(),
));
}
scanner.additional_sql_filters.push(filter.to_string());
Ok(0)
}

/// Register a callback that receives execution statistics after the scan stream
/// is fully consumed to EOF.
///
Expand Down
Loading
Loading