diff --git a/datafusion/core/tests/parquet/external_access_plan.rs b/datafusion/core/tests/parquet/external_access_plan.rs index 8fd9689ae3a8d..27387448296e2 100644 --- a/datafusion/core/tests/parquet/external_access_plan.rs +++ b/datafusion/core/tests/parquet/external_access_plan.rs @@ -23,6 +23,7 @@ use std::sync::Arc; use crate::parquet::utils::MetricsFinder; use crate::parquet::{Scenario, create_data_batch}; +use arrow::buffer::BooleanBuffer; use arrow::datatypes::SchemaRef; use arrow::util::pretty::pretty_format_batches; use datafusion::common::Result; @@ -218,6 +219,79 @@ async fn row_selection_extension_spanning_row_groups() { assert_ne!(bytes_scanned, 0, "metrics : {parquet_metrics:#?}",); } +#[tokio::test] +async fn bitmap_row_selection_extension_spanning_row_groups() { + // Repeat the cross-row-group case with a bitmap-backed selection. + let parquet_metrics = TestFull { + access_plan: None, + row_selection: Some(ParquetRowSelection::new(RowSelection::from_boolean_buffer( + BooleanBuffer::from(vec![ + false, false, false, false, true, true, true, false, false, false, + ]), + ))), + expected_rows: 3, + expected_output: Some(&[ + "+------+------------+", + "| utf8 | large_utf8 |", + "+------+------------+", + "| | |", + "| e | e |", + "| f | f |", + "+------+------------+", + ]), + predicate: None, + } + .run() + .await + .unwrap(); + + let bytes_scanned = metric_value(&parquet_metrics, "bytes_scanned").unwrap(); + assert_ne!(bytes_scanned, 0, "metrics : {parquet_metrics:#?}",); +} + +#[tokio::test] +async fn bitmap_row_selection_extension_with_predicate() { + // Pruning removes row group 0, leaving the bitmap-selected rows in group 1. + let parquet_metrics = TestFull { + access_plan: None, + row_selection: Some(ParquetRowSelection::new(RowSelection::from_boolean_buffer( + BooleanBuffer::from(vec![ + false, false, true, false, true, true, true, false, false, false, + ]), + ))), + expected_rows: 2, + expected_output: Some(&[ + "+------+------------+", + "| utf8 | large_utf8 |", + "+------+------------+", + "| e | e |", + "| f | f |", + "+------+------------+", + ]), + predicate: Some(col("utf8").eq(lit("e"))), + } + .run() + .await + .unwrap(); + + // Verify that statistics pruned row group 0. + let row_groups_pruned_statistics = parquet_metrics + .sum_by_name("row_groups_pruned_statistics") + .unwrap(); + if let MetricValue::PruningMetrics { + pruning_metrics, .. + } = row_groups_pruned_statistics + { + assert_eq!( + pruning_metrics.pruned(), + 1, + "metrics : {parquet_metrics:#?}" + ); + } else { + unreachable!("metrics `row_groups_pruned_statistics` should exist") + } +} + #[tokio::test] async fn bad_row_selection_extension() { // selection specifies fewer rows than the file actually contains diff --git a/datafusion/datasource-parquet/src/access_plan.rs b/datafusion/datasource-parquet/src/access_plan.rs index 1e9bae0ff6ba3..544aa9a2e0554 100644 --- a/datafusion/datasource-parquet/src/access_plan.rs +++ b/datafusion/datasource-parquet/src/access_plan.rs @@ -16,6 +16,8 @@ // under the License. use crate::sort::reverse_row_selection; +use arrow::array::BooleanBufferBuilder; +use arrow::buffer::BooleanBuffer; use arrow::datatypes::Schema; use datafusion_common::{Result, assert_eq_or_internal_err, exec_err}; use datafusion_physical_expr::expressions::Column; @@ -160,16 +162,69 @@ impl RowGroupAccess { } } -/// Single-pass cursor over a file-level [`RowSelection`]. +/// Splits an overall selection into row groups in one pass. +enum OverallRowSelectionCursor { + Mask { mask: BooleanBuffer, offset: usize }, + Selectors(SelectorRowSelectionCursor), +} + +impl OverallRowSelectionCursor { + fn new(selection: RowSelection) -> Self { + match selection.as_mask() { + Some(mask) => Self::Mask { + mask: mask.clone(), + offset: 0, + }, + None => Self::Selectors(SelectorRowSelectionCursor::new(selection)), + } + } + + /// Takes the next row group, or `None` if too few rows remain. + fn take_row_group(&mut self, row_group_rows: usize) -> Option { + match self { + Self::Mask { mask, offset } => { + if row_group_rows > mask.len() - *offset { + return None; + } + + let end = *offset + row_group_rows; + let row_group_mask = mask.slice(*offset, row_group_rows); + *offset = end; + let selected_rows = row_group_mask.count_set_bits(); + + Some(if selected_rows == 0 { + RowGroupAccess::Skip + } else if selected_rows == row_group_rows { + RowGroupAccess::Scan + } else { + RowGroupAccess::Selection(RowSelection::from_boolean_buffer( + row_group_mask, + )) + }) + } + Self::Selectors(cursor) => cursor.take_row_group(row_group_rows), + } + } + + fn total_rows(self) -> usize { + match self { + Self::Mask { mask, .. } => mask.len(), + Self::Selectors(cursor) => cursor.total_rows(), + } + } +} + +/// Cursor for a selector-backed selection. /// /// `take` returns the next selector fragment capped to the requested row count, /// splitting the current selector when it straddles a row group boundary. -struct OverallRowSelectionCursor { +struct SelectorRowSelectionCursor { selector_iter: std::vec::IntoIter, current: Option, + consumed_rows: usize, } -impl OverallRowSelectionCursor { +impl SelectorRowSelectionCursor { fn new(selection: RowSelection) -> Self { let selectors: Vec = selection.into(); let mut selector_iter = selectors.into_iter(); @@ -177,6 +232,7 @@ impl OverallRowSelectionCursor { Self { selector_iter, current, + consumed_rows: 0, } } @@ -189,6 +245,7 @@ impl OverallRowSelectionCursor { fn take(&mut self, max_rows: usize) -> Option { let sel = self.current?; let row_count = sel.row_count.min(max_rows); + self.consumed_rows += row_count; self.current = if row_count < sel.row_count { Some(RowSelector { row_count: sel.row_count - row_count, @@ -204,9 +261,21 @@ impl OverallRowSelectionCursor { }) } - fn remaining_rows(self) -> usize { - self.current.map_or(0, |s| s.row_count) - + self.selector_iter.map(|s| s.row_count).sum::() + fn take_row_group(&mut self, row_group_rows: usize) -> Option { + let mut builder = RowGroupAccessBuilder::new(row_group_rows); + while builder.remaining > 0 { + builder.push(self.take(builder.remaining)?); + } + Some(builder.into_access()) + } + + fn total_rows(self) -> usize { + self.consumed_rows + + self.current.map_or(0, |selector| selector.row_count) + + self + .selector_iter + .map(|selector| selector.row_count) + .sum::() } } @@ -256,6 +325,29 @@ impl RowGroupAccessBuilder { } } +/// Returns the selection length. +fn row_selection_len(selection: &RowSelection) -> usize { + match selection.as_mask() { + Some(mask) => mask.len(), + None => selection.iter().map(|selector| selector.row_count).sum(), + } +} + +/// Converts to bitmap backing if needed. +fn into_mask_backed(selection: RowSelection) -> RowSelection { + match selection.as_mask() { + Some(_) => selection, + None => { + let total_rows = row_selection_len(&selection); + let mut mask = BooleanBufferBuilder::new(total_rows); + for selector in selection.iter() { + mask.append_n(selector.row_count, !selector.skip); + } + RowSelection::from_boolean_buffer(mask.finish()) + } + } +} + impl ParquetAccessPlan { /// Create a new `ParquetAccessPlan` that scans all row groups pub fn new_all(row_group_count: usize) -> Self { @@ -298,35 +390,20 @@ impl ParquetAccessPlan { selection: RowSelection, row_group_meta_data: &[RowGroupMetaData], ) -> Result { - // Keep this as a single pass over the selector stream rather than - // repeatedly calling `RowSelection::split_off` per row group. The - // `split_off` version is simpler, but it clones/retains substantially - // more selector buffer capacity for highly fragmented selections. - let mut cursor = OverallRowSelectionCursor::new(selection); - - let mut selection_rows = 0usize; - let mut file_rows = 0usize; - - let mut row_groups = Vec::with_capacity(row_group_meta_data.len()); - for rg_meta in row_group_meta_data { - let rg_rows = rg_meta.num_rows() as usize; - file_rows += rg_rows; - - let mut builder = RowGroupAccessBuilder::new(rg_rows); - while builder.remaining > 0 { - let Some(selector) = cursor.take(builder.remaining) else { - break; - }; - selection_rows += selector.row_count; - builder.push(selector); - } - - row_groups.push(builder.into_access()); - } + let file_rows = row_group_meta_data + .iter() + .map(|rg| rg.num_rows() as usize) + .sum::(); - selection_rows += cursor.remaining_rows(); + let mut cursor = OverallRowSelectionCursor::new(selection); + let row_groups = row_group_meta_data + .iter() + .map_while(|rg_meta| cursor.take_row_group(rg_meta.num_rows() as usize)) + .collect::>(); - if selection_rows != file_rows { + // Count selector rows during traversal. + let selection_rows = cursor.total_rows(); + if row_groups.len() != row_group_meta_data.len() || selection_rows != file_rows { return exec_err!( "Invalid Parquet RowSelection. File has {file_rows} rows, \ but selection specifies {selection_rows} rows." @@ -392,6 +469,11 @@ impl ParquetAccessPlan { RowGroupAccess::Skip => RowGroupAccess::Skip, RowGroupAccess::Scan => RowGroupAccess::Selection(selection), RowGroupAccess::Selection(existing_selection) => { + // Keep intersections bitmap-backed when the existing selection is. + let selection = match existing_selection.as_mask() { + Some(_) => into_mask_backed(selection), + None => selection, + }; RowGroupAccess::Selection(existing_selection.intersection(&selection)) } } @@ -419,6 +501,8 @@ impl ParquetAccessPlan { /// is returned for *all* the rows in the row groups that are not skipped. /// Thus it includes a `Select` selection for any [`RowGroupAccess::Scan`]. /// + /// The result is bitmap-backed if any row-group selection is bitmap-backed. + /// /// # Errors /// /// Returns an error if any specified row selection does not specify @@ -494,10 +578,7 @@ impl ParquetAccessPlan { let RowGroupAccess::Selection(selection) = rg else { continue; }; - let rows_in_selection = selection - .iter() - .map(|selection| selection.row_count) - .sum::(); + let rows_in_selection = row_selection_len(selection); let row_group_row_count = rg_meta.num_rows(); assert_eq_or_internal_err!( @@ -509,24 +590,53 @@ impl ParquetAccessPlan { ); } - let total_selection: RowSelection = self - .row_groups - .into_iter() - .zip(row_group_meta_data.iter()) - .flat_map(|(rg, rg_meta)| { + // Preserve bitmap backing across mixed row-group selections. + let any_selection_mask_backed = self.row_groups.iter().any(|rg| { + matches!(rg, RowGroupAccess::Selection(selection) if selection.as_mask().is_some()) + }); + + let total_selection = if any_selection_mask_backed { + let total_rows = self + .row_groups + .iter() + .zip(row_group_meta_data.iter()) + .filter(|(rg, _)| rg.should_scan()) + .map(|(_, rg_meta)| rg_meta.num_rows() as usize) + .sum(); + let mut mask = BooleanBufferBuilder::new(total_rows); + + for (rg, rg_meta) in self.row_groups.into_iter().zip(row_group_meta_data) { match rg { + RowGroupAccess::Skip => {} + RowGroupAccess::Scan => { + mask.append_n(rg_meta.num_rows() as usize, true) + } + RowGroupAccess::Selection(selection) => match selection.as_mask() { + Some(buffer) => mask.append_buffer(buffer), + None => { + for selector in selection.iter() { + mask.append_n(selector.row_count, !selector.skip); + } + } + }, + } + } + + RowSelection::from_boolean_buffer(mask.finish()) + } else { + // Keep selector backing when possible. + self.row_groups + .into_iter() + .zip(row_group_meta_data.iter()) + .flat_map(|(rg, rg_meta)| match rg { RowGroupAccess::Skip => vec![], RowGroupAccess::Scan => { - // need a row group access to scan the entire row group (need row group counts) vec![RowSelector::select(rg_meta.num_rows() as usize)] } - RowGroupAccess::Selection(selection) => { - let selection: Vec = selection.into(); - selection - } - } - }) - .collect(); + RowGroupAccess::Selection(selection) => selection.into(), + }) + .collect() + }; Ok(Some(total_selection)) } @@ -787,6 +897,7 @@ impl PreparedAccessPlan { #[cfg(test)] mod test { use super::*; + use arrow::buffer::BooleanBuffer; use datafusion_common::assert_contains; use parquet::basic::LogicalType; use parquet::file::metadata::ColumnChunkMetaData; @@ -910,6 +1021,95 @@ mod test { ); } + #[test] + fn test_mixed_backings_promote_to_mask() { + // Mixed backing produces a bitmap-backed overall selection. + let access_plan = ParquetAccessPlan::new(vec![ + RowGroupAccess::Selection(RowSelection::from_boolean_buffer( + BooleanBuffer::from(vec![ + true, false, false, false, false, false, false, false, false, false, + ]), + )), + RowGroupAccess::Selection(RowSelection::from(vec![ + RowSelector::select(10), + RowSelector::skip(10), + ])), + RowGroupAccess::Skip, + RowGroupAccess::Skip, + ]); + + let row_selection = access_plan + .into_overall_row_selection(&ROW_GROUP_METADATA) + .unwrap() + .unwrap(); + + let mut expected = vec![true]; + expected.extend(vec![false; 9]); + expected.extend(vec![true; 10]); + expected.extend(vec![false; 10]); + let expected = BooleanBuffer::from(expected); + assert_eq!(row_selection.as_mask(), Some(&expected)); + } + + #[test] + fn test_scan_selection_preserves_mask_backing() { + let mut access_plan = ParquetAccessPlan::new(vec![RowGroupAccess::Selection( + RowSelection::from_boolean_buffer(BooleanBuffer::from(vec![ + true, true, false, false, true, true, false, false, true, true, + ])), + )]); + + // Simulate selector-backed page pruning. + access_plan.scan_selection( + 0, + RowSelection::from(vec![RowSelector::select(5), RowSelector::skip(5)]), + ); + // Exercise the already bitmap-backed path. + access_plan.scan_selection( + 0, + RowSelection::from_boolean_buffer(BooleanBuffer::new_set(10)), + ); + + let selection = access_plan + .into_overall_row_selection(&ROW_GROUP_METADATA[..1]) + .unwrap() + .unwrap(); + let expected = BooleanBuffer::from(vec![ + true, true, false, false, true, false, false, false, false, false, + ]); + assert_eq!(selection.as_mask(), Some(&expected)); + } + + #[test] + fn test_scan_selection_preserves_selector_backing() { + let mut access_plan = ParquetAccessPlan::new(vec![RowGroupAccess::Selection( + RowSelection::from(vec![RowSelector::select(6), RowSelector::skip(4)]), + )]); + + access_plan.scan_selection( + 0, + RowSelection::from(vec![ + RowSelector::skip(2), + RowSelector::select(5), + RowSelector::skip(3), + ]), + ); + + let selection = access_plan + .into_overall_row_selection(&ROW_GROUP_METADATA[..1]) + .unwrap() + .unwrap(); + assert_eq!(selection.as_mask(), None); + assert_eq!( + selection, + RowSelection::from(vec![ + RowSelector::skip(2), + RowSelector::select(4), + RowSelector::skip(4), + ]) + ); + } + #[test] fn test_new_from_overall_row_selection() { let row_selection = RowSelection::from(vec![ @@ -946,19 +1146,25 @@ mod test { #[test] fn test_new_from_overall_row_selection_invalid_row_count() { - let row_selection = RowSelection::from(vec![RowSelector::select(99)]); + for selection_rows in [99, 101] { + let row_selection = + RowSelection::from(vec![RowSelector::select(selection_rows)]); - let err = ParquetAccessPlan::try_new_from_overall_row_selection( - row_selection, - &ROW_GROUP_METADATA, - ) - .unwrap_err() - .to_string(); + let err = ParquetAccessPlan::try_new_from_overall_row_selection( + row_selection, + &ROW_GROUP_METADATA, + ) + .unwrap_err() + .to_string(); - assert_contains!( - err, - "Invalid Parquet RowSelection. File has 100 rows, but selection specifies 99 rows" - ); + assert_contains!( + err, + format!( + "Invalid Parquet RowSelection. File has 100 rows, \ + but selection specifies {selection_rows} rows" + ) + ); + } } #[test] @@ -994,6 +1200,71 @@ mod test { ); } + #[test] + fn test_new_from_overall_mask_preserves_bitmap_backing() { + let partial_mask = vec![ + false, true, false, true, false, true, false, true, false, true, false, true, + false, true, false, true, false, true, false, true, false, true, false, true, + false, true, false, true, false, true, + ]; + let mut file_mask = vec![true; 10]; + file_mask.extend(vec![false; 20]); + file_mask.extend(&partial_mask); + file_mask.extend(vec![true; 40]); + + let access_plan = ParquetAccessPlan::try_new_from_overall_row_selection( + RowSelection::from_boolean_buffer(BooleanBuffer::from(file_mask)), + &ROW_GROUP_METADATA, + ) + .unwrap(); + + assert_eq!( + access_plan, + ParquetAccessPlan::new(vec![ + RowGroupAccess::Scan, + RowGroupAccess::Skip, + RowGroupAccess::Selection(RowSelection::from_boolean_buffer( + BooleanBuffer::from(partial_mask.clone()), + )), + RowGroupAccess::Scan, + ]) + ); + + // Skipped groups are omitted; scanned groups become set bits. + let overall = access_plan + .into_overall_row_selection(&ROW_GROUP_METADATA) + .unwrap() + .unwrap(); + let mut expected_overall = vec![true; 10]; + expected_overall.extend(partial_mask); + expected_overall.extend(vec![true; 40]); + let expected_overall = BooleanBuffer::from(expected_overall); + assert_eq!(overall.as_mask(), Some(&expected_overall)); + } + + #[test] + fn test_new_from_overall_mask_invalid_row_count() { + for selection_rows in [99, 101] { + let row_selection = + RowSelection::from_boolean_buffer(BooleanBuffer::new_set(selection_rows)); + + let err = ParquetAccessPlan::try_new_from_overall_row_selection( + row_selection, + &ROW_GROUP_METADATA, + ) + .unwrap_err() + .to_string(); + + assert_contains!( + err, + format!( + "Invalid Parquet RowSelection. File has 100 rows, \ + but selection specifies {selection_rows} rows" + ) + ); + } + } + #[test] fn test_invalid_too_few() { let access_plan = ParquetAccessPlan::new(vec![ diff --git a/datafusion/datasource-parquet/src/sort.rs b/datafusion/datasource-parquet/src/sort.rs index ea33fb0e2ecb2..ec846168ca074 100644 --- a/datafusion/datasource-parquet/src/sort.rs +++ b/datafusion/datasource-parquet/src/sort.rs @@ -54,6 +54,26 @@ pub fn reverse_row_selection( ) -> Result { let rg_metadata = parquet_metadata.row_groups(); + // Reverse bitmap slices without materializing selectors. + if row_selection.as_mask().is_some() { + let mut remaining = row_selection.clone(); + let mut row_group_selections = Vec::with_capacity(row_groups_to_scan.len()); + + for &rg_idx in row_groups_to_scan { + let num_rows = rg_metadata[rg_idx].num_rows() as usize; + row_group_selections.push(remaining.split_off(num_rows)); + } + + // No rows should remain after splitting all row groups. + debug_assert_eq!( + remaining.row_count() + remaining.skipped_row_count(), + 0, + "row selection covers more rows than the scanned row groups" + ); + + return Ok(row_group_selections.into_iter().rev().collect()); + } + // Build a mapping of row group index to its row range, but ONLY for // the row groups that are actually being scanned. // @@ -244,6 +264,7 @@ fn file_min_value(file: &PartitionedFile, col_idx: usize) -> Option mod tests { use crate::ParquetAccessPlan; use crate::RowGroupAccess; + use arrow::buffer::BooleanBuffer; use arrow::datatypes::{DataType, Field, Schema}; use bytes::Bytes; use parquet::arrow::ArrowWriter; @@ -358,6 +379,47 @@ mod tests { ); } + #[test] + fn test_prepared_access_plan_reverse_preserves_bitmap_backing() { + let metadata = create_test_metadata(vec![4, 3, 5, 2]); + let first_mask = vec![true, false, true, false]; + let third_mask = vec![false, true, true, false, true]; + let access_plan = ParquetAccessPlan::new(vec![ + RowGroupAccess::Selection(RowSelection::from_boolean_buffer( + BooleanBuffer::from(first_mask.clone()), + )), + RowGroupAccess::Skip, + RowGroupAccess::Selection(RowSelection::from_boolean_buffer( + BooleanBuffer::from(third_mask.clone()), + )), + RowGroupAccess::Scan, + ]); + + let prepared_plan = access_plan.prepare(metadata.row_groups()).unwrap(); + assert!( + prepared_plan + .row_selection + .as_ref() + .unwrap() + .as_mask() + .is_some() + ); + + let reversed_plan = prepared_plan.reverse(&metadata).unwrap(); + assert_eq!(reversed_plan.row_group_indexes, vec![3, 2, 0]); + let reversed_selection = reversed_plan.row_selection.unwrap(); + + // Fully scanned row group 3 becomes the leading set bits. + let expected = BooleanBuffer::from( + [true, true] + .into_iter() + .chain(third_mask) + .chain(first_mask) + .collect::>(), + ); + assert_eq!(reversed_selection.as_mask(), Some(&expected)); + } + #[test] fn test_prepared_access_plan_reverse_multi_row_group_selection() { // Test: row selection spanning multiple row groups