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
53 changes: 46 additions & 7 deletions rust/lance/src/io/exec/filtered_read.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2770,6 +2770,15 @@ impl FilteredReadExec {
lazy_stream,
)))
}

fn retained_physical_row_count(&self, fragments: &[Fragment]) -> Option<u64> {
self.dataset.manifest().writer_version.as_ref()?;
fragments
.iter()
.map(|fragment| fragment.physical_rows)
.sum::<Option<usize>>()
.map(|physical_rows| physical_rows as u64)
}
}

/// How many batches run concurrently. Each batch's read already carries
Expand Down Expand Up @@ -3353,13 +3362,24 @@ impl ExecutionPlan for FilteredReadExec {
.clone()
.unwrap_or_else(|| self.dataset.fragments().clone());

if fragments.iter().any(|f| f.num_rows().is_none()) {
return Err(DataFusionError::Internal(
"Fragments are missing row count stats".to_string(),
));
}

let total_rows: u64 = fragments.iter().map(|f| f.num_rows().unwrap() as u64).sum();
let total_rows = if self.options.with_deleted_rows {
if self.options.scan_range_before_filter.is_some()
|| self.options.scan_range_after_filter.is_some()
{
return Ok(Arc::new(Statistics::new_unknown(self.schema().as_ref())));
}
let Some(total_rows) = self.retained_physical_row_count(fragments.as_ref()) else {
return Ok(Arc::new(Statistics::new_unknown(self.schema().as_ref())));
};
total_rows
} else {
if fragments.iter().any(|f| f.num_rows().is_none()) {
return Err(DataFusionError::Internal(
"Fragments are missing row count stats".to_string(),
));
}
fragments.iter().map(|f| f.num_rows().unwrap() as u64).sum()
};

let Some(filter) = self.options.full_filter.as_ref() else {
// If there is no filter, we just return the total number of rows (sans any before-filter range)
Expand Down Expand Up @@ -4709,13 +4729,32 @@ mod tests {
.unwrap()
.with_projection(fixture.dataset.empty_projection().with_row_id());
let plan = fixture.make_plan(options).await;
assert_eq!(
plan.partition_statistics(None).unwrap().num_rows,
Precision::Exact(300)
);
let stream = plan.execute(0, Arc::new(TaskContext::default())).unwrap();
let num_rows = stream
.map_ok(|batch| batch.num_rows())
.try_fold(0, |acc, val| std::future::ready(Ok(acc + val)))
.await
.unwrap();
assert_eq!(num_rows, 300);
let mut scanner = fixture.dataset.scan();
scanner.with_row_id().include_deleted_rows();
assert_eq!(scanner.count_rows().await.unwrap(), 300);

let filter = fixture.filter_plan("not_indexed >= 250", false).await;
let options = base_options
.with_deleted_rows()
.unwrap()
.with_filter_plan(filter);
let plan = fixture.make_plan(options.clone()).await;
assert_eq!(
plan.partition_statistics(None).unwrap().num_rows,
Precision::Inexact(300)
);
fixture.test_plan(options, &u32s(vec![250..400])).await;
}

/// A stale (not rebuilt after a delete) index hit drops on the live view
Expand Down
158 changes: 158 additions & 0 deletions rust/lance/src/io/exec/scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -809,6 +809,25 @@ impl ExecutionPlan for LanceScanExec {
}

fn partition_statistics(&self, _partition: Option<usize>) -> Result<Arc<Statistics>> {
if self.config.with_make_deletions_null {
// Ranges use visible-row offsets. Trust `physical_rows` only with a writer version.
if self.range.is_some() || self.dataset.manifest().writer_version.is_none() {
return Ok(Arc::new(Statistics {
num_rows: Precision::Absent,
..Statistics::new_unknown(self.schema().as_ref())
}));
}
let num_rows = self
.fragments
.iter()
.map(|fragment| fragment.physical_rows)
.sum::<Option<usize>>()
.map_or(Precision::Absent, Precision::Exact);
return Ok(Arc::new(Statistics {
num_rows,
..Statistics::new_unknown(self.schema().as_ref())
}));
}
// Some fragments from older datasets might have the row count stats missing.
let (row_count, is_exact) =
self.fragments
Expand Down Expand Up @@ -982,6 +1001,145 @@ mod tests {
assert_eq!(scanned_rows(&scan).await, TOTAL_ROWS);
}

#[rstest]
#[case::known(None, false, Precision::Exact(TOTAL_ROWS))]
#[case::missing_rows(None, true, Precision::Absent)]
#[case::ranged(Some(0..10), false, Precision::Absent)]
#[tokio::test]
async fn deletion_null_statistics(
#[case] range: Option<Range<u64>>,
#[case] has_missing_row_counts: bool,
#[case] expected: Precision<usize>,
) {
let mut dataset = ranged_scan_dataset(LanceFileVersion::Stable)
.await
.as_ref()
.clone();
dataset.delete("x = 0").await.unwrap();
let mut fragments = dataset.fragments().clone();
if has_missing_row_counts {
Arc::make_mut(&mut fragments)[0].physical_rows = None;
}
let dataset = Arc::new(dataset);
let scan = LanceScanExec::new(
dataset.clone(),
fragments,
range.clone(),
Arc::new(dataset.schema().clone()),
LanceScanConfig {
with_row_id: true,
with_make_deletions_null: true,
..Default::default()
},
);

assert_eq!(scan.partition_statistics(None).unwrap().num_rows, expected);
if range.is_none() {
assert_eq!(scanned_rows(&scan).await, TOTAL_ROWS);
}
}

#[derive(Debug)]
enum PhysicalRowMetadata {
Present,
Untrusted,
Missing,
}

#[rstest]
#[case::legacy(
LanceFileVersion::Legacy,
"LanceScan:",
PhysicalRowMetadata::Present,
true,
Precision::Exact(8)
)]
#[case::stable(
LanceFileVersion::Stable,
"LanceRead:",
PhysicalRowMetadata::Present,
true,
Precision::Exact(8)
)]
#[case::legacy_untrusted(
LanceFileVersion::Legacy,
"LanceScan:",
PhysicalRowMetadata::Untrusted,
false,
Precision::Absent
)]
#[case::stable_untrusted(
LanceFileVersion::Stable,
"LanceRead:",
PhysicalRowMetadata::Untrusted,
false,
Precision::Absent
)]
#[case::legacy_missing_physical_rows(
LanceFileVersion::Legacy,
"LanceScan:",
PhysicalRowMetadata::Missing,
false,
Precision::Absent
)]
#[case::stable_missing_physical_rows(
LanceFileVersion::Stable,
"LanceRead:",
PhysicalRowMetadata::Missing,
false,
Precision::Absent
)]
#[tokio::test]
async fn include_deleted_rows_statistics(
#[case] version: LanceFileVersion,
#[case] scan_node: &str,
#[case] metadata: PhysicalRowMetadata,
#[case] deleted: bool,
#[case] expected: Precision<usize>,
) {
let mut dataset = gen_batch()
.col("x", array::step::<Int32Type>())
.into_ram_dataset_with_params(
FragmentCount::from(2),
FragmentRowCount::from(4),
Some(WriteParams {
max_rows_per_file: 4,
data_storage_version: Some(version),
..Default::default()
}),
)
.await
.unwrap();
if deleted {
dataset.delete("x = 0").await.unwrap();
}
let mut fragments = dataset.fragments().as_ref().clone();
match metadata {
PhysicalRowMetadata::Present => {}
PhysicalRowMetadata::Untrusted => {
Arc::make_mut(&mut dataset.manifest).writer_version = None;
fragments[0].physical_rows = Some(999);
}
PhysicalRowMetadata::Missing => {
fragments[0].physical_rows = None;
}
}
let mut scanner = dataset.scan();
scanner
.with_fragments(fragments)
.with_row_id()
.include_deleted_rows();
let plan = scanner.create_plan().await.unwrap();
let description = datafusion::physical_plan::displayable(plan.as_ref())
.indent(true)
.to_string();
assert!(description.contains(scan_node), "{description}");
assert_eq!(plan.partition_statistics(None).unwrap().num_rows, expected);
let batch = scanner.try_into_batch().await.unwrap();
assert_eq!(batch.num_rows(), 8);
assert_eq!(batch["_rowid"].null_count(), usize::from(deleted));
}

/// Verify that executing with target_partitions=1 produces the same row count as the
/// default context. Regression guard for the parallelism cap.
#[tokio::test]
Expand Down
Loading