From b5765a2d7bf3d9a79b7e711f98eaefbc568147f6 Mon Sep 17 00:00:00 2001 From: geruh Date: Fri, 2 Oct 2026 20:29:44 -0500 Subject: [PATCH 1/2] fix(scan): report retained physical row counts Co-authored-by: Cursor --- rust/lance/src/io/exec/filtered_read.rs | 55 +++++++-- rust/lance/src/io/exec/scan.rs | 158 ++++++++++++++++++++++++ 2 files changed, 206 insertions(+), 7 deletions(-) diff --git a/rust/lance/src/io/exec/filtered_read.rs b/rust/lance/src/io/exec/filtered_read.rs index 4f6c74d70c4..750eea3c6b9 100644 --- a/rust/lance/src/io/exec/filtered_read.rs +++ b/rust/lance/src/io/exec/filtered_read.rs @@ -2770,6 +2770,17 @@ impl FilteredReadExec { lazy_stream, ))) } + + fn retained_physical_row_count(&self, fragments: &[Fragment]) -> Option { + if self.dataset.manifest().writer_version.is_none() { + return None; + } + fragments + .iter() + .map(|fragment| fragment.physical_rows) + .sum::>() + .map(|physical_rows| physical_rows as u64) + } } /// How many batches run concurrently. Each batch's read already carries @@ -3353,13 +3364,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) @@ -4709,6 +4731,10 @@ 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()) @@ -4716,6 +4742,21 @@ mod tests { .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 diff --git a/rust/lance/src/io/exec/scan.rs b/rust/lance/src/io/exec/scan.rs index 1eeec1824e3..d795a42dba3 100644 --- a/rust/lance/src/io/exec/scan.rs +++ b/rust/lance/src/io/exec/scan.rs @@ -809,6 +809,25 @@ impl ExecutionPlan for LanceScanExec { } fn partition_statistics(&self, _partition: Option) -> Result> { + 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::>() + .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 @@ -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>, + #[case] has_missing_row_counts: bool, + #[case] expected: Precision, + ) { + 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, + ) { + let mut dataset = gen_batch() + .col("x", array::step::()) + .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] From d232d38f3b71e1a8b4dcd4f0f6e03c7db5d1b790 Mon Sep 17 00:00:00 2001 From: geruh Date: Fri, 2 Oct 2026 20:42:10 -0500 Subject: [PATCH 2/2] clippy happy --- rust/lance/src/io/exec/filtered_read.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/rust/lance/src/io/exec/filtered_read.rs b/rust/lance/src/io/exec/filtered_read.rs index 750eea3c6b9..2e05828c52f 100644 --- a/rust/lance/src/io/exec/filtered_read.rs +++ b/rust/lance/src/io/exec/filtered_read.rs @@ -2772,9 +2772,7 @@ impl FilteredReadExec { } fn retained_physical_row_count(&self, fragments: &[Fragment]) -> Option { - if self.dataset.manifest().writer_version.is_none() { - return None; - } + self.dataset.manifest().writer_version.as_ref()?; fragments .iter() .map(|fragment| fragment.physical_rows)