alamb commented on code in PR #10859:
URL: https://github.com/apache/arrow-rs/pull/10859#discussion_r4123380960
##########
parquet/src/arrow/arrow_reader/selection/cursor.rs:
##########
@@ -51,6 +50,38 @@ impl Default for RowSelectionPolicy {
}
}
+impl RowSelectionPolicy {
+ /// Resolve this policy for a selection without changing its backing.
Review Comment:
This API is becoming quite nice 👌
##########
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 just tested it, and the change is a bit too big. and I don’t find much
difference in performance.
Yeah, I agree this would be a larger change
It would only matter if we have predicates that could natively create
RowSelections backed by `Vec<...>` -- but by definition no such predicates
exist at the moment (or at least they aren't wired in)
##########
parquet/src/arrow/arrow_reader/read_plan.rs:
##########
@@ -157,17 +157,7 @@ impl ReadPlanBuilder {
///
/// Guarantees to return either `Selectors` or `Mask`, never `Auto`.
pub(crate) fn resolve_selection_strategy(&self) -> RowSelectionStrategy {
- match self.row_selection_policy {
- RowSelectionPolicy::Selectors => RowSelectionStrategy::Selectors,
- RowSelectionPolicy::Mask => RowSelectionStrategy::Mask,
- RowSelectionPolicy::Auto { threshold, .. } => {
- let Some(selection) = self.selection.as_ref() else {
- return RowSelectionStrategy::Selectors;
- };
-
- selection.auto_selection_strategy(threshold)
- }
- }
+ self.row_selection_policy.resolve(self.selection.as_ref())
Review Comment:
As a follow on PR we could probably just remove this level of indirection as
call `self.row_selection_policy.resolve(self.selection.as_ref())` directly
##########
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:
This is looking really nice now -- thank you
--
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]