haohuaijin commented on code in PR #10141:
URL: https://github.com/apache/arrow-rs/pull/10141#discussion_r3643176544
##########
parquet/src/arrow/arrow_reader/selection/mod.rs:
##########
@@ -135,12 +149,297 @@ impl RowSelector {
/// * Consecutive [`RowSelector`]s alternate skipping or selecting rows
///
/// [`PageIndex`]: crate::file::page_index::column_index::ColumnIndexMetaData
-#[derive(Debug, Clone, Default, Eq, PartialEq)]
+#[derive(Default, Clone)]
pub struct RowSelection {
- selectors: Vec<RowSelector>,
+ inner: RowSelectionInner,
+}
+
+/// Internal storage for [`RowSelection`].
+#[derive(Debug, Clone)]
+pub(crate) enum RowSelectionInner {
+ Selectors(Vec<RowSelector>),
+ Mask(Box<MaskSelection>),
+}
+
+impl Default for RowSelectionInner {
+ fn default() -> Self {
+ Self::Selectors(Vec::new())
+ }
+}
+
+impl std::fmt::Debug for RowSelection {
+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ match &self.inner {
+ RowSelectionInner::Selectors(s) => f
+ .debug_struct("RowSelection")
+ .field("selectors", s)
+ .finish(),
+ RowSelectionInner::Mask(m) => f
+ .debug_struct("RowSelection")
+ .field("mask_len", &m.mask().len())
+ .finish_non_exhaustive(),
+ }
+ }
+}
+
+impl PartialEq for RowSelection {
+ fn eq(&self, other: &Self) -> bool {
+ match (&self.inner, &other.inner) {
+ (RowSelectionInner::Selectors(a), RowSelectionInner::Selectors(b))
=> a == b,
+ (RowSelectionInner::Mask(a), RowSelectionInner::Mask(b)) =>
a.mask() == b.mask(),
+ (RowSelectionInner::Mask(mask),
RowSelectionInner::Selectors(selectors))
+ | (RowSelectionInner::Selectors(selectors),
RowSelectionInner::Mask(mask)) => {
+ if selectors
+ .iter()
+ .try_fold(0usize, |acc, selector|
acc.checked_add(selector.row_count))
+ != Some(mask.mask().len())
+ {
+ return false;
+ }
+
+ let mut slices = mask.mask().set_slices().peekable();
+ let mut cursor = 0usize;
+
+ for selector in selectors {
+ let end = cursor + selector.row_count;
+
+ if selector.skip {
+ if slices.peek().is_some_and(|(start, _)| *start <
end) {
+ return false;
+ }
+ } else {
+ match slices.next() {
+ Some((start, slice_end)) if start == cursor &&
slice_end == end => {}
+ _ => return false,
+ }
+ }
+
+ cursor = end;
+ }
+
+ slices.next().is_none()
+ }
+ }
+ }
+}
+
+impl Eq for RowSelection {}
+
+/// Borrowed iterator over the [`RowSelector`]s of a [`RowSelection`].
+#[derive(Debug)]
+pub struct RowSelectionIter<'a>(std::slice::Iter<'a, RowSelector>);
+
+impl<'a> Iterator for RowSelectionIter<'a> {
+ type Item = &'a RowSelector;
+
+ #[inline]
+ fn next(&mut self) -> Option<Self::Item> {
+ self.0.next()
+ }
+}
+
+#[inline]
+fn scan_ranges_from_selectors<I>(selectors: I, page_locations:
&[PageLocation]) -> Vec<Range<u64>>
+where
+ I: IntoIterator<Item = RowSelector>,
+{
+ let mut ranges: Vec<Range<u64>> = vec![];
+ let mut row_offset = 0;
+
+ let mut pages = page_locations.iter().peekable();
+ let mut selectors = selectors.into_iter();
+ let mut current_selector = selectors.next();
+ let mut current_page = pages.next();
+
+ let mut current_page_included = false;
+
+ while let Some((selector, page)) =
current_selector.as_mut().zip(current_page) {
+ if !(selector.skip || current_page_included) {
+ let start = page.offset as u64;
+ let end = start + page.compressed_page_size as u64;
+ ranges.push(start..end);
+ current_page_included = true;
+ }
+
+ if let Some(next_page) = pages.peek() {
+ if row_offset + selector.row_count > next_page.first_row_index as
usize {
+ let remaining_in_page = next_page.first_row_index as usize -
row_offset;
+ selector.row_count -= remaining_in_page;
+ row_offset += remaining_in_page;
+ current_page = pages.next();
+ current_page_included = false;
+
+ continue;
+ } else {
+ if row_offset + selector.row_count ==
next_page.first_row_index as usize {
+ current_page = pages.next();
+ current_page_included = false;
+ }
+ row_offset += selector.row_count;
+ current_selector = selectors.next();
+ }
+ } else {
+ if !(selector.skip || current_page_included) {
+ let start = page.offset as u64;
+ let end = start + page.compressed_page_size as u64;
+ ranges.push(start..end);
+ }
+ current_selector = selectors.next()
+ }
+ }
+
+ ranges
+}
+
+#[inline]
+fn expand_to_batch_boundaries_from_selectors<I>(
+ selectors: I,
+ batch_size: usize,
+ total_rows: usize,
+) -> RowSelection
+where
+ I: IntoIterator<Item = RowSelector>,
+{
+ let mut expanded_ranges = Vec::new();
+ let mut row_offset = 0;
+
+ for selector in selectors {
+ if selector.skip {
+ row_offset += selector.row_count;
+ } else {
+ let start = row_offset;
+ let end = row_offset + selector.row_count;
+
+ // Expand start to batch boundary
+ let expanded_start = (start / batch_size) * batch_size;
+ // Expand end to batch boundary
+ let expanded_end = end.div_ceil(batch_size) * batch_size;
+ let expanded_end = expanded_end.min(total_rows);
+
+ expanded_ranges.push(expanded_start..expanded_end);
+ row_offset += selector.row_count;
+ }
+ }
+
+ // Sort ranges by start position
+ expanded_ranges.sort_by_key(|range| range.start);
+
+ // Merge overlapping or consecutive ranges
+ let mut merged_ranges: Vec<Range<usize>> = Vec::new();
+ for range in expanded_ranges {
+ if let Some(last) = merged_ranges.last_mut() {
+ if range.start <= last.end {
+ // Overlapping or consecutive - merge them
+ last.end = last.end.max(range.end);
+ } else {
+ // No overlap - add new range
+ merged_ranges.push(range);
+ }
+ } else {
+ // First range
+ merged_ranges.push(range);
+ }
+ }
+
+ RowSelection::from_consecutive_ranges(merged_ranges.into_iter(),
total_rows)
}
impl RowSelection {
+ fn from_selectors(selectors: Vec<RowSelector>) -> Self {
+ Self {
+ inner: RowSelectionInner::Selectors(selectors),
+ }
+ }
+
+ /// Create a [`RowSelection`] from a packed [`BooleanBuffer`].
+ ///
+ /// Each set bit selects a row, each unset bit skips one. Unlike
+ /// [`Self::from_filters`], the bitmap is kept as-is rather than
+ /// eagerly run-length-encoded. [`Self::iter`] materializes and caches the
+ /// RLE form on first use; use [`MaskRunIter`] to stream the RLE form
+ /// directly from the bitmap.
+ pub fn from_boolean_buffer(mask: BooleanBuffer) -> Self {
+ Self {
+ inner: RowSelectionInner::Mask(Box::new(MaskSelection::new(mask))),
+ }
+ }
+
+ /// Returns the underlying mask if this selection is mask-backed.
+ ///
+ /// Public so that engines composing selections (e.g. DataFusion's
+ /// `ParquetAccessPlan::into_overall_row_selection`) can concatenate
+ /// mask-backed selections without materialising the RLE form.
+ pub fn as_mask(&self) -> Option<&BooleanBuffer> {
+ match &self.inner {
+ RowSelectionInner::Mask(m) => Some(m.mask()),
+ _ => None,
+ }
+ }
+
+ /// Consume the selection and return its internal storage.
+ pub(crate) fn into_inner(self) -> RowSelectionInner {
+ self.inner
+ }
+
+ /// Choose the automatic materialisation strategy without converting
between
+ /// selector and mask backing.
+ #[inline]
+ pub(crate) fn auto_selection_strategy(&self, threshold: usize) ->
RowSelectionStrategy {
+ let (total_rows, effective_count) = match &self.inner {
+ RowSelectionInner::Selectors(selectors) => {
+ selectors.iter().fold((0usize, 0usize), |(rows, count), s| {
+ if s.row_count > 0 {
+ (rows + s.row_count, count + 1)
+ } else {
+ (rows, count)
+ }
+ })
+ }
+ RowSelectionInner::Mask(mask) => {
+ let mask = mask.mask();
+ let total_rows = mask.len();
+ (total_rows, mask_run_count(mask))
+ }
+ };
+
+ if effective_count == 0 {
+ return RowSelectionStrategy::Mask;
+ }
+
+ if total_rows < effective_count.saturating_mul(threshold) {
+ RowSelectionStrategy::Mask
+ } else {
+ RowSelectionStrategy::Selectors
+ }
+ }
+
+ #[cfg(test)]
+ fn selectors(&self) -> Vec<RowSelector> {
+ self.iter().copied().collect()
+ }
+
+ fn into_selectors_vec(self) -> Vec<RowSelector> {
Review Comment:
track in https://github.com/apache/arrow-rs/issues/10422
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]