alamb commented on code in PR #10859:
URL: https://github.com/apache/arrow-rs/pull/10859#discussion_r4104718800
##########
parquet/src/arrow/arrow_reader/filter.rs:
##########
@@ -153,7 +158,10 @@ where
/// This design has a couple of implications:
///
/// * [`RowFilter`] can be used to skip entire pages, and thus IO, in addition
to CPU decode overheads
-/// * Columns may be decoded multiple times if they appear in multiple
[`ProjectionMask`]
+/// * Columns may be decoded multiple times if they appear in multiple
[`ProjectionMask`].
Review Comment:
👍
##########
parquet/src/arrow/arrow_reader/filter.rs:
##########
@@ -198,4 +206,540 @@ impl RowFilter {
pub fn into_predicates(self) -> Vec<Box<dyn ArrowPredicate>> {
self.predicates
}
+
+ /// Fuse consecutive predicates on the same single top-level, non-repeated
leaf.
+ /// This avoids repeated decoding or predicate-cache replay of that column.
+ pub(crate) fn fuse_same_projection(
+ self,
+ parquet_schema: &SchemaDescriptor,
+ row_selection_policy: RowSelectionPolicy,
+ ) -> Self {
+ let mut predicates: Vec<Box<dyn ArrowPredicate>> =
+ Vec::with_capacity(self.predicates.len());
+ let mut group: Vec<Box<dyn ArrowPredicate>> = Vec::new();
+ let mut flush_group = |group: &mut Vec<Box<dyn ArrowPredicate>>| {
+ if group.len() > 1 && can_fuse_projection(group[0].projection(),
parquet_schema) {
+ let group = std::mem::take(group);
+ predicates.push(Box::new(FusedPredicate::new(group,
row_selection_policy)));
+ } else {
+ predicates.append(group);
+ }
+ };
+
+ for predicate in self.predicates {
+ if group
+ .last()
+ .is_some_and(|last| last.projection() !=
predicate.projection())
+ {
+ flush_group(&mut group);
+ }
+ group.push(predicate);
+ }
+ flush_group(&mut group);
+
+ Self { predicates }
+ }
+}
+
+/// Restrict fusion to one top-level, non-repeated leaf to limit compaction
costs.
+fn can_fuse_projection(projection: &ProjectionMask, parquet_schema:
&SchemaDescriptor) -> bool {
+ let mut leaf_indices =
+ (0..parquet_schema.num_columns()).filter(|idx|
projection.leaf_included(*idx));
+ let Some(leaf_idx) = leaf_indices.next() else {
+ return false;
+ };
+ if leaf_indices.next().is_some() {
+ return false;
+ }
+
+ let column = parquet_schema.column(leaf_idx);
+ column.path().parts().len() == 1 && column.max_rep_level() == 0
+}
+
+/// Evaluate same-projection predicates in order on one decoded batch.
+/// Later predicates see only surviving rows, sliced when contiguous or
compacted
+/// otherwise. Selections follow the reader's row selection policy.
+struct FusedPredicate {
+ /// At least two predicates, all with the same projection.
+ predicates: Vec<Box<dyn ArrowPredicate>>,
+ row_selection_policy: RowSelectionPolicy,
+}
+
+impl FusedPredicate {
+ /// Create a fused predicate from at least two predicates that share one
+ /// projection.
+ fn new(
+ predicates: Vec<Box<dyn ArrowPredicate>>,
+ row_selection_policy: RowSelectionPolicy,
+ ) -> Self {
+ debug_assert!(predicates.len() > 1);
+ debug_assert!(
+ predicates
+ .windows(2)
+ .all(|pair| pair[0].projection() == pair[1].projection())
+ );
+ Self {
+ predicates,
+ row_selection_policy,
+ }
+ }
+}
+
+impl Debug for FusedPredicate {
+ fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
+ write!(
+ f,
+ "FusedPredicate {{ {} predicates }}",
+ self.predicates.len()
+ )
+ }
+}
+
+impl ArrowPredicate for FusedPredicate {
+ fn projection(&self) -> &ProjectionMask {
+ self.predicates[0].projection()
+ }
+
+ fn evaluate(&mut self, batch: RecordBatch) -> Result<BooleanArray,
ArrowError> {
+ let num_rows = batch.num_rows();
+ // Positions in the original batch; None until a predicate rejects
rows.
+ let mut selection: Option<RowSelection> = None;
+ let mut filtered_batch = batch;
+ let last_predicate_idx = self.predicates.len() - 1;
+
+ for (idx, predicate) in self.predicates.iter_mut().enumerate() {
Review Comment:
I don't fully understand the need for a row_selection_policy here. I was
expecting that when a predicate was fused, it would basically evaluate `pred1
AND pred2 AND ...`, combining the results with `and` -- or more likely the `&=`
operator that knows how to reuse the bitmaps
It is really strange to me to try and use the RowSelection API *inside* a
`ArrowPredicate` whose external API is in terms of BooleanArray (because then
we have to turn the result back to BooleanArray, to turn it right back into a
RowSeelction...)
Is the idea to avoid evaluating expensive subsequent filters when the first
filter in a fused filter was selective? I wonder if we should look into doing
this without RowSelection if possible 🤔
Or maybe we should change the ArrowPredicate API to return a RowSelection
(with a default implementation that returns the BooleanArray 🤔 )
##########
parquet/src/arrow/arrow_reader/filter.rs:
##########
@@ -198,4 +206,540 @@ impl RowFilter {
pub fn into_predicates(self) -> Vec<Box<dyn ArrowPredicate>> {
self.predicates
}
+
+ /// Fuse consecutive predicates on the same single top-level, non-repeated
leaf.
+ /// This avoids repeated decoding or predicate-cache replay of that column.
+ pub(crate) fn fuse_same_projection(
+ self,
+ parquet_schema: &SchemaDescriptor,
+ row_selection_policy: RowSelectionPolicy,
+ ) -> Self {
+ let mut predicates: Vec<Box<dyn ArrowPredicate>> =
+ Vec::with_capacity(self.predicates.len());
+ let mut group: Vec<Box<dyn ArrowPredicate>> = Vec::new();
+ let mut flush_group = |group: &mut Vec<Box<dyn ArrowPredicate>>| {
+ if group.len() > 1 && can_fuse_projection(group[0].projection(),
parquet_schema) {
+ let group = std::mem::take(group);
+ predicates.push(Box::new(FusedPredicate::new(group,
row_selection_policy)));
+ } else {
+ predicates.append(group);
+ }
+ };
+
+ for predicate in self.predicates {
+ if group
+ .last()
+ .is_some_and(|last| last.projection() !=
predicate.projection())
+ {
+ flush_group(&mut group);
+ }
+ group.push(predicate);
+ }
+ flush_group(&mut group);
+
+ Self { predicates }
+ }
+}
+
+/// Restrict fusion to one top-level, non-repeated leaf to limit compaction
costs.
+fn can_fuse_projection(projection: &ProjectionMask, parquet_schema:
&SchemaDescriptor) -> bool {
+ let mut leaf_indices =
+ (0..parquet_schema.num_columns()).filter(|idx|
projection.leaf_included(*idx));
+ let Some(leaf_idx) = leaf_indices.next() else {
+ return false;
+ };
+ if leaf_indices.next().is_some() {
+ return false;
+ }
+
+ let column = parquet_schema.column(leaf_idx);
+ column.path().parts().len() == 1 && column.max_rep_level() == 0
+}
+
+/// Evaluate same-projection predicates in order on one decoded batch.
+/// Later predicates see only surviving rows, sliced when contiguous or
compacted
+/// otherwise. Selections follow the reader's row selection policy.
+struct FusedPredicate {
+ /// At least two predicates, all with the same projection.
+ predicates: Vec<Box<dyn ArrowPredicate>>,
+ row_selection_policy: RowSelectionPolicy,
+}
+
+impl FusedPredicate {
+ /// Create a fused predicate from at least two predicates that share one
+ /// projection.
+ fn new(
+ predicates: Vec<Box<dyn ArrowPredicate>>,
+ row_selection_policy: RowSelectionPolicy,
+ ) -> Self {
+ debug_assert!(predicates.len() > 1);
+ debug_assert!(
+ predicates
+ .windows(2)
+ .all(|pair| pair[0].projection() == pair[1].projection())
+ );
+ Self {
+ predicates,
+ row_selection_policy,
+ }
+ }
+}
+
+impl Debug for FusedPredicate {
+ fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
+ write!(
+ f,
+ "FusedPredicate {{ {} predicates }}",
+ self.predicates.len()
+ )
+ }
+}
+
+impl ArrowPredicate for FusedPredicate {
+ fn projection(&self) -> &ProjectionMask {
+ self.predicates[0].projection()
+ }
+
+ fn evaluate(&mut self, batch: RecordBatch) -> Result<BooleanArray,
ArrowError> {
+ let num_rows = batch.num_rows();
+ // Positions in the original batch; None until a predicate rejects
rows.
+ let mut selection: Option<RowSelection> = None;
+ let mut filtered_batch = batch;
+ let last_predicate_idx = self.predicates.len() - 1;
+
+ for (idx, predicate) in self.predicates.iter_mut().enumerate() {
+ let filter = evaluate_predicate(predicate.as_mut(),
filtered_batch.clone())?;
+ // No mapping is needed if all preceding predicates accepted every
row.
+ if idx == last_predicate_idx && selection.is_none() {
+ return Ok(filter);
+ }
+ let true_count = filter.true_count();
+ if true_count == 0 {
+ return
Ok(BooleanArray::new(BooleanBuffer::new_unset(num_rows), None));
+ }
+ if true_count == filter.len() {
+ continue;
+ }
+
+ let predicate_selection =
RowSelection::from_boolean_buffer(filter.values().clone());
+ selection = Some(match selection.take() {
+ // Only the accumulated selection drives the composition
+ // algorithm, so adapt it once, right before it is used.
+ Some(prev) => adapt_fusion_selection(prev,
self.row_selection_policy)
+ .and_then(&predicate_selection),
+ None => predicate_selection,
+ });
+ if idx != last_predicate_idx {
+ filtered_batch = narrow_batch(&filtered_batch, &filter,
true_count)?;
+ }
+ }
+
+ let mask = match selection {
+ Some(selection) => selection.into_boolean_buffer(),
+ None => BooleanBuffer::new_set(num_rows),
+ };
+ debug_assert_eq!(mask.len(), num_rows);
+ Ok(BooleanArray::new(mask, None))
+ }
+}
+
+/// Apply the reader's row selection policy to a fusion selection.
+fn adapt_fusion_selection(selection: RowSelection, policy: RowSelectionPolicy)
-> RowSelection {
Review Comment:
This feels very similar to
[ReadPlanBuilder::resolve_selection_strategy](https://github.com/apache/arrow-rs/blob/d506e145f42169e9c188b627f341da499bf09859/parquet/src/arrow/arrow_reader/read_plan.rs#L159-L171)
and
[build_cursor](https://github.com/apache/arrow-rs/blob/d506e145f42169e9c188b627f341da499bf09859/parquet/src/arrow/arrow_reader/read_plan.rs#L320-L339)
would it be possible to re-use that logic here? For example by making
`RowSelectionStrategy::build_cursor(selection)` or adding some way to apply a
row selection to an array 🤔
##########
parquet/src/arrow/arrow_reader/filter.rs:
##########
@@ -198,4 +206,540 @@ impl RowFilter {
pub fn into_predicates(self) -> Vec<Box<dyn ArrowPredicate>> {
self.predicates
}
+
+ /// Fuse consecutive predicates on the same single top-level, non-repeated
leaf.
+ /// This avoids repeated decoding or predicate-cache replay of that column.
+ pub(crate) fn fuse_same_projection(
+ self,
+ parquet_schema: &SchemaDescriptor,
+ row_selection_policy: RowSelectionPolicy,
+ ) -> Self {
+ let mut predicates: Vec<Box<dyn ArrowPredicate>> =
Review Comment:
we could potentially scan through the predicates first and find if any
needed to be fused before allocating, though I am not sure that matters
--
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]