Skip to content

Commit c3ef346

Browse files
authored
feat: preserve bitmap-backed Parquet row selections (#25731)
## Which issue does this PR close? - Closes #23883. - Closes #24488. - Follows up on #24186 using the row-group-local scan path merged in #25608. ## Rationale for this change An external index can supply a bitmap-backed Parquet row selection, but splitting a file-level selection into row groups or intersecting it with selector-backed page pruning still converts it to selectors. Keeping the bitmap avoids materializing fragmented selector runs and lets pruning use a bitwise intersection. ## What changes are included in this PR? - Split bitmap-backed file selections using bitmap slices, preserving partial-group masks and recognizing fully selected or skipped groups. Keep the existing single-pass selector path. - Promote incoming selectors to a bitmap when intersecting with an existing bitmap-backed selection. - Use Parquet's `RowSelection::total_row_count()` directly instead of introducing the temporary `row_selection_len` helper from #24186. - Add regression coverage for bitmap preservation and external selections when a dynamic predicate changes between row groups. Consolidate opener test imports at module scope. The deprecated `into_overall_row_selection` implementation is unchanged. Preparation and reverse scans use the local selections introduced by #25608. Runtime pruning remains disabled for scans with selections. ## What is the testing strategy for this PR? Unit tests cover selector/bitmap intersections, empty intersections, file-length validation, non-byte-aligned bitmap slicing, and preservation through preparation and reversal. Integration tests cover bitmap selections spanning row groups and statistics pruning. The dynamic-pruning regression covers both selector and bitmap external selections, with an unselected control scan that verifies runtime pruning occurs. Temporarily removing the selection guard makes the test fail because rebuilding with `None` returns unselected rows; the guard was restored after this check. Passed targeted tests: - `cargo test -p datafusion-datasource-parquet --lib access_plan::test --offline` (26 tests) - `cargo test -p datafusion-datasource-parquet --lib opener::test --offline` (56 tests) - `cargo test -p datafusion --test parquet_integration external_access_plan --offline` (16 tests) Pre-submission checks passed: `cargo fmt --all`, `cargo clippy --all-targets --all-features -- -D warnings`, and the checks in `uv run ./dev/rust_lint.sh`. The final HTML documentation check was rerun successfully with `uv run ./ci/scripts/check_docs_html.sh` after installing its missing `cargo-depgraph` prerequisite. ## Are there any user-facing changes? No public API or selected-row changes. Bitmap-backed external selections retain their representation through splitting and intersection rather than being converted to selectors.
1 parent 76f9fde commit c3ef346

3 files changed

Lines changed: 446 additions & 98 deletions

File tree

‎datafusion/core/tests/parquet/external_access_plan.rs‎

Lines changed: 64 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ use std::sync::Arc;
2323
use crate::parquet::utils::MetricsFinder;
2424
use crate::parquet::{Scenario, create_data_batch};
2525

26+
use arrow::buffer::BooleanBuffer;
2627
use arrow::datatypes::SchemaRef;
2728
use arrow::util::pretty::pretty_format_batches;
2829
use datafusion::common::Result;
@@ -189,33 +190,82 @@ async fn row_selection_extension() {
189190

190191
#[tokio::test]
191192
async fn row_selection_extension_spanning_row_groups() {
192-
// A selection whose selectors straddle the row group boundary (row 4 is the
193-
// last row of group 0, rows 5-6 are the first rows of group 1).
194-
let parquet_metrics = TestFull {
195-
access_plan: None,
196-
row_selection: Some(ParquetRowSelection::new(RowSelection::from(vec![
193+
// Row 4 is the last row of group 0; rows 5-6 begin group 1.
194+
for selection in [
195+
RowSelection::from(vec![
197196
RowSelector::skip(4),
198197
RowSelector::select(3),
199198
RowSelector::skip(3),
200-
]))),
201-
expected_rows: 3,
199+
]),
200+
RowSelection::from_boolean_buffer(BooleanBuffer::from(vec![
201+
false, false, false, false, true, true, true, false, false, false,
202+
])),
203+
] {
204+
let parquet_metrics = TestFull {
205+
access_plan: None,
206+
row_selection: Some(ParquetRowSelection::new(selection)),
207+
expected_rows: 3,
208+
expected_output: Some(&[
209+
"+------+------------+",
210+
"| utf8 | large_utf8 |",
211+
"+------+------------+",
212+
"| | |",
213+
"| e | e |",
214+
"| f | f |",
215+
"+------+------------+",
216+
]),
217+
predicate: None,
218+
}
219+
.run()
220+
.await
221+
.unwrap();
222+
223+
let bytes_scanned = metric_value(&parquet_metrics, "bytes_scanned").unwrap();
224+
assert_ne!(bytes_scanned, 0, "metrics : {parquet_metrics:#?}");
225+
}
226+
}
227+
228+
#[tokio::test]
229+
async fn bitmap_row_selection_extension_with_predicate() {
230+
// Pruning removes row group 0, leaving the bitmap-selected rows in group 1.
231+
let parquet_metrics = TestFull {
232+
access_plan: None,
233+
row_selection: Some(ParquetRowSelection::new(RowSelection::from_boolean_buffer(
234+
BooleanBuffer::from(vec![
235+
false, false, true, false, true, true, true, false, false, false,
236+
]),
237+
))),
238+
expected_rows: 2,
202239
expected_output: Some(&[
203240
"+------+------------+",
204241
"| utf8 | large_utf8 |",
205242
"+------+------------+",
206-
"| | |",
207243
"| e | e |",
208244
"| f | f |",
209245
"+------+------------+",
210246
]),
211-
predicate: None,
247+
predicate: Some(col("utf8").eq(lit("e"))),
212248
}
213249
.run()
214250
.await
215251
.unwrap();
216252

217-
let bytes_scanned = metric_value(&parquet_metrics, "bytes_scanned").unwrap();
218-
assert_ne!(bytes_scanned, 0, "metrics : {parquet_metrics:#?}");
253+
// Verify that statistics pruned row group 0.
254+
let row_groups_pruned_statistics = parquet_metrics
255+
.sum_by_name("row_groups_pruned_statistics")
256+
.unwrap();
257+
if let MetricValue::PruningMetrics {
258+
pruning_metrics, ..
259+
} = row_groups_pruned_statistics
260+
{
261+
assert_eq!(
262+
pruning_metrics.pruned(),
263+
1,
264+
"metrics : {parquet_metrics:#?}"
265+
);
266+
} else {
267+
unreachable!("metrics `row_groups_pruned_statistics` should exist")
268+
}
219269
}
220270

221271
#[tokio::test]
@@ -316,9 +366,9 @@ async fn mixed_bitmap_and_selector_selections() {
316366
// Each representation must keep coordinates local to its row group.
317367
TestFull {
318368
access_plan: Some(ParquetAccessPlan::new(vec![
319-
RowGroupAccess::Selection(RowSelection::from(
320-
arrow::buffer::BooleanBuffer::from(vec![true, false, true, false, false]),
321-
)),
369+
RowGroupAccess::Selection(RowSelection::from(BooleanBuffer::from(vec![
370+
true, false, true, false, false,
371+
]))),
322372
RowGroupAccess::Selection(select_two_rows()),
323373
])),
324374
row_selection: None,

‎datafusion/datasource-parquet/src/access_plan.rs‎

Lines changed: 191 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
// specific language governing permissions and limitations
1616
// under the License.
1717

18+
use arrow::array::BooleanBufferBuilder;
1819
use arrow::datatypes::Schema;
1920
use datafusion_common::{Result, assert_eq_or_internal_err, exec_err};
2021
use datafusion_physical_expr::expressions::Column;
@@ -297,6 +298,54 @@ impl ParquetAccessPlan {
297298
pub fn try_new_from_overall_row_selection(
298299
selection: RowSelection,
299300
row_group_meta_data: &[RowGroupMetaData],
301+
) -> Result<Self> {
302+
if selection.as_mask().is_some() {
303+
Self::try_new_from_overall_row_selection_mask(selection, row_group_meta_data)
304+
} else {
305+
Self::try_new_from_overall_row_selection_selectors(
306+
selection,
307+
row_group_meta_data,
308+
)
309+
}
310+
}
311+
312+
fn try_new_from_overall_row_selection_mask(
313+
mut selection: RowSelection,
314+
row_group_meta_data: &[RowGroupMetaData],
315+
) -> Result<Self> {
316+
let selection_rows = selection.total_row_count();
317+
let file_rows = row_group_meta_data
318+
.iter()
319+
.map(|rg| rg.num_rows() as usize)
320+
.sum::<usize>();
321+
if selection_rows != file_rows {
322+
return exec_err!(
323+
"Invalid Parquet RowSelection. File has {file_rows} rows, \
324+
but selection specifies {selection_rows} rows."
325+
);
326+
}
327+
328+
// For masks, split_off uses bitmap slices without materializing
329+
// selectors. Use row_count() to cache each partial group's count
330+
// for later preparation.
331+
let row_groups = row_group_meta_data
332+
.iter()
333+
.map(|rg| {
334+
let row_count = rg.num_rows() as usize;
335+
let group_selection = selection.split_off(row_count);
336+
match group_selection.row_count() {
337+
0 => RowGroupAccess::Skip,
338+
selected if selected == row_count => RowGroupAccess::Scan,
339+
_ => RowGroupAccess::Selection(group_selection),
340+
}
341+
})
342+
.collect();
343+
Ok(Self::new(row_groups))
344+
}
345+
346+
fn try_new_from_overall_row_selection_selectors(
347+
selection: RowSelection,
348+
row_group_meta_data: &[RowGroupMetaData],
300349
) -> Result<Self> {
301350
// Keep this as a single pass over the selector stream rather than
302351
// repeatedly calling `RowSelection::split_off` per row group. The
@@ -392,6 +441,22 @@ impl ParquetAccessPlan {
392441
RowGroupAccess::Skip => RowGroupAccess::Skip,
393442
RowGroupAccess::Scan => RowGroupAccess::Selection(selection),
394443
RowGroupAccess::Selection(existing_selection) => {
444+
// Parquet preserves bitmap backing only when both operands
445+
// are masks. Promote selector-backed page pruning to retain
446+
// an external index's bitmap and use a bitwise intersection.
447+
// Revisit this conversion once Parquet optimizes mixed-backed
448+
// intersections: https://github.com/apache/arrow-rs/issues/10423
449+
let selection = if existing_selection.as_mask().is_some()
450+
&& selection.as_mask().is_none()
451+
{
452+
let mut mask = BooleanBufferBuilder::new(selection.total_row_count());
453+
for selector in selection.iter() {
454+
mask.append_n(selector.row_count, !selector.skip);
455+
}
456+
RowSelection::from(mask.finish())
457+
} else {
458+
selection
459+
};
395460
RowGroupAccess::Selection(existing_selection.intersection(&selection))
396461
}
397462
}
@@ -838,6 +903,7 @@ impl PreparedAccessPlan {
838903
#[cfg(test)]
839904
mod test {
840905
use super::*;
906+
use arrow::buffer::BooleanBuffer;
841907
use datafusion_common::assert_contains;
842908
use parquet::basic::LogicalType;
843909
use parquet::file::metadata::ColumnChunkMetaData;
@@ -862,7 +928,7 @@ mod test {
862928
// Check both input representations retain the same conversion behavior.
863929
let selectors =
864930
RowSelection::from(vec![RowSelector::skip(10), RowSelector::select(20)]);
865-
let bitmap = RowSelection::from(arrow::buffer::BooleanBuffer::from(
931+
let bitmap = RowSelection::from(BooleanBuffer::from(
866932
(0..30).map(|i| i >= 10).collect::<Vec<_>>(),
867933
));
868934
for selection in [selectors, bitmap] {
@@ -986,6 +1052,110 @@ mod test {
9861052
}
9871053
}
9881054

1055+
#[test]
1056+
fn test_scan_selection_preserves_mask_backing() {
1057+
let mask = BooleanBuffer::from(vec![
1058+
true, true, false, false, true, true, false, false, true, true,
1059+
]);
1060+
let selectors =
1061+
RowSelection::from(vec![RowSelector::select(5), RowSelector::skip(5)]);
1062+
// Both selector-backed page pruning and bitmap intersections retain
1063+
// the existing mask, including when no rows survive.
1064+
for incoming in [
1065+
selectors,
1066+
RowSelection::from(BooleanBuffer::from(vec![
1067+
true, true, true, true, true, false, false, false, false, false,
1068+
])),
1069+
RowSelection::from(vec![RowSelector::skip(10)]),
1070+
] {
1071+
let empty = incoming.row_count() == 0;
1072+
let mut plan = ParquetAccessPlan::new(vec![RowGroupAccess::Selection(
1073+
RowSelection::from(mask.clone()),
1074+
)]);
1075+
plan.scan_selection(0, incoming);
1076+
let RowGroupAccess::Selection(selection) = &plan.inner()[0] else {
1077+
panic!("expected selection");
1078+
};
1079+
let expected = if empty {
1080+
BooleanBuffer::new_unset(10)
1081+
} else {
1082+
BooleanBuffer::from(vec![
1083+
true, true, false, false, true, false, false, false, false, false,
1084+
])
1085+
};
1086+
assert_eq!(selection.as_mask(), Some(&expected));
1087+
}
1088+
}
1089+
1090+
#[test]
1091+
fn test_scan_selection_preserves_selector_backing() {
1092+
let mut plan = ParquetAccessPlan::new(vec![RowGroupAccess::Selection(
1093+
RowSelection::from(vec![RowSelector::select(6), RowSelector::skip(4)]),
1094+
)]);
1095+
plan.scan_selection(
1096+
0,
1097+
RowSelection::from(vec![
1098+
RowSelector::skip(2),
1099+
RowSelector::select(5),
1100+
RowSelector::skip(3),
1101+
]),
1102+
);
1103+
let RowGroupAccess::Selection(selection) = &plan.inner()[0] else {
1104+
panic!("expected selection");
1105+
};
1106+
assert!(selection.as_mask().is_none());
1107+
assert_eq!(
1108+
selection,
1109+
&RowSelection::from(vec![
1110+
RowSelector::skip(2),
1111+
RowSelector::select(4),
1112+
RowSelector::skip(4),
1113+
])
1114+
);
1115+
}
1116+
1117+
#[test]
1118+
fn test_new_from_overall_mask_preserves_bitmap_backing() {
1119+
// Include all-selected, all-skipped, and fragmented groups. Start at
1120+
// a non-byte-aligned offset to exercise slicing an existing bitmap.
1121+
let mut bits = vec![false; 3];
1122+
bits.extend(vec![true; 10]);
1123+
bits.extend(vec![false; 20]);
1124+
bits.extend((0..30).map(|i| i % 2 == 1));
1125+
bits.extend(vec![true; 40]);
1126+
let mask = BooleanBuffer::from(bits).slice(3, 100);
1127+
for cache_row_count in [false, true] {
1128+
let selection = RowSelection::from(mask.clone());
1129+
if cache_row_count {
1130+
// Exercise split_off's propagation of an already cached count.
1131+
assert_eq!(selection.row_count(), 65);
1132+
}
1133+
let plan = ParquetAccessPlan::try_new_from_overall_row_selection(
1134+
selection,
1135+
&ROW_GROUP_METADATA,
1136+
)
1137+
.unwrap();
1138+
assert_eq!(plan.inner()[0], RowGroupAccess::Scan);
1139+
assert_eq!(plan.inner()[1], RowGroupAccess::Skip);
1140+
assert_eq!(plan.inner()[3], RowGroupAccess::Scan);
1141+
let RowGroupAccess::Selection(selection) = &plan.inner()[2] else {
1142+
panic!("expected selection");
1143+
};
1144+
assert_eq!(selection.row_count(), 15);
1145+
let group_mask = mask.slice(30, 30);
1146+
assert!(selection.as_mask().unwrap().ptr_eq(&group_mask));
1147+
1148+
// The local selection must also survive preparation and reversal.
1149+
let prepared = plan.prepare(&ROW_GROUP_METADATA).unwrap().reverse();
1150+
assert_eq!(prepared.row_group_indexes(), vec![3, 2, 0]);
1151+
assert!(prepared.row_groups[0].selection.selection().is_none());
1152+
assert!(prepared.row_groups[2].selection.selection().is_none());
1153+
let selection = prepared.row_groups[1].selection.selection().unwrap();
1154+
assert_eq!(selection.row_count(), 15);
1155+
assert!(selection.as_mask().unwrap().ptr_eq(&group_mask));
1156+
}
1157+
}
1158+
9891159
#[test]
9901160
fn test_new_from_overall_row_selection() {
9911161
let row_selection = RowSelection::from(vec![
@@ -1022,19 +1192,26 @@ mod test {
10221192

10231193
#[test]
10241194
fn test_new_from_overall_row_selection_invalid_row_count() {
1025-
let row_selection = RowSelection::from(vec![RowSelector::select(99)]);
1026-
1027-
let err = ParquetAccessPlan::try_new_from_overall_row_selection(
1028-
row_selection,
1029-
&ROW_GROUP_METADATA,
1030-
)
1031-
.unwrap_err()
1032-
.to_string();
1033-
1034-
assert_contains!(
1035-
err,
1036-
"Invalid Parquet RowSelection. File has 100 rows, but selection specifies 99 rows"
1037-
);
1195+
for selection_rows in [99, 101] {
1196+
for selection in [
1197+
RowSelection::from(vec![RowSelector::select(selection_rows)]),
1198+
RowSelection::from(BooleanBuffer::new_set(selection_rows)),
1199+
] {
1200+
let err = ParquetAccessPlan::try_new_from_overall_row_selection(
1201+
selection,
1202+
&ROW_GROUP_METADATA,
1203+
)
1204+
.unwrap_err()
1205+
.to_string();
1206+
assert_contains!(
1207+
err,
1208+
format!(
1209+
"Invalid Parquet RowSelection. File has 100 rows, \
1210+
but selection specifies {selection_rows} rows"
1211+
)
1212+
);
1213+
}
1214+
}
10381215
}
10391216

10401217
#[test]

0 commit comments

Comments
 (0)