parquet/arrow/arrow_reader/selection/
cursor.rs1use super::boolean::boolean_mask_from_selectors;
26use super::{RowSelection, RowSelector};
27use crate::errors::ParquetError;
28use arrow_array::BooleanArray;
29use arrow_buffer::BooleanBuffer;
30use std::collections::VecDeque;
31use std::ops::Range;
32use std::sync::Arc;
33
34#[derive(Clone, Copy, Debug, Eq, PartialEq)]
36pub enum RowSelectionPolicy {
37 Selectors,
39 Mask,
41 Auto {
43 threshold: usize,
45 },
46}
47
48impl Default for RowSelectionPolicy {
49 fn default() -> Self {
50 Self::Auto { threshold: 32 }
51 }
52}
53
54#[derive(Clone, Copy, Debug, Eq, PartialEq)]
59pub(crate) enum RowSelectionStrategy {
60 Selectors,
62 Mask,
64}
65
66#[derive(Debug)]
72pub enum RowSelectionCursor {
73 All,
75 Mask(MaskCursor),
77 Selectors(SelectorsCursor),
79}
80
81impl RowSelectionCursor {
82 pub(crate) fn new_mask_from_selectors(
84 selectors: Vec<RowSelector>,
85 loaded_row_ranges: Option<Arc<LoadedRowRanges>>,
86 ) -> Self {
87 debug_assert!(
88 selectors
89 .last()
90 .map(|selector| !selector.skip)
91 .unwrap_or(true),
92 "Mask selectors must not end with a skip"
93 );
94 Self::Mask(MaskCursor {
95 mask: boolean_mask_from_selectors(&selectors),
96 position: 0,
97 loaded_row_ranges,
98 })
99 }
100
101 pub(crate) fn new_mask_from_buffer(
103 mask: BooleanBuffer,
104 loaded_row_ranges: Option<Arc<LoadedRowRanges>>,
105 ) -> Self {
106 debug_assert!(
107 mask.is_empty() || mask.value(mask.len() - 1),
108 "Mask selections must not end with a skip"
109 );
110 Self::Mask(MaskCursor {
111 mask,
112 position: 0,
113 loaded_row_ranges,
114 })
115 }
116
117 pub(crate) fn new_selectors(selectors: Vec<RowSelector>) -> Self {
119 Self::Selectors(SelectorsCursor {
120 selectors: selectors.into(),
121 position: 0,
122 })
123 }
124
125 pub(crate) fn new_all() -> Self {
127 Self::All
128 }
129}
130
131#[derive(Debug)]
136pub struct SelectorsCursor {
137 selectors: VecDeque<RowSelector>,
138 position: usize,
140}
141
142impl SelectorsCursor {
143 pub fn is_empty(&self) -> bool {
145 self.selectors.is_empty()
146 }
147
148 pub(crate) fn next_selector(&mut self) -> RowSelector {
150 let selector = self.selectors.pop_front().unwrap();
151 self.position += selector.row_count;
152 selector
153 }
154
155 pub(crate) fn return_selector(&mut self, selector: RowSelector) {
157 self.position = self.position.saturating_sub(selector.row_count);
158 self.selectors.push_front(selector);
159 }
160}
161
162#[derive(Debug)]
187pub struct MaskCursor {
188 mask: BooleanBuffer,
189 position: usize,
191 loaded_row_ranges: Option<Arc<LoadedRowRanges>>,
193}
194
195impl MaskCursor {
196 pub fn is_empty(&self) -> bool {
198 self.position >= self.mask.len()
199 }
200
201 pub fn next_mask_chunk(&mut self, batch_size: usize) -> Option<MaskChunk> {
203 if self.is_empty() {
204 return None;
205 }
206
207 Some(self.next_mask_chunk_non_empty(batch_size))
208 }
209
210 fn next_mask_chunk_non_empty(&mut self, batch_size: usize) -> MaskChunk {
212 debug_assert!(!self.is_empty());
213
214 let (initial_skip, chunk_rows, selected_rows, mask_start, end_position) = {
215 let mask = &self.mask;
216 let start_position = self.position;
217 let mut cursor = start_position;
218 let mut initial_skip = 0;
219
220 while cursor < mask.len() && !mask.value(cursor) {
221 initial_skip += 1;
222 cursor += 1;
223 }
224 debug_assert!(
225 cursor < mask.len(),
226 "ReadPlan must remove trailing skips from Mask selections"
227 );
228
229 let mask_start = cursor;
230 let mut chunk_rows = 0;
231 let mut selected_rows = 0;
232
233 while cursor < mask.len() && selected_rows < batch_size {
237 chunk_rows += 1;
238 if mask.value(cursor) {
239 selected_rows += 1;
240 }
241 cursor += 1;
242 }
243
244 (initial_skip, chunk_rows, selected_rows, mask_start, cursor)
245 };
246
247 self.position = end_position;
248
249 MaskChunk {
250 initial_skip,
251 chunk_rows,
252 selected_rows,
253 mask_start,
254 }
255 }
256
257 pub(crate) fn next_chunk(&mut self, batch_size: usize) -> Result<MaskChunk, ParquetError> {
266 debug_assert!(batch_size > 0);
267 debug_assert!(!self.is_empty());
268
269 if self.loaded_row_ranges.is_none() {
270 return Ok(self.next_mask_chunk_non_empty(batch_size));
271 }
272
273 let start_position = self.position;
274 let mut cursor = start_position;
275 while cursor < self.mask.len() && !self.mask.value(cursor) {
276 cursor += 1;
277 }
278
279 if cursor == self.mask.len() {
280 return Err(ParquetError::General(
281 "Internal Error: Mask cursor reached the end without finding a selected row; \
282 ReadPlan must remove trailing skips"
283 .to_string(),
284 ));
285 }
286
287 let loaded_range_end = self
288 .loaded_row_ranges
289 .as_ref()
290 .and_then(|ranges| ranges.end_containing(cursor))
291 .ok_or_else(|| {
292 ParquetError::General(format!(
293 "Internal Error: selected row {cursor} has no loaded page range"
294 ))
295 })?;
296
297 let mask_start = cursor;
298 let mut selected_rows = 0;
299 let mut chunk_end = cursor;
300 while cursor < loaded_range_end && cursor < self.mask.len() && selected_rows < batch_size {
301 if self.mask.value(cursor) {
302 selected_rows += 1;
303 chunk_end = cursor + 1;
304 }
305 cursor += 1;
306 }
307
308 self.position = chunk_end;
309 Ok(MaskChunk {
310 initial_skip: mask_start - start_position,
311 chunk_rows: chunk_end - mask_start,
312 selected_rows,
313 mask_start,
314 })
315 }
316
317 pub fn mask_values_for(&self, chunk: &MaskChunk) -> Result<BooleanArray, ParquetError> {
319 if chunk.mask_start.saturating_add(chunk.chunk_rows) > self.mask.len() {
320 return Err(ParquetError::General(
321 "Internal Error: MaskChunk exceeds mask length".to_string(),
322 ));
323 }
324 Ok(BooleanArray::from(
325 self.mask.slice(chunk.mask_start, chunk.chunk_rows),
326 ))
327 }
328}
329
330#[derive(Debug)]
332pub struct MaskChunk {
333 pub initial_skip: usize,
335 pub chunk_rows: usize,
337 pub selected_rows: usize,
339 pub mask_start: usize,
341}
342
343#[derive(Clone, Debug)]
345pub(crate) struct LoadedRowRanges(Vec<Range<usize>>);
346
347impl LoadedRowRanges {
348 pub(crate) fn from_selection(selection: RowSelection) -> Self {
349 let selectors: Vec<RowSelector> = selection.into();
350 let mut position = 0;
351 let ranges = selectors
352 .into_iter()
353 .filter_map(|selector| {
354 let start = position;
355 position += selector.row_count;
356 (!selector.skip).then_some(start..position)
357 })
358 .collect();
359 Self(ranges)
360 }
361
362 fn end_containing(&self, row: usize) -> Option<usize> {
363 let idx = self.0.partition_point(|range| range.end <= row);
364 self.0
365 .get(idx)
366 .filter(|range| range.start <= row)
367 .map(|range| range.end)
368 }
369
370 #[cfg(test)]
371 pub(crate) fn ranges(&self) -> &[Range<usize>] {
372 &self.0
373 }
374}
375
376#[cfg(test)]
377mod tests {
378 use super::*;
379
380 #[test]
381 fn test_loaded_mask_chunk_stops_at_trimmed_mask_end() {
382 let loaded = LoadedRowRanges::from_selection(RowSelection::from_consecutive_ranges(
383 std::iter::once(0..5),
384 10,
385 ));
386 let RowSelectionCursor::Mask(mut cursor) = RowSelectionCursor::new_mask_from_selectors(
387 vec![RowSelector::select(1)],
388 Some(loaded.into()),
389 ) else {
390 unreachable!()
391 };
392
393 let chunk = cursor.next_chunk(10).unwrap();
394 assert_eq!(chunk.chunk_rows, 1);
395 assert!(cursor.is_empty());
396 }
397
398 #[test]
399 fn test_next_mask_chunk_until_cursor_is_empty() {
400 let RowSelectionCursor::Mask(mut cursor) = RowSelectionCursor::new_mask_from_selectors(
401 vec![
402 RowSelector::skip(2),
403 RowSelector::select(2),
404 RowSelector::skip(1),
405 RowSelector::select(1),
406 ],
407 None,
408 ) else {
409 unreachable!()
410 };
411
412 let first = cursor.next_mask_chunk(2).unwrap();
413 assert_eq!(first.initial_skip, 2);
414 assert_eq!(first.chunk_rows, 2);
415 assert_eq!(first.selected_rows, 2);
416
417 let second = cursor.next_mask_chunk(2).unwrap();
418 assert_eq!(second.initial_skip, 1);
419 assert_eq!(second.chunk_rows, 1);
420 assert_eq!(second.selected_rows, 1);
421
422 assert!(cursor.next_mask_chunk(2).is_none());
423 }
424}