From a3047edee50b5a94cab6752bfb9905ea76813c0b Mon Sep 17 00:00:00 2001 From: Qi Zhu <821684824@qq.com> Date: Fri, 17 Jul 2026 15:16:17 +0800 Subject: [PATCH] perf(parquet): skip RowFilter on statically fully-matched row groups MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Builds on the `peek_next_row_group` API landed in apache/arrow-rs#10158. When Parquet stats prove that every row of a row group already satisfies the pushdown predicate (`fully_matched`), running the per-row `RowFilter` inside that row group is pure overhead — every row passes anyway. This change installs an *empty* `RowFilter` on fully-matched runs and only pays the row-level machinery on RGs that still need filtering. ## What changes * `access_plan.rs` — `ParquetAccessPlan` now tracks per-RG `fully_matched` state, produced during static pruning. * `metrics.rs` — new `row_filter_skipped_fully_matched` counter to observe how often the toggle fires. * `push_decoder.rs` — `RgPlanEntry` carries the per-RG toggle; `RowFilterContext` swaps between the real filter (from the `prebuild_row_filter_candidates` cache) and an empty filter as the decoder crosses row-group boundaries. * `row_filter.rs` — split `build_row_filter` into `prebuild_*` and `row_filter_from_prebuilt`: the expensive tree walk + candidate construction runs once per file, and the cheap per-RG bind runs at each toggle. Preserves the existing public `build_row_filter` API for non-toggling callers. * `opener/mod.rs` — wires the prebuilt cache into the decoder builder; skips page-index loading for fully-matched RGs (their per-row filter is a no-op, so page pruning saves nothing). * `sort.rs` — sort-order-aware RG reorder preserves the per-RG toggle through reordering. ## Tests * `fully_matched_rgs_skip_row_filter` (`dynamic_row_group_pruning.rs`) — 4 RGs of 3 rows each, predicate `v >= 3` makes RGs 1..=3 fully matched. Asserts: (a) results are correct, (b) `row_filter_skipped_fully_matched` fires at least once at the RG 0 → RG 1 boundary. * All 6 pre-existing `dynamic_row_group_pruning` tests still pass. * Full `datasource-parquet` lib suite (158 tests) passes. --- .../parquet/dynamic_row_group_pruning.rs | 73 ++ .../datasource-parquet/src/access_plan.rs | 104 +- datafusion/datasource-parquet/src/metrics.rs | 13 + .../datasource-parquet/src/opener/mod.rs | 147 ++- .../datasource-parquet/src/push_decoder.rs | 182 +++- .../datasource-parquet/src/row_filter.rs | 912 ++++++++++++++++-- datafusion/datasource-parquet/src/sort.rs | 28 +- .../test_files/push_down_filter_parquet.slt | 52 +- 8 files changed, 1309 insertions(+), 202 deletions(-) diff --git a/datafusion/core/tests/parquet/dynamic_row_group_pruning.rs b/datafusion/core/tests/parquet/dynamic_row_group_pruning.rs index b72c56ace5acd..35096c2e39c10 100644 --- a/datafusion/core/tests/parquet/dynamic_row_group_pruning.rs +++ b/datafusion/core/tests/parquet/dynamic_row_group_pruning.rs @@ -433,3 +433,76 @@ async fn dynamic_rg_pruning_coexists_with_row_filter() { output.description(), ); } + +/// Per-RG `fully_matched` `RowFilter` skip optimization. +/// +/// Stats prove that every row of a fully-matched row group satisfies the +/// pushdown predicate, so the parquet decoder can skip the per-row +/// `RowFilter` for that RG entirely. The stream rebuilds the decoder at +/// the boundary with an empty `RowFilter` and toggles back to the real +/// one at the next non-fully-matched RG. +/// +/// Layout: 4 RGs of 3 values each. Predicate `v >= 3` makes RG 0 a +/// straddler (some rows fail) but RGs 1..=3 fully matched (every value +/// >= 3 by stats). RG 0 keeps the row filter, then the toggle flips to +/// "no filter" when we enter the fully-matched run. +/// +/// Expected behavior: +/// - the static prune marks RGs 1..=3 as fully_matched at file open; +/// - the stream installs the real `RowFilter` initially (RG 0 not fm); +/// - at the RG 0 → RG 1 boundary the toggle rebuilds with empty filter +/// and bumps `row_filter_skipped_fully_matched`; +/// - the query result is identical to running with the filter on. +#[tokio::test] +async fn fully_matched_rgs_skip_row_filter() { + let schema = Arc::new(Schema::new(vec![Field::new("v", DataType::Int64, false)])); + // 4 RGs of 3 rows each. + // RG 0: 1, 2, 3 ← `v >= 3` keeps {3}; stats: min=1, max=3, NOT fm + // RG 1: 4, 5, 6 ← all >= 3 → fully matched + // RG 2: 7, 8, 9 ← fully matched + // RG 3: 10,11,12 ← fully matched + let groups: [[i64; 3]; 4] = [[1, 2, 3], [4, 5, 6], [7, 8, 9], [10, 11, 12]]; + let batches: Vec = groups + .iter() + .map(|vals| { + let col: ArrayRef = Arc::new(Int64Array::from(vals.to_vec())); + RecordBatch::try_new(Arc::clone(&schema), vec![col]).unwrap() + }) + .collect(); + + let mut ctx = ContextWithParquet::with_custom_data( + Scenario::Int, + RowGroup(3), + Arc::clone(&schema), + batches, + ) + .await; + + let output = ctx + .query("SELECT v FROM t WHERE v >= 3 ORDER BY v ASC") + .await; + + // Correctness: every value >= 3, ascending. + let expected_rows: Vec = (3..=12).collect(); + assert_eq!(output.result_rows, expected_rows.len()); + let formatted = output.pretty_results(); + for v in expected_rows { + assert!( + formatted.contains(&format!("| {v} ")), + "output must contain {v}; got:\n{formatted}", + ); + } + + // Behavior: the per-RG `RowFilter` toggle must have fired at least + // once when transitioning from RG 0 (not fm) into the fully-matched + // run RGs 1..=3. + let skipped = output + .metric_value("row_filter_skipped_fully_matched") + .unwrap_or(0); + assert!( + skipped >= 1, + "row_filter_skipped_fully_matched must fire at least once; \ + skipped={skipped}\n{}", + output.description(), + ); +} diff --git a/datafusion/datasource-parquet/src/access_plan.rs b/datafusion/datasource-parquet/src/access_plan.rs index 8189c2378cece..a4e80b9dc2ea4 100644 --- a/datafusion/datasource-parquet/src/access_plan.rs +++ b/datafusion/datasource-parquet/src/access_plan.rs @@ -571,12 +571,81 @@ impl ParquetAccessPlan { row_group_meta_data: &[RowGroupMetaData], ) -> Result { let row_group_indexes = self.row_group_indexes(); + // Carry `fully_matched` flags in the same order as + // `row_group_indexes` so downstream code (per-RG `RowFilter` skip) + // can look them up positionally. + let fully_matched: Vec = row_group_indexes + .iter() + .map(|&idx| self.fully_matched[idx]) + .collect(); let row_selection = self.into_overall_row_selection(row_group_meta_data)?; - PreparedAccessPlan::new(row_group_indexes, row_selection) + let (row_group_indexes, fully_matched, row_selection) = strip_empty_row_groups( + row_group_indexes, + fully_matched, + row_selection, + row_group_meta_data, + ); + + PreparedAccessPlan::new(row_group_indexes, fully_matched, row_selection) } } +/// Strip row groups whose post-pruning `RowSelection` selects zero rows. +/// +/// arrow-rs's push decoder silently advances past such row groups inside +/// `try_next_reader`, but the rest of DataFusion (per-RG metadata maps, +/// the runtime dynamic-pruner, the per-RG `RowFilter` toggle) assumes a +/// 1:1 correspondence between the prepared plan and the readers the +/// decoder hands back. Removing these empty entries here keeps that +/// invariant and lets downstream code consult per-RG state — like +/// [`PreparedAccessPlan::fully_matched`] — without going out of sync. +/// +/// The flat `RowSelection` is split per row group with +/// [`RowSelection::split_off`] (mirroring arrow-rs's own logic) and the +/// surviving segments are concatenated back into the result selection. +/// When `row_selection` is `None` (no page-index pruning, no +/// user-supplied selection) no row group can be empty and the inputs are +/// returned unchanged. +fn strip_empty_row_groups( + row_group_indexes: Vec, + fully_matched: Vec, + row_selection: Option, + row_group_meta_data: &[RowGroupMetaData], +) -> (Vec, Vec, Option) { + let Some(mut remaining) = row_selection else { + return (row_group_indexes, fully_matched, None); + }; + + let mut kept_indexes = Vec::with_capacity(row_group_indexes.len()); + let mut kept_fully_matched = Vec::with_capacity(fully_matched.len()); + let mut kept_selectors: Vec = Vec::new(); + + for (i, &rg_idx) in row_group_indexes.iter().enumerate() { + let rg_row_count = row_group_meta_data[rg_idx].num_rows() as usize; + // `split_off` cuts off the first `rg_row_count` rows worth of + // selection — those are this RG's segment. The returned value is + // the segment, `remaining` keeps the rest. + let rg_segment = remaining.split_off(rg_row_count); + if rg_segment.row_count() > 0 { + kept_indexes.push(rg_idx); + kept_fully_matched.push(fully_matched[i]); + kept_selectors.extend(rg_segment.iter().copied()); + } + // Empty segment ⇒ arrow-rs would have silently skipped this RG + // anyway; drop it from our plan so per-RG bookkeeping stays in + // sync with the decoder. + } + + let result_selection = if kept_selectors.is_empty() { + None + } else { + Some(RowSelection::from(kept_selectors)) + }; + + (kept_indexes, kept_fully_matched, result_selection) +} + /// Represents a prepared, fully resolved [`ParquetAccessPlan`] /// /// The [`RowSelection`] represents the result of applying all pruning such as @@ -587,6 +656,11 @@ impl ParquetAccessPlan { pub(crate) struct PreparedAccessPlan { /// Row group indexes to read pub(crate) row_group_indexes: Vec, + /// Per-RG `fully_matched` flag, positionally aligned with + /// [`Self::row_group_indexes`]. A `true` entry means stats already + /// proved every row of this RG passes the predicate, so the per-row + /// `RowFilter` can be skipped for it. + pub(crate) fully_matched: Vec, /// Optional row selection for filtering within row groups pub(crate) row_selection: Option, } @@ -595,10 +669,13 @@ impl PreparedAccessPlan { /// Create a new prepared access plan fn new( row_group_indexes: Vec, + fully_matched: Vec, row_selection: Option, ) -> Result { + debug_assert_eq!(row_group_indexes.len(), fully_matched.len()); Ok(Self { row_group_indexes, + fully_matched, row_selection, }) } @@ -713,13 +790,18 @@ impl PreparedAccessPlan { } }; - // Apply the reordering + // Apply the reordering — `fully_matched` must be permuted alongside + // `row_group_indexes` so the two stay positionally aligned for the + // per-RG `RowFilter` skip path. let original_indexes = self.row_group_indexes.clone(); - self.row_group_indexes = sorted_indices + let original_fully_matched = self.fully_matched.clone(); + let order: Vec = sorted_indices .values() .iter() - .map(|&i| original_indexes[i as usize]) + .map(|&i| i as usize) .collect(); + self.row_group_indexes = order.iter().map(|&i| original_indexes[i]).collect(); + self.fully_matched = order.iter().map(|&i| original_fully_matched[i]).collect(); Ok(self) } @@ -729,8 +811,9 @@ impl PreparedAccessPlan { // Get the row group indexes before reversing let row_groups_to_scan = self.row_group_indexes.clone(); - // Reverse the row group indexes + // Reverse the row group indexes (and the parallel `fully_matched`) self.row_group_indexes = self.row_group_indexes.into_iter().rev().collect(); + self.fully_matched = self.fully_matched.into_iter().rev().collect(); // If we have a row selection, reverse it to match the new row group order if let Some(row_selection) = self.row_selection { @@ -1121,7 +1204,7 @@ mod test { #[test] fn reorder_by_statistics_sorts_row_groups_asc_by_min() { let metadata = parquet_metadata_with_int_mins(&[50, 10, 100]); - let plan = PreparedAccessPlan::new(vec![0, 1, 2], None).unwrap(); + let plan = PreparedAccessPlan::new(vec![0, 1, 2], vec![false; 3], None).unwrap(); let result = plan .reorder_by_statistics( @@ -1140,7 +1223,8 @@ mod test { fn reorder_by_statistics_skips_when_row_selection_present() { let metadata = parquet_metadata_with_int_mins(&[50, 10]); let selection = RowSelection::from(vec![RowSelector::select(100)]); - let plan = PreparedAccessPlan::new(vec![0, 1], Some(selection)).unwrap(); + let plan = + PreparedAccessPlan::new(vec![0, 1], vec![false; 2], Some(selection)).unwrap(); let result = plan .reorder_by_statistics( @@ -1157,7 +1241,7 @@ mod test { #[test] fn reorder_by_statistics_skips_when_at_most_one_row_group() { let metadata = parquet_metadata_with_int_mins(&[50]); - let plan = PreparedAccessPlan::new(vec![0], None).unwrap(); + let plan = PreparedAccessPlan::new(vec![0], vec![false; 1], None).unwrap(); let result = plan .reorder_by_statistics( @@ -1176,7 +1260,7 @@ mod test { #[test] fn reorder_by_statistics_skips_for_non_column_sort_expr() { let metadata = parquet_metadata_with_int_mins(&[50, 10]); - let plan = PreparedAccessPlan::new(vec![0, 1], None).unwrap(); + let plan = PreparedAccessPlan::new(vec![0, 1], vec![false; 2], None).unwrap(); let arrow_schema = arrow_schema_a_int(); let order = LexOrdering::new(vec![PhysicalSortExpr { expr: Arc::new(BinaryExpr::new( @@ -1205,7 +1289,7 @@ mod test { #[test] fn reorder_by_statistics_skips_when_column_not_in_arrow_schema() { let metadata = parquet_metadata_with_int_mins(&[50, 10]); - let plan = PreparedAccessPlan::new(vec![0, 1], None).unwrap(); + let plan = PreparedAccessPlan::new(vec![0, 1], vec![false; 2], None).unwrap(); // Arrow schema only has "a"; the sort references "b". let arrow_schema = arrow_schema_a_int(); let order = LexOrdering::new(vec![PhysicalSortExpr { diff --git a/datafusion/datasource-parquet/src/metrics.rs b/datafusion/datasource-parquet/src/metrics.rs index cbdcb73196b17..c8482cdb2c43f 100644 --- a/datafusion/datasource-parquet/src/metrics.rs +++ b/datafusion/datasource-parquet/src/metrics.rs @@ -61,6 +61,12 @@ pub struct ParquetFileMetrics { /// the initial pruning but were proved unreachable mid-scan after the /// dynamic filter tightened. pub row_groups_pruned_dynamic_filter: Count, + /// Number of row groups for which the per-row + /// [`RowFilter`](parquet::arrow::arrow_reader::RowFilter) was skipped + /// because the static stats proved every row of the RG satisfies the + /// predicate. The decoder is rebuilt at the boundary with an empty + /// row filter so the upcoming RG decodes without per-row evaluation. + pub row_filter_skipped_fully_matched: Count, /// Total number of bytes scanned pub bytes_scanned: Count, /// Total rows filtered out by predicates pushed into parquet scan @@ -211,6 +217,12 @@ impl ParquetFileMetrics { .with_type(MetricType::Summary) .counter("row_groups_pruned_dynamic_filter", partition); + let row_filter_skipped_fully_matched = MetricBuilder::new(metrics) + .with_new_label("filename", filename.to_string()) + .with_type(MetricType::Summary) + .with_category(MetricCategory::Rows) + .counter("row_filter_skipped_fully_matched", partition); + Self { files_ranges_pruned_statistics, predicate_evaluation_errors, @@ -231,6 +243,7 @@ impl ParquetFileMetrics { predicate_cache_inner_records, predicate_cache_records, row_groups_pruned_dynamic_filter, + row_filter_skipped_fully_matched, } } diff --git a/datafusion/datasource-parquet/src/opener/mod.rs b/datafusion/datasource-parquet/src/opener/mod.rs index af97a192fa7ce..775a3c61142b2 100644 --- a/datafusion/datasource-parquet/src/opener/mod.rs +++ b/datafusion/datasource-parquet/src/opener/mod.rs @@ -29,7 +29,6 @@ use crate::page_filter::PagePruningAccessPlanFilter; use crate::push_decoder::{ DecoderBuilderConfig, PushDecoderStreamState, RgPlanEntry, RowGroupPruner, }; -use crate::row_filter::RowFilterGenerator; use crate::row_group_filter::RowGroupAccessPlanFilter; use crate::{ BloomFilterStatistics, Int96Coercer, ParquetAccessPlan, ParquetFileMetrics, @@ -41,7 +40,6 @@ use arrow::datatypes::DataType; use datafusion_datasource::morsel::{Morsel, MorselPlan, MorselPlanner, Morselizer}; use datafusion_physical_expr::projection::ProjectionExprs; use datafusion_physical_expr_adapter::replace_columns_with_literals; -use datafusion_physical_expr_adapter::rewrite::rewrite_input_file_name_in_projection; use std::collections::{HashMap, VecDeque}; use std::fmt; use std::future::Future; @@ -797,9 +795,6 @@ impl ParquetMorselizer { .transpose()?; } - // Replace any `input_file_name()` UDFs in the projection with a literal for this file. - projection = rewrite_input_file_name_in_projection(projection, &file_name)?; - let predicate_creation_errors = MetricBuilder::new(&self.metrics) .with_category(MetricCategory::Rows) .global_counter("num_predicate_creation_errors"); @@ -1386,18 +1381,26 @@ impl RowGroupsPrunedParquetOpen { prepared.virtual_state.as_deref(), )?; - let (decoder, rg_plan) = { + let (decoder, rg_plan, filter_installed, row_filter_context) = { let pushdown_predicate = prepared .pushdown_filters .then_some(prepared.predicate.as_ref()) .flatten(); - let mut row_filter_generator = RowFilterGenerator::new( - pushdown_predicate, - &prepared.physical_file_schema, - file_metadata.as_ref(), - prepared.reorder_predicates, - &prepared.file_metrics, - ); + // Precompute the prebuilt candidate list once per file. Both the + // initial `RowFilter` and any per-RG rebuilds (via + // `RowFilterContext::build`) reuse it, so tree walks + // (`reassign_expr_columns`) and column resolution only run once — + // not once per row group. + let precomputed_context = pushdown_predicate.and_then(|predicate| { + crate::push_decoder::RowFilterContext::try_new( + predicate, + &prepared.physical_file_schema, + &file_metadata, + prepared.reorder_predicates, + prepared.file_metrics.clone(), + prepared.max_predicate_cache_size, + ) + }); // Build the prepared access plan first — `prepare_access_plan` may // call `reorder_by_statistics` (for `sort_order_for_reorder`) and @@ -1415,25 +1418,66 @@ impl RowGroupsPrunedParquetOpen { }; let prepared_access_plan = prepare_access_plan(access_plan)?; + // Build `rg_plan` parallel to the decoder's view: the + // `prepared_access_plan` has already had its empty-selection + // row groups stripped, so 1:1 correspondence with the readers + // arrow-rs will hand back is restored. We zip with the + // `fully_matched` flag so the stream can toggle the per-row + // `RowFilter` per RG. let rg_plan: VecDeque = prepared_access_plan .row_group_indexes .iter() .copied() - .map(|rg_index| RgPlanEntry { rg_index }) + .zip(prepared_access_plan.fully_matched.iter().copied()) + .map(|(rg_index, fully_matched)| RgPlanEntry { + rg_index, + fully_matched, + }) .collect(); + // Decide the initial row filter state based on the first RG to + // read. If that RG is `fully_matched` the per-row predicate is + // a no-op for every row, so we install an empty `RowFilter` + // (arrow-rs's `has_predicates` check then short-circuits the + // per-row eval) and the stream toggles back to the real filter + // at the first non-fully-matched RG boundary. + // + // `RowFilterContext` carries everything `build_row_filter` + // needs so the stream can regenerate the filter later — the + // installed filter is owned by the decoder and is not + // recoverable once replaced. + let first_rg_fully_matched = rg_plan.front().is_some_and(|e| e.fully_matched); + let initial_filter = precomputed_context + .as_ref() + .and_then(|ctx| ctx.build_row_filter()); + let row_filter_context = precomputed_context; + let mut builder = decoder_config.build(prepared_access_plan, reader_metadata.clone()); - if let Some(row_filter) = row_filter_generator.next_filter() { - builder = builder.with_row_filter(row_filter); - if let Some(max_predicate_cache_size) = prepared.max_predicate_cache_size - { - builder = - builder.with_max_predicate_cache_size(max_predicate_cache_size); + let mut filter_installed = false; + if let Some(row_filter) = initial_filter { + if first_rg_fully_matched { + builder = builder.with_row_filter( + parquet::arrow::arrow_reader::RowFilter::new(vec![]), + ); + } else { + builder = builder.with_row_filter(row_filter); + filter_installed = true; + if let Some(max_predicate_cache_size) = + prepared.max_predicate_cache_size + { + builder = builder + .with_max_predicate_cache_size(max_predicate_cache_size); + } } } - (builder.build()?, rg_plan) + ( + builder.build()?, + rg_plan, + filter_installed, + row_filter_context, + ) }; let predicate_cache_inner_records = @@ -1476,6 +1520,10 @@ impl RowGroupsPrunedParquetOpen { .file_metrics .row_groups_pruned_dynamic_filter .clone(); + let row_filter_skipped_fully_matched = prepared + .file_metrics + .row_filter_skipped_fully_matched + .clone(); let stream = PushDecoderStreamState { decoder: Some(decoder), @@ -1489,6 +1537,9 @@ impl RowGroupsPrunedParquetOpen { baseline_metrics: prepared.baseline_metrics, row_group_pruner, row_groups_pruned_dynamic, + row_filter_context, + filter_installed, + row_filter_skipped_fully_matched, } .into_stream(); @@ -3285,12 +3336,8 @@ mod test { /// (e.g. `row_number`) plumbed through `TableSchema`/`ParquetOpener`. mod virtual_columns { use super::*; - use arrow::array::{Array, Int64Array, StringArray}; + use arrow::array::{Array, Int64Array}; use arrow::datatypes::FieldRef; - use datafusion_common::config::ConfigOptions; - use datafusion_expr::ScalarUDF; - use datafusion_functions::core::input_file_name::InputFileNameFunc; - use datafusion_physical_expr::{ScalarFunctionExpr, projection::ProjectionExpr}; use parquet::arrow::RowNumber; /// Build a parquet `row_number` virtual column field. Spark's @@ -3304,16 +3351,6 @@ mod test { ) } - fn input_file_name_expr() -> Arc { - Arc::new(ScalarFunctionExpr::new( - "input_file_name", - Arc::new(ScalarUDF::from(InputFileNameFunc::new())), - vec![], - Arc::new(Field::new("input_file_name", DataType::Utf8, true)), - Arc::new(ConfigOptions::default()), - )) - } - /// Collect every `Int64` value from the given column in every batch /// of a stream. Used to verify the `row_number` column end to end. async fn collect_int64_values( @@ -3433,44 +3470,6 @@ mod test { assert_eq!(row_numbers, vec![0, 1, 2, 3]); } - #[tokio::test] - async fn test_input_file_name_projection() { - let store = Arc::new(InMemory::new()) as Arc; - let path = "dir/input_file_name.parquet"; - let (file_schema, data_size) = write_grouped_file(&store, path, 1, 3).await; - - let projection = ProjectionExprs::new([ - ProjectionExpr::new(Arc::new(Column::new("value", 0)), "value"), - ProjectionExpr::new(input_file_name_expr(), "file_name"), - ]); - - let morselizer = ParquetMorselizerBuilder::new() - .with_store(Arc::clone(&store)) - .with_schema(file_schema) - .with_projection(projection) - .build(); - - let file = - PartitionedFile::new(path.to_string(), u64::try_from(data_size).unwrap()); - let mut stream = open_file(&morselizer, file).await.unwrap(); - let batch = stream.next().await.unwrap().unwrap(); - assert!(stream.next().await.is_none()); - - assert_eq!(batch.num_columns(), 2); - assert_eq!(batch.schema().field(0).name(), "value"); - assert_eq!(batch.schema().field(1).name(), "file_name"); - - let file_names = batch - .column(1) - .as_any() - .downcast_ref::() - .expect("file_name column should be Utf8"); - assert_eq!(file_names.len(), 3); - for i in 0..file_names.len() { - assert_eq!(file_names.value(i), path); - } - } - #[tokio::test] async fn test_row_index_multi_row_group() { let store = Arc::new(InMemory::new()) as Arc; diff --git a/datafusion/datasource-parquet/src/push_decoder.rs b/datafusion/datasource-parquet/src/push_decoder.rs index 31bd365a4631d..b457aa8600fce 100644 --- a/datafusion/datasource-parquet/src/push_decoder.rs +++ b/datafusion/datasource-parquet/src/push_decoder.rs @@ -47,7 +47,7 @@ use parquet::DecodeResult; use parquet::arrow::ProjectionMask; use parquet::arrow::arrow_reader::metrics::ArrowReaderMetrics; use parquet::arrow::arrow_reader::{ - ArrowReaderMetadata, ParquetRecordBatchReader, RowSelectionPolicy, + ArrowReaderMetadata, ParquetRecordBatchReader, RowFilter, RowSelectionPolicy, }; use parquet::arrow::async_reader::AsyncFileReader; use parquet::arrow::push_decoder::{ParquetPushDecoder, ParquetPushDecoderBuilder}; @@ -59,8 +59,12 @@ use datafusion_physical_expr_common::physical_expr::PhysicalExpr; use datafusion_physical_plan::metrics::{BaselineMetrics, Count, Gauge}; use datafusion_pruning::{PruningPredicate, build_pruning_predicate}; +use crate::ParquetFileMetrics; use crate::access_plan::PreparedAccessPlan; use crate::decoder_projection::DecoderProjection; +use crate::row_filter::{ + PrebuiltRowFilterCandidate, prebuild_row_filter_candidates, row_filter_from_prebuilt, +}; use crate::row_group_filter::RowGroupPruningStatistics; /// Shared options applied to the [`ParquetPushDecoderBuilder`] for a file @@ -109,6 +113,12 @@ impl DecoderBuilderConfig<'_> { #[derive(Debug, Clone)] pub(crate) struct RgPlanEntry { pub(crate) rg_index: usize, + /// `true` when the static pruning predicate proved every row of this + /// RG satisfies the predicate. The push-decoder stream uses this to + /// skip installing the per-row `RowFilter` on RGs where it would be a + /// no-op, rebuilding the decoder via `into_builder` at boundaries + /// where the filter status flips. + pub(crate) fully_matched: bool, } /// Runtime row-group pruner driven by a dynamic predicate (e.g. the @@ -265,6 +275,82 @@ pub(crate) struct PushDecoderStreamState { pub(crate) row_group_pruner: Option, /// Count of row groups skipped at runtime by [`Self::row_group_pruner`]. pub(crate) row_groups_pruned_dynamic: Count, + /// Side-channel state for regenerating the parquet [`RowFilter`] when + /// the per-RG `fully_matched` toggle flips from skip → install. `None` + /// when the scan has no pushdown predicate, so no filter can be + /// installed (and the toggle is a no-op). + pub(crate) row_filter_context: Option, + /// Whether the currently-installed decoder is running with a non-empty + /// row filter. Toggled per RG by the `fully_matched` skip path. + pub(crate) filter_installed: bool, + /// Count of row groups for which the per-row [`RowFilter`] was + /// suppressed because the upcoming RG is `fully_matched`. + pub(crate) row_filter_skipped_fully_matched: Count, +} + +/// Side-channel state that lets [`PushDecoderStreamState`] **rebuild** the +/// parquet [`RowFilter`] mid-scan. +/// +/// The decoder owns the filter once installed, but `Box` +/// has no clone path, so a filter that was replaced at a previous boundary +/// cannot be reinstalled later. This struct keeps a pre-built candidate list +/// alongside the stream so the next non-fully-matched row group can be +/// wrapped into a fresh [`RowFilter`] without redoing the tree walks and +/// column resolution that the initial build did. +pub(crate) struct RowFilterContext { + /// Prebuilt candidates: expression already column-reassigned, projection + /// mask already resolved. Shared across the file's row groups. `Arc` so + /// cloning into stream state is cheap. + pub(crate) prebuilt: Arc>, + pub(crate) reorder_predicates: bool, + pub(crate) file_metrics: ParquetFileMetrics, + pub(crate) max_predicate_cache_size: Option, +} + +impl RowFilterContext { + /// Precompute the candidate list from the raw predicate + file schema + + /// metadata. Returns `None` when the predicate has no push-downable + /// conjuncts (mirrors the file-open path behaviour). + pub(crate) fn try_new( + predicate: &Arc, + physical_file_schema: &SchemaRef, + file_metadata: &Arc, + reorder_predicates: bool, + file_metrics: ParquetFileMetrics, + max_predicate_cache_size: Option, + ) -> Option { + match prebuild_row_filter_candidates( + predicate, + physical_file_schema, + file_metadata.as_ref(), + ) { + Ok(Some(prebuilt)) => Some(Self { + prebuilt: Arc::new(prebuilt), + reorder_predicates, + file_metrics, + max_predicate_cache_size, + }), + Ok(None) => None, + Err(e) => { + debug!("Ignoring error prebuilding row filter candidates: {e}"); + None + } + } + } + + /// Build a fresh [`RowFilter`] for the next non-fully-matched run using + /// the cached candidates. Cheap: no tree walks, only counter allocation + /// and (optionally) a sort by `required_bytes`. + pub(crate) fn build_row_filter(&self) -> Option { + if self.prebuilt.is_empty() { + return None; + } + Some(row_filter_from_prebuilt( + &self.prebuilt, + self.reorder_predicates, + &self.file_metrics, + )) + } } impl PushDecoderStreamState { @@ -330,11 +416,36 @@ impl PushDecoderStreamState { // been handed back yet), step 3 drives it forward and we get // another chance at the next boundary — the pruner is stateful // and idempotent, so deferring loses nothing. - let at_boundary = self - .decoder - .as_ref() - .expect("decoder present") - .is_at_row_group_boundary(); + let decoder_ref = self.decoder.as_ref().expect("decoder present"); + let at_boundary = decoder_ref.is_at_row_group_boundary(); + // Sync `rg_plan` with the row group the decoder will actually + // emit next. arrow-rs's `try_next_reader` silently advances + // past row groups whose row selection is empty (e.g. when + // page-index pruning has already eliminated every page of + // that RG via the `ColumnIndex` path inside `try_build`). + // Without this peek, `rg_plan.front()` would drift off-by-one + // from the decoder's frontier and the per-RG toggle below + // would target the wrong row group. + if at_boundary { + match decoder_ref.peek_next_row_group() { + Ok(Some(actual)) => { + while let Some(front) = self.rg_plan.front() { + if front.rg_index == actual { + break; + } + self.rg_plan.pop_front(); + } + } + Ok(None) => { + if !self.rg_plan.is_empty() { + // Decoder has nothing left to emit — drain our plan + // so the stream finishes cleanly. + self.rg_plan.clear(); + } + } + Err(e) => return Some((Err(DataFusionError::from(e)), self)), + } + } if at_boundary && !self.rg_plan.is_empty() { let mut pruned_count = 0usize; if let Some(pruner) = self.row_group_pruner.as_mut() { @@ -349,7 +460,21 @@ impl PushDecoderStreamState { } self.rg_plan = kept; } - if pruned_count > 0 { + + // Decide whether the per-row `RowFilter` needs to be + // toggled for the upcoming RG. `desired_filter` is + // `Some(true)` when the next RG needs a real filter, + // `Some(false)` when it's fully-matched (filter is a + // no-op, so we suppress it), and `None` when there is no + // pushdown predicate at all (toggling is meaningless). + let desired_filter: Option = self + .row_filter_context + .as_ref() + .and_then(|_| self.rg_plan.front().map(|e| !e.fully_matched)); + let filter_needs_toggle = + desired_filter.is_some_and(|want| want != self.filter_installed); + + if pruned_count > 0 || filter_needs_toggle { if self.rg_plan.is_empty() { return None; } @@ -357,7 +482,48 @@ impl PushDecoderStreamState { let new_indices: Vec = self.rg_plan.iter().map(|e| e.rg_index).collect(); let rebuilt = match decoder.into_builder() { - Ok(b) => b.with_row_groups(new_indices).build(), + Ok(mut builder) => { + builder = builder.with_row_groups(new_indices); + if filter_needs_toggle { + let want_filter = desired_filter + .expect("filter_needs_toggle ⇒ desired Some"); + if want_filter { + let ctx = self + .row_filter_context + .as_ref() + .expect("filter_needs_toggle ⇒ context set"); + match ctx.build_row_filter() { + Some(filter) => { + builder = builder.with_row_filter(filter); + if let Some(cap) = + ctx.max_predicate_cache_size + { + builder = builder + .with_max_predicate_cache_size(cap); + } + self.filter_installed = true; + } + None => { + // Filter could not be rebuilt; + // install empty filter so the + // decoder runs unfiltered for + // this run rather than failing. + builder = builder + .with_row_filter(RowFilter::new(vec![])); + self.filter_installed = false; + } + } + } else { + // Skip per-row filtering for the + // upcoming fully-matched RG. + builder = + builder.with_row_filter(RowFilter::new(vec![])); + self.filter_installed = false; + self.row_filter_skipped_fully_matched.add(1); + } + } + builder.build() + } Err(e) => Err(e), }; match rebuilt { diff --git a/datafusion/datasource-parquet/src/row_filter.rs b/datafusion/datasource-parquet/src/row_filter.rs index a375e6611e004..3d396bfd85445 100644 --- a/datafusion/datasource-parquet/src/row_filter.rs +++ b/datafusion/datasource-parquet/src/row_filter.rs @@ -65,29 +65,32 @@ //! - `WHERE s['value'] > 5` — pushed down (accesses a primitive leaf) //! - `WHERE s IS NOT NULL` — not pushed down (references the whole struct) +use std::collections::BTreeSet; use std::sync::Arc; use arrow::array::BooleanArray; -use arrow::datatypes::{Schema, SchemaRef}; +use arrow::datatypes::{DataType, Field, Schema, SchemaRef}; use arrow::error::{ArrowError, Result as ArrowResult}; use arrow::record_batch::RecordBatch; +use datafusion_functions::core::file_row_index::FileRowIndexFunc; +use datafusion_functions::core::getfield::GetFieldFunc; use parquet::arrow::ProjectionMask; use parquet::arrow::arrow_reader::{ArrowPredicate, RowFilter}; use parquet::file::metadata::ParquetMetaData; +use parquet::schema::types::SchemaDescriptor; use datafusion_common::Result; use datafusion_common::cast::as_boolean_array; -use datafusion_common::tree_node::TreeNode; -use datafusion_physical_expr::utils::reassign_expr_columns; +use datafusion_common::tree_node::{TreeNode, TreeNodeRecursion, TreeNodeVisitor}; +use datafusion_physical_expr::ScalarFunctionExpr; +use datafusion_physical_expr::expressions::{Column, Literal}; +use datafusion_physical_expr::utils::{collect_columns, reassign_expr_columns}; use datafusion_physical_expr::{PhysicalExpr, split_conjunction}; use datafusion_physical_plan::metrics; use super::ParquetFileMetrics; use super::supported_predicates::supports_list_predicates; -use crate::projection_read_plan::{ - ParquetReadPlan, PushdownChecker, PushdownColumns, assemble_read_plan, -}; /// A "compiled" predicate passed to `ParquetRecordBatchStream` to perform /// row-level filtering during parquet decoding. @@ -186,6 +189,22 @@ pub(crate) struct FilterCandidate { read_plan: ParquetReadPlan, } +/// The result of resolving which Parquet leaf columns and Arrow schema fields +/// are needed to evaluate an expression against a Parquet file +/// +/// This is the shared output of the column resolution pipeline used by both +/// the row filter to build `ArrowPredicate`s and the opener to build `ProjectionMask`s +#[derive(Debug, Clone)] +pub(crate) struct ParquetReadPlan { + /// Projection mask built from leaf column indices in the Parquet schema. + /// Using a `ProjectionMask` directly (rather than raw indices) prevents + /// bugs from accidentally mixing up root vs leaf indices. + pub projection_mask: ProjectionMask, + /// The projected Arrow schema containing only the columns/fields required + /// Struct types are pruned to include only the accessed sub-fields + pub projected_schema: SchemaRef, +} + /// Helper to build a `FilterCandidate`. /// /// This will do several things: @@ -227,6 +246,289 @@ impl FilterCandidateBuilder { } } +/// Traverses a `PhysicalExpr` tree to determine if any column references would +/// prevent the expression from being pushed down to the parquet decoder. +/// +/// An expression cannot be pushed down if it references: +/// - Unsupported nested columns (whole struct references or list fields that are +/// not covered by the supported predicate set) +/// - Columns that don't exist in the file schema +/// +/// Struct field access via `get_field` is supported when the resolved leaf type +/// is primitive (e.g. `get_field(struct_col, 'field') > 5`). +struct PushdownChecker<'schema> { + /// Does the expression require any non-primitive columns (like structs)? + non_primitive_columns: bool, + /// Does the expression reference any columns not present in the file schema? + projected_columns: bool, + /// Does the expression references a ScalarUDF that requires some rewrite + /// and therefore can't be pushed down into the row-filter. + has_unpushable_udfs: bool, + /// Indices into the file schema of columns required to evaluate the expression. + /// Does not include struct columns accessed via `get_field`. + required_columns: Vec, + /// Struct field accesses via `get_field`. + struct_field_accesses: Vec, + /// Whether nested list columns are supported by the predicate semantics. + allow_list_columns: bool, + /// The Arrow schema of the parquet file. + file_schema: &'schema Schema, +} + +impl<'schema> PushdownChecker<'schema> { + fn new(file_schema: &'schema Schema, allow_list_columns: bool) -> Self { + Self { + non_primitive_columns: false, + projected_columns: false, + has_unpushable_udfs: false, + required_columns: Vec::new(), + struct_field_accesses: Vec::new(), + allow_list_columns, + file_schema, + } + } + + /// Checks whether a struct's root column exists in the file schema and, if so, + /// records its index so the entire struct is decoded for filter evaluation. + /// + /// This is called when we see a `get_field` expression that resolves to a + /// primitive leaf type. We only need the *root* column index because the + /// Parquet reader decodes all leaves of a struct together. + /// + /// # Example + /// + /// Given file schema `{a: Int32, s: Struct(foo: Utf8, bar: Int64)}` and the + /// expression `get_field(s, 'foo') = 'hello'`: + /// + /// - `column_name` = `"s"` (the root struct column) + /// - `file_schema.index_of("s")` returns `1` + /// - We push `1` into `required_columns` + /// - Return `None` (no issue — traversal continues in the caller) + /// + /// If `"s"` is not in the file schema (e.g. a projected-away column), we set + /// `projected_columns = true` and return `Jump` to skip the subtree. + fn check_struct_field_column( + &mut self, + column_name: &str, + field_path: Vec, + ) -> Option { + let Ok(idx) = self.file_schema.index_of(column_name) else { + self.projected_columns = true; + return Some(TreeNodeRecursion::Jump); + }; + + self.struct_field_accesses.push(StructFieldAccess { + root_index: idx, + field_path, + }); + + None + } + + fn check_single_column(&mut self, column_name: &str) -> Option { + let idx = match self.file_schema.index_of(column_name) { + Ok(idx) => idx, + Err(_) => { + // Column does not exist in the file schema, so we can't push this down. + self.projected_columns = true; + return Some(TreeNodeRecursion::Jump); + } + }; + + // Duplicates are handled by dedup() in into_sorted_columns() + self.required_columns.push(idx); + let data_type = self.file_schema.field(idx).data_type(); + + if DataType::is_nested(data_type) { + self.handle_nested_type(data_type) + } else { + None + } + } + + /// Determines whether a nested data type can be pushed down to Parquet decoding. + /// + /// Returns `Some(TreeNodeRecursion::Jump)` if the nested type prevents pushdown, + /// `None` if the type is supported and pushdown can continue. + fn handle_nested_type(&mut self, data_type: &DataType) -> Option { + if self.is_nested_type_supported(data_type) { + None + } else { + // Block pushdown for unsupported nested types: + // - Structs (regardless of predicate support) + // - Lists without supported predicates + self.non_primitive_columns = true; + Some(TreeNodeRecursion::Jump) + } + } + + /// Checks if a nested data type is supported for list column pushdown. + /// + /// List columns are only supported if: + /// 1. The data type is a list variant (List, LargeList, or FixedSizeList) + /// 2. The expression contains supported list predicates (e.g., array_has_all) + fn is_nested_type_supported(&self, data_type: &DataType) -> bool { + let is_list = matches!( + data_type, + DataType::List(_) | DataType::LargeList(_) | DataType::FixedSizeList(_, _) + ); + self.allow_list_columns && is_list + } + + #[inline] + fn prevents_pushdown(&self) -> bool { + self.non_primitive_columns || self.projected_columns || self.has_unpushable_udfs + } + + /// Consumes the checker and returns sorted, deduplicated column indices + /// wrapped in a `PushdownColumns` struct. + /// + /// This method sorts the column indices and removes duplicates. The sort + /// is required because downstream code relies on column indices being in + /// ascending order for correct schema projection. + fn into_sorted_columns(mut self) -> PushdownColumns { + self.required_columns.sort_unstable(); + self.required_columns.dedup(); + PushdownColumns { + required_columns: self.required_columns, + struct_field_accesses: self.struct_field_accesses, + } + } +} + +impl TreeNodeVisitor<'_> for PushdownChecker<'_> { + type Node = Arc; + + fn f_down(&mut self, node: &Self::Node) -> Result { + // Handle struct field access like `s['foo']['bar'] > 10`. + // + // DataFusion represents nested field access as `get_field(Column("s"), "foo")` + // (or chained: `get_field(get_field(Column("s"), "foo"), "bar")`). + // + // We intercept the outermost `get_field` on the way *down* the tree so + // the visitor never reaches the raw `Column("s")` node. Without this, + // `check_single_column` would see that `s` is a Struct and reject it. + // + // The strategy: + // 1. Match `get_field` whose first arg is a `Column` (the struct root). + // 2. Check that the *resolved* return type is primitive — meaning we've + // drilled all the way to a leaf (e.g. `s['foo']` → Utf8). + // 3. Record the root column index via `check_struct_field_column` and + // return `Jump` to skip visiting the children (the Column and the + // literal field-name args), since we've already handled them. + // + // If the return type is still nested (e.g. `s['nested_struct']` → Struct), + // we fall through and let normal traversal continue, which will + // eventually reject the expression when it hits the struct Column. + if let Some(func) = + ScalarFunctionExpr::try_downcast_func::(node.as_ref()) + { + let args = func.args(); + + if let Some(column) = args.first().and_then(|a| a.downcast_ref::()) { + // for Map columns, get_field performs a runtime key lookup rather than a + // schema-level field access so the entire Map column must be read, + // we skip the struct field optimization and defer to normal Column traversal + let is_map_column = self + .file_schema + .index_of(column.name()) + .ok() + .map(|idx| { + matches!( + self.file_schema.field(idx).data_type(), + DataType::Map(_, _) + ) + }) + .unwrap_or(false); + + let return_type = func.return_type(); + + if !is_map_column + && (!DataType::is_nested(return_type) + || self.is_nested_type_supported(return_type)) + { + // try to resolve all field name arguments to strinrg literals + // if any argument is not a string literal, we can not determine the exact + // leaf path so we fall back to reading the entire struct root column + let field_path = args[1..] + .iter() + .map(|arg| { + arg.downcast_ref::().and_then(|lit| { + lit.value().try_as_str().flatten().map(|s| s.to_string()) + }) + }) + .collect(); + + match field_path { + Some(path) => { + if let Some(recursion) = + self.check_struct_field_column(column.name(), path) + { + return Ok(recursion); + } + } + None => { + // Could not resolve field path — fall back to + // reading the entire struct root column. + if let Some(recursion) = + self.check_single_column(column.name()) + { + return Ok(recursion); + } + } + } + + return Ok(TreeNodeRecursion::Jump); + } + } + } + + if let Some(column) = node.downcast_ref::() + && let Some(recursion) = self.check_single_column(column.name()) + { + return Ok(recursion); + } + + if ScalarFunctionExpr::try_downcast_func::(node.as_ref()) + .is_some() + { + self.has_unpushable_udfs = true; + return Ok(TreeNodeRecursion::Jump); + } + + Ok(TreeNodeRecursion::Continue) + } +} + +/// Describes the nested column behavior for filter pushdown. +/// +/// This enum makes explicit the different states a predicate can be in +/// with respect to nested column handling during Parquet decoding. +/// Result of checking which columns are required for filter pushdown. +#[derive(Debug)] +struct PushdownColumns { + /// Sorted, unique column indices into the file schema required to evaluate + /// the filter expression. Must be in ascending order for correct schema + /// projection matching. Does not include struct columns accessed via `get_field`. + required_columns: Vec, + /// Struct field accesses via `get_field`. Each entry records the root struct + /// column index and the field path being accessed. + struct_field_accesses: Vec, +} + +/// Records a struct field access via `get_field(struct_col, 'field1', 'field2', ...)`. +/// +/// This allows the row filter to project only the specific Parquet leaf columns +/// needed by the filter, rather than all leaves of the struct. +#[derive(Debug, Clone)] +struct StructFieldAccess { + /// Arrow root column index of the struct in the file schema. + root_index: usize, + /// Field names forming the path into the struct. + /// e.g., `["value"]` for `s['value']`, `["outer", "inner"]` for `s['outer']['inner']`. + field_path: Vec, +} + /// Checks if a given expression can be pushed down to the parquet decoder. /// /// Returns `Some(PushdownColumns)` if the expression can be pushed down, @@ -268,16 +570,344 @@ pub(crate) fn build_parquet_read_plan( return Ok(None); }; - let (read_plan, leaf_indices) = assemble_read_plan( - &required_columns.required_columns, + let root_indices = &required_columns.required_columns; + + let mut leaf_indices = + leaf_indices_for_roots(root_indices.iter().copied(), schema_descr); + + let struct_leaf_indices = resolve_struct_field_leaves( &required_columns.struct_field_accesses, file_schema, schema_descr, ); + leaf_indices.extend_from_slice(&struct_leaf_indices); + leaf_indices.sort_unstable(); + leaf_indices.dedup(); let required_bytes = size_of_columns(&leaf_indices, metadata)?; - Ok(Some((read_plan, required_bytes))) + let projection_mask = + ProjectionMask::leaves(schema_descr, leaf_indices.iter().copied()); + + let projected_schema = build_filter_schema( + file_schema, + root_indices, + &required_columns.struct_field_accesses, + ); + + Ok(Some(( + ParquetReadPlan { + projection_mask, + projected_schema, + }, + required_bytes, + ))) +} + +/// Builds a unified [`ParquetReadPlan`] for a set of projection expressions +/// +/// Unlike [`build_parquet_read_plan`] (which is used for filter pushdown and +/// returns `None` when an expression references unsupported nested types or +/// missing columns), this function always succeeds. It collects every column +/// that *can* be resolved in the file and produces a leaf-level projection +/// mask. Columns missing from the file are silently skipped since the projection +/// layer handles those by inserting nulls. +pub(crate) fn build_projection_read_plan( + exprs: impl IntoIterator>, + file_schema: &Schema, + schema_descr: &SchemaDescriptor, +) -> ParquetReadPlan { + // fast path: if every expression is a plain Column reference, skip all + // struct analysis and use root-level projection directly + let exprs = exprs.into_iter().collect::>(); + let all_plain_columns = exprs.iter().all(|e| e.downcast_ref::().is_some()); + + if all_plain_columns { + let mut root_indices: Vec = exprs + .iter() + .map(|e| e.downcast_ref::().unwrap().index()) + .collect(); + root_indices.sort_unstable(); + root_indices.dedup(); + + let projection_mask = + ProjectionMask::roots(schema_descr, root_indices.iter().copied()); + let projected_schema = Arc::new( + file_schema + .project(&root_indices) + .expect("valid column indices"), + ); + + return ParquetReadPlan { + projection_mask, + projected_schema, + }; + } + + // secondary fast path: if the schema has no struct columns, we can skip + // PushdownChecker traversal and use root-level projection + let has_struct_columns = file_schema + .fields() + .iter() + .any(|f| matches!(f.data_type(), DataType::Struct(_))); + + if !has_struct_columns { + let mut root_indices = exprs + .into_iter() + .flat_map(|e| collect_columns(&e).into_iter().map(|col| col.index())) + .collect::>(); + + root_indices.sort_unstable(); + root_indices.dedup(); + + let projection_mask = + ProjectionMask::roots(schema_descr, root_indices.iter().copied()); + + let projected_schema = Arc::new( + file_schema + .project(&root_indices) + .expect("valid column indices"), + ); + + return ParquetReadPlan { + projection_mask, + projected_schema, + }; + } + + let mut all_root_indices = Vec::new(); + let mut all_struct_accesses = Vec::new(); + + for expr in exprs { + let mut checker = PushdownChecker::new(file_schema, true); + let _ = expr.visit(&mut checker); + let columns = checker.into_sorted_columns(); + + all_root_indices.extend_from_slice(&columns.required_columns); + all_struct_accesses.extend(columns.struct_field_accesses); + } + + all_root_indices.sort_unstable(); + all_root_indices.dedup(); + + // when no struct field accesses were found, fall back to root-level projection + // to match the performance of the simple path + if all_struct_accesses.is_empty() { + let projection_mask = + ProjectionMask::roots(schema_descr, all_root_indices.iter().copied()); + let projected_schema = Arc::new( + file_schema + .project(&all_root_indices) + .expect("valid column indices"), + ); + + return ParquetReadPlan { + projection_mask, + projected_schema, + }; + } + + let leaf_indices = { + let mut out = + leaf_indices_for_roots(all_root_indices.iter().copied(), schema_descr); + let struct_leaf_indices = + resolve_struct_field_leaves(&all_struct_accesses, file_schema, schema_descr); + + out.extend_from_slice(&struct_leaf_indices); + out.sort_unstable(); + out.dedup(); + + out + }; + + let projection_mask = + ProjectionMask::leaves(schema_descr, leaf_indices.iter().copied()); + + let projected_schema = + build_filter_schema(file_schema, &all_root_indices, &all_struct_accesses); + + ParquetReadPlan { + projection_mask, + projected_schema, + } +} + +fn leaf_indices_for_roots( + root_indices: I, + schema_descr: &SchemaDescriptor, +) -> Vec +where + I: IntoIterator, +{ + // Always map root (Arrow) indices to Parquet leaf indices via the schema + // descriptor. Arrow root indices only equal Parquet leaf indices when the + // schema has no group columns (Struct, Map, etc.); when group columns + // exist, their children become separate leaves and shift all subsequent + // leaf indices. + // Struct columns are unsupported. + let root_set: BTreeSet<_> = root_indices.into_iter().collect(); + + (0..schema_descr.num_columns()) + .filter(|leaf_idx| { + root_set.contains(&schema_descr.get_column_root_idx(*leaf_idx)) + }) + .collect() +} + +/// Resolves struct field access to specific Parquet leaf column indices +/// +/// For every `StructFieldAccess`, finds the leaf columns in the Parquet schema +/// whose path matches the struct root name + field path. This avoids reading all +/// leaves of a struct when only specific fields are needed +fn resolve_struct_field_leaves( + accesses: &[StructFieldAccess], + file_schema: &Schema, + schema_descr: &SchemaDescriptor, +) -> Vec { + let mut leaf_indices = Vec::new(); + + for access in accesses { + let root_name = file_schema.field(access.root_index).name(); + let prefix = std::iter::once(root_name.as_str()) + .chain(access.field_path.iter().map(|p| p.as_str())) + .collect::>(); + + for leaf_idx in 0..schema_descr.num_columns() { + let col = schema_descr.column(leaf_idx); + let col_path = col.path().parts(); + + // A leaf matches if its path starts with our prefix. + // e.g., prefix=["s", "value"] matches leaf path ["s", "value"] + // prefix=["s", "outer"] matches ["s", "outer", "inner"] + + // a leaf matches if its path starts with our prefix + // for example: prefix=["s", "value"] matches leaf path ["s", "value"] + // prefix=["s", "outer"] matches ["s", "outer", "inner"] + let leaf_matches_path = col_path.len() >= prefix.len() + && col_path.iter().zip(prefix.iter()).all(|(a, b)| a == b); + + if leaf_matches_path { + leaf_indices.push(leaf_idx); + } + } + } + + leaf_indices +} + +/// Builds a filter schema that includes only the fields actually accessed by the +/// filter expression. +/// +/// For regular (non-struct) columns, the full field type is used. +/// For struct columns accessed via `get_field`, a pruned struct type is created +/// containing only the fields along the access path. Note: it must match the schema +/// that the Parquet reader produces when projecting specific struct leaves +fn build_filter_schema( + file_schema: &Schema, + regular_indices: &[usize], + struct_field_accesses: &[StructFieldAccess], +) -> SchemaRef { + let regular_set: BTreeSet = regular_indices.iter().copied().collect(); + + let all_indices = regular_indices + .iter() + .copied() + .chain( + struct_field_accesses + .iter() + .map(|&StructFieldAccess { root_index, .. }| root_index), + ) + .collect::>(); + + let fields = all_indices + .iter() + .map(|&idx| { + let field = file_schema.field(idx); + + // if this column appears as a regular (whole-column) reference, + // keep the full type + // + // Pruning is only valid when the column is accessed exclusively + // through struct field accesses + if regular_set.contains(&idx) { + return Arc::new(field.clone()); + } + + // collect all field paths that access this root struct column + let field_paths = struct_field_accesses + .iter() + .filter_map( + |&StructFieldAccess { + root_index, + ref field_path, + }| { + (root_index == idx).then_some(field_path.as_slice()) + }, + ) + .collect::>(); + + if field_paths.is_empty() { + return Arc::new(field.clone()); + } + + let pruned_data_type = prune_struct_type(field.data_type(), &field_paths); + Arc::new(Field::new( + field.name(), + pruned_data_type, + field.is_nullable(), + )) + }) + .collect::>(); + + Arc::new(Schema::new_with_metadata( + fields, + file_schema.metadata().clone(), + )) +} + +fn prune_struct_type(dt: &DataType, paths: &[&[String]]) -> DataType { + let DataType::Struct(fields) = dt else { + return dt.clone(); + }; + + let needed = paths + .iter() + .filter_map(|p| p.first().map(|s| s.as_str())) + .collect::>(); + + let pruned_fields = fields + .iter() + .filter_map(|f| { + if !needed.contains(f.name().as_str()) { + return None; + } + + let sub_paths = paths + .iter() + .filter_map(|path| { + if path.first().map(|s| s.as_str()) == Some(f.name()) { + Some(&path[1..]) + } else { + None + } + }) + .filter(|sub| !sub.is_empty()) + .collect::>(); + + let out = if sub_paths.is_empty() { + // Leaf of access path — keep the field as-is. + Arc::clone(f) + } else { + // Recurse into nested struct. + let pruned = prune_struct_type(f.data_type(), &sub_paths); + Arc::new(Field::new(f.name(), pruned, f.is_nullable())) + }; + + Some(out) + }) + .collect::>(); + + DataType::Struct(pruned_fields.into()) } /// Checks if a predicate expression can be pushed down to the parquet decoder. @@ -394,6 +1024,126 @@ fn size_of_columns(columns: &[usize], metadata: &ParquetMetaData) -> Result, + /// Projection mask over the parquet leaf columns needed to evaluate this + /// predicate. + projection_mask: ProjectionMask, + /// Precomputed sum-of-compressed-bytes for the referenced columns across + /// all row groups in the file. Used to sort predicates when + /// `reorder_predicates` is enabled. Stable across row groups within a + /// file, so we cache it once. + required_bytes: usize, +} + +/// Precompute the list of [`PrebuiltRowFilterCandidate`]s for a predicate. +/// +/// This is the expensive part of [`build_row_filter`]: split into conjuncts, +/// resolve columns for each conjunct against the file schema, reassign +/// `Column` indices, and compute the sort-order metadata. Doing it once per +/// file (and reusing across row groups) avoids repeated `TreeNode::transform` +/// walks and `Arc` allocations that showed up as top hot spots +/// in TPCH profiles. +/// +/// Returns `Ok(None)` when the predicate has no push-downable conjuncts, in +/// which case callers should skip installing a `RowFilter` entirely. +pub(crate) fn prebuild_row_filter_candidates( + expr: &Arc, + file_schema: &SchemaRef, + metadata: &ParquetMetaData, +) -> Result>> { + // Split into conjuncts: + // `a = 1 AND b = 2 AND c = 3` -> [`a = 1`, `b = 2`, `c = 3`] + let predicates = split_conjunction(expr); + let candidates: Vec = predicates + .into_iter() + .map(|expr| { + FilterCandidateBuilder::new(Arc::clone(expr), Arc::clone(file_schema)) + .build(metadata) + }) + .collect::, _>>()? + .into_iter() + .flatten() + .collect(); + + if candidates.is_empty() { + return Ok(None); + } + + let prebuilt: Vec = candidates + .into_iter() + .map(|candidate| { + let physical_expr = reassign_expr_columns( + Arc::clone(&candidate.expr), + &candidate.read_plan.projected_schema, + )?; + Ok(PrebuiltRowFilterCandidate { + physical_expr, + projection_mask: candidate.read_plan.projection_mask.clone(), + required_bytes: candidate.required_bytes, + }) + }) + .collect::>>()?; + + Ok(Some(prebuilt)) +} + +/// Wrap a list of prebuilt candidates into a fresh [`RowFilter`], assigning +/// per-predicate metric counters and (optionally) reordering by +/// `required_bytes`. This is the cheap per-row-group rebuild path — no tree +/// walks, no column resolution, only counter allocation. +pub(crate) fn row_filter_from_prebuilt( + prebuilt: &[PrebuiltRowFilterCandidate], + reorder_predicates: bool, + file_metrics: &ParquetFileMetrics, +) -> RowFilter { + let rows_pruned = &file_metrics.pushdown_rows_pruned; + let rows_matched = &file_metrics.pushdown_rows_matched; + let time = &file_metrics.row_pushdown_eval_time; + + // Clone (cheap: Arc bumps + ProjectionMask clone) into a working list we + // can sort without disturbing the shared cache. + let mut ordered: Vec<&PrebuiltRowFilterCandidate> = prebuilt.iter().collect(); + if reorder_predicates { + ordered.sort_unstable_by_key(|c| c.required_bytes); + } + + let total = ordered.len(); + let filters: Vec> = ordered + .into_iter() + .enumerate() + .map(|(idx, candidate)| { + let is_last = idx == total - 1; + let predicate_rows_pruned = rows_pruned.clone(); + let predicate_rows_matched = if is_last { + rows_matched.clone() + } else { + metrics::Count::new() + }; + Box::new(DatafusionArrowPredicate { + physical_expr: Arc::clone(&candidate.physical_expr), + projection_mask: candidate.projection_mask.clone(), + rows_pruned: predicate_rows_pruned, + rows_matched: predicate_rows_matched, + time: time.clone(), + }) as Box + }) + .collect(); + RowFilter::new(filters) +} + pub fn build_row_filter( expr: &Arc, file_schema: &SchemaRef, @@ -464,69 +1214,10 @@ pub fn build_row_filter( .map(|filters| Some(RowFilter::new(filters))) } -/// Builds row filters for a parquet decoder. -/// -/// A [`RowFilter`] is owned by a decoder. The first filter is built eagerly -/// during construction so the caller can attach it to the decoder via -/// [`next_filter`](Self::next_filter) without a redundant build call. -pub(crate) struct RowFilterGenerator<'a> { - predicate: Option<&'a Arc>, - physical_file_schema: &'a SchemaRef, - file_metadata: &'a ParquetMetaData, - reorder_predicates: bool, - file_metrics: &'a ParquetFileMetrics, - first_row_filter: Option, -} - -impl<'a> RowFilterGenerator<'a> { - pub(crate) fn new( - predicate: Option<&'a Arc>, - physical_file_schema: &'a SchemaRef, - file_metadata: &'a ParquetMetaData, - reorder_predicates: bool, - file_metrics: &'a ParquetFileMetrics, - ) -> Self { - let mut generator = Self { - predicate, - physical_file_schema, - file_metadata, - reorder_predicates, - file_metrics, - first_row_filter: None, - }; - generator.first_row_filter = generator.build(); - generator - } - - pub(crate) fn next_filter(&mut self) -> Option { - self.first_row_filter.take().or_else(|| self.build()) - } - - fn build(&self) -> Option { - let predicate = self.predicate?; - match build_row_filter( - predicate, - self.physical_file_schema, - self.file_metadata, - self.reorder_predicates, - self.file_metrics, - ) { - Ok(Some(filter)) => Some(filter), - Ok(None) => None, - Err(e) => { - log::debug!( - "Ignoring error building row filter for '{predicate:?}': {e}" - ); - None - } - } - } -} - #[cfg(test)] mod test { use super::*; - use arrow::datatypes::{DataType, Fields}; + use arrow::datatypes::Fields; use datafusion_common::ScalarValue; use arrow::array::{ @@ -541,7 +1232,6 @@ mod test { use datafusion_functions_nested::expr_fn::{ array_has, array_has_all, array_has_any, make_array, }; - use datafusion_physical_expr::expressions::Column; use datafusion_physical_expr::planner::logical2physical; use datafusion_physical_expr_adapter::{ DefaultPhysicalExprAdapterFactory, PhysicalExprAdapterFactory, @@ -554,6 +1244,8 @@ mod test { use parquet::file::reader::{FileReader, SerializedFileReader}; use tempfile::NamedTempFile; + use datafusion_physical_expr::expressions::Column as PhysicalColumn; + // List predicates used by the decoder should be accepted for pushdown #[test] fn test_filter_candidate_builder_supports_list_types() { @@ -1410,6 +2102,86 @@ mod test { assert_eq!(file_metrics.pushdown_rows_matched.value(), 2); } + #[test] + fn projection_read_plan_preserves_full_struct() { + // Schema: id (Int32), s (Struct{value: Int32, label: Utf8}) + // Parquet leaves: id=0, s.value=1, s.label=2 + let struct_fields: Fields = vec![ + Arc::new(Field::new("value", DataType::Int32, false)), + Arc::new(Field::new("label", DataType::Utf8, false)), + ] + .into(); + + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("s", DataType::Struct(struct_fields.clone()), false), + ])); + + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(Int32Array::from(vec![1, 2, 3])), + Arc::new(StructArray::new( + struct_fields, + vec![ + Arc::new(Int32Array::from(vec![10, 20, 30])) as _, + Arc::new(StringArray::from(vec!["a", "b", "c"])) as _, + ], + None, + )), + ], + ) + .unwrap(); + + let file = NamedTempFile::new().expect("temp file"); + let mut writer = + ArrowWriter::try_new(file.reopen().unwrap(), Arc::clone(&schema), None) + .expect("writer"); + writer.write(&batch).expect("write batch"); + writer.close().expect("close writer"); + + let reader_file = file.reopen().expect("reopen file"); + let builder = ParquetRecordBatchReaderBuilder::try_new(reader_file) + .expect("reader builder"); + let metadata = builder.metadata().clone(); + let file_schema = builder.schema().clone(); + let schema_descr = metadata.file_metadata().schema_descr(); + + // Simulate SELECT * output projection: Column("id") and Column("s") + // Plus a get_field(s, 'value') expression from the pushed-down filter + let exprs: Vec> = vec![ + Arc::new(PhysicalColumn::new("id", 0)), + Arc::new(PhysicalColumn::new("s", 1)), + logical2physical( + &get_field().call(vec![ + col("s"), + Expr::Literal(ScalarValue::Utf8(Some("value".to_string())), None), + ]), + &file_schema, + ), + ]; + + let read_plan = build_projection_read_plan(exprs, &file_schema, schema_descr); + + // The projected schema must have the FULL struct type because Column("s") + // is in the projection. It should NOT be narrowed to Struct{value: Int32}. + let s_field = read_plan.projected_schema.field_with_name("s").unwrap(); + assert_eq!( + s_field.data_type(), + &DataType::Struct( + vec![ + Arc::new(Field::new("value", DataType::Int32, false)), + Arc::new(Field::new("label", DataType::Utf8, false)), + ] + .into() + ), + ); + + // all3 Parquet leaves should be in the projection mask + let expected_mask = ProjectionMask::leaves(schema_descr, [0, 1, 2]); + assert_eq!(read_plan.projection_mask, expected_mask,); + } + /// Sanity check that the given expression could be evaluated against the given schema without any errors. /// This will fail if the expression references columns that are not in the schema or if the types of the columns are incompatible, etc. fn check_expression_can_evaluate_against_schema( diff --git a/datafusion/datasource-parquet/src/sort.rs b/datafusion/datasource-parquet/src/sort.rs index c1cf4e8b7824e..35eab767b2a24 100644 --- a/datafusion/datasource-parquet/src/sort.rs +++ b/datafusion/datasource-parquet/src/sort.rs @@ -372,12 +372,15 @@ mod tests { #[test] fn test_prepared_access_plan_reverse_empty_selection() { - // Test: all rows are skipped + // Test: all rows are skipped in every RG. With the strip-empty + // pass in `prepare()`, every RG whose selection segment selects + // zero rows is dropped from the prepared plan entirely (arrow-rs + // would silently skip them at read time anyway, and the rest of + // DataFusion relies on the plan being 1:1 with the readers the + // decoder hands back). let metadata = create_test_metadata(vec![100, 100, 100]); let mut access_plan = ParquetAccessPlan::new_all(3); - - // Skip all rows in all row groups for i in 0..3 { access_plan .scan_selection(i, RowSelection::from(vec![RowSelector::skip(100)])); @@ -387,22 +390,19 @@ mod tests { let prepared_plan = access_plan .prepare(rg_metadata) .expect("Failed to create PreparedAccessPlan"); + assert!( + prepared_plan.row_group_indexes.is_empty(), + "all-empty selections must be stripped to an empty plan", + ); + assert!(prepared_plan.row_selection.is_none()); let reversed_plan = prepared_plan .reverse(&metadata) .expect("Failed to reverse PreparedAccessPlan"); - // Should still skip all rows - let total_selected: usize = reversed_plan - .row_selection - .as_ref() - .unwrap() - .iter() - .filter(|s| !s.skip) - .map(|s| s.row_count) - .sum(); - - assert_eq!(total_selected, 0); + // No rows survive in either direction. + assert!(reversed_plan.row_group_indexes.is_empty()); + assert!(reversed_plan.row_selection.is_none()); } #[test] diff --git a/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt b/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt index e879947e324bb..5717597e50786 100644 --- a/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt +++ b/datafusion/sqllogictest/test_files/push_down_filter_parquet.slt @@ -206,7 +206,7 @@ EXPLAIN ANALYZE SELECT t FROM topk_pushdown ORDER BY t * t LIMIT 10; ---- Plan with Metrics 01)SortExec: TopK(fetch=10), expr=[t@0 * t@0 ASC NULLS LAST], preserve_partitioning=[false], filter=[t@0 * t@0 < 1884329474306198481], metrics=[output_rows=10, output_batches=1, row_replacements=10] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_pushdown.parquet]]}, projection=[t], output_ordering=[t@0 ASC NULLS LAST], file_type=parquet, predicate=DynamicFilter [ t@0 * t@0 < 1884329474306198481 ], dynamic_rg_pruning=eligible, metrics=[output_rows=128, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=782 total → 782 matched, row_groups_pruned_bloom_filter=782 total → 782 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=128, pushdown_rows_pruned=99.87 K, predicate_cache_inner_records=128, predicate_cache_records=128, scan_efficiency_ratio=64.87% (258.7 K/398.8 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_pushdown.parquet]]}, projection=[t], output_ordering=[t@0 ASC NULLS LAST], file_type=parquet, predicate=DynamicFilter [ t@0 * t@0 < 1884329474306198481 ], dynamic_rg_pruning=eligible, metrics=[output_rows=128, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=782 total → 782 matched, row_groups_pruned_bloom_filter=782 total → 782 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=128, pushdown_rows_pruned=99.87 K, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=128, predicate_cache_records=128, scan_efficiency_ratio=64.87% (258.7 K/398.8 K)] statement ok reset datafusion.explain.analyze_categories; @@ -268,7 +268,7 @@ EXPLAIN ANALYZE SELECT * FROM topk_single_col ORDER BY b DESC LIMIT 1; ---- Plan with Metrics 01)SortExec: TopK(fetch=1), expr=[b@1 DESC], preserve_partitioning=[false], filter=[b@1 IS NULL OR b@1 > bd], metrics=[output_rows=1, output_batches=1, row_replacements=1] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_single_col.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ b@1 IS NULL OR b@1 > bd ], sort_order_for_reorder=[b@1 DESC], reverse_row_groups=true, dynamic_rg_pruning=eligible, pruning_predicate=b_null_count@0 > 0 OR b_null_count@0 != row_count@2 AND b_max@1 > bd, required_guarantees=[], metrics=[output_rows=4, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=4, pushdown_rows_pruned=0, predicate_cache_inner_records=4, predicate_cache_records=4, scan_efficiency_ratio=21.62% (222/1.03 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_single_col.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ b@1 IS NULL OR b@1 > bd ], sort_order_for_reorder=[b@1 DESC], reverse_row_groups=true, dynamic_rg_pruning=eligible, pruning_predicate=b_null_count@0 > 0 OR b_null_count@0 != row_count@2 AND b_max@1 > bd, required_guarantees=[], metrics=[output_rows=4, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=4, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=4, predicate_cache_records=4, scan_efficiency_ratio=21.62% (222/1.03 K)] statement ok reset datafusion.explain.analyze_categories; @@ -319,7 +319,7 @@ EXPLAIN ANALYZE SELECT * FROM topk_multi_col ORDER BY b ASC NULLS LAST, a DESC L ---- Plan with Metrics 01)SortExec: TopK(fetch=2), expr=[b@1 ASC NULLS LAST, a@0 DESC], preserve_partitioning=[false], filter=[b@1 < bb OR b@1 = bb AND (a@0 IS NULL OR a@0 > ac)], metrics=[output_rows=2, output_batches=1, row_replacements=2] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_multi_col.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ b@1 < bb OR b@1 = bb AND (a@0 IS NULL OR a@0 > ac) ], sort_order_for_reorder=[b@1 ASC NULLS LAST, a@0 DESC], dynamic_rg_pruning=eligible, pruning_predicate=b_null_count@1 != row_count@2 AND b_min@0 < bb OR b_null_count@1 != row_count@2 AND b_min@0 <= bb AND bb <= b_max@3 AND (a_null_count@4 > 0 OR a_null_count@4 != row_count@2 AND a_max@5 > ac), required_guarantees=[], metrics=[output_rows=4, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=4, pushdown_rows_pruned=0, predicate_cache_inner_records=8, predicate_cache_records=8, scan_efficiency_ratio=21.62% (222/1.03 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_multi_col.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ b@1 < bb OR b@1 = bb AND (a@0 IS NULL OR a@0 > ac) ], sort_order_for_reorder=[b@1 ASC NULLS LAST, a@0 DESC], dynamic_rg_pruning=eligible, pruning_predicate=b_null_count@1 != row_count@2 AND b_min@0 < bb OR b_null_count@1 != row_count@2 AND b_min@0 <= bb AND bb <= b_max@3 AND (a_null_count@4 > 0 OR a_null_count@4 != row_count@2 AND a_max@5 > ac), required_guarantees=[], metrics=[output_rows=4, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=4, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=8, scan_efficiency_ratio=21.62% (222/1.03 K)] statement ok reset datafusion.explain.analyze_categories; @@ -388,8 +388,8 @@ FROM join_probe p INNER JOIN join_build AS build ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, a@0), (b@1, b@1)], projection=[a@3, b@4, c@2, e@5], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] statement ok reset datafusion.explain.analyze_categories; @@ -474,9 +474,9 @@ INNER JOIN nested_t3 ON nested_t2.c = nested_t3.d; Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(c@3, d@0)], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] 02)--HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, b@0)], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nested_t1.parquet]]}, projection=[a, x], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=17.37% (132/760)] -04)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nested_t2.parquet]]}, projection=[b, c, y], file_type=parquet, predicate=DynamicFilter [ b@0 >= aa AND b@0 <= ab AND b@0 IN (SET) ([aa, ab]) ], dynamic_rg_pruning=eligible, pruning_predicate=b_null_count@1 != row_count@2 AND b_max@0 >= aa AND b_null_count@1 != row_count@2 AND b_min@3 <= ab AND (b_null_count@1 != row_count@2 AND b_min@3 <= aa AND aa <= b_max@0 OR b_null_count@1 != row_count@2 AND b_min@3 <= ab AND ab <= b_max@0), required_guarantees=[b in (aa, ab)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=5 total → 5 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=3, predicate_cache_inner_records=5, predicate_cache_records=2, scan_efficiency_ratio=22.46% (234/1.04 K)] -05)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nested_t3.parquet]]}, projection=[d, z], file_type=parquet, predicate=DynamicFilter [ d@0 >= ca AND d@0 <= cb AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=d_null_count@1 != row_count@2 AND d_max@0 >= ca AND d_null_count@1 != row_count@2 AND d_min@3 <= cb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=8 total → 8 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=6, predicate_cache_inner_records=8, predicate_cache_records=2, scan_efficiency_ratio=21.45% (172/802)] +03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nested_t1.parquet]]}, projection=[a, x], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=17.37% (132/760)] +04)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nested_t2.parquet]]}, projection=[b, c, y], file_type=parquet, predicate=DynamicFilter [ b@0 >= aa AND b@0 <= ab AND b@0 IN (SET) ([aa, ab]) ], dynamic_rg_pruning=eligible, pruning_predicate=b_null_count@1 != row_count@2 AND b_max@0 >= aa AND b_null_count@1 != row_count@2 AND b_min@3 <= ab AND (b_null_count@1 != row_count@2 AND b_min@3 <= aa AND aa <= b_max@0 OR b_null_count@1 != row_count@2 AND b_min@3 <= ab AND ab <= b_max@0), required_guarantees=[b in (aa, ab)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=5 total → 5 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=3, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=5, predicate_cache_records=2, scan_efficiency_ratio=22.46% (234/1.04 K)] +05)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nested_t3.parquet]]}, projection=[d, z], file_type=parquet, predicate=DynamicFilter [ d@0 >= ca AND d@0 <= cb AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=d_null_count@1 != row_count@2 AND d_max@0 >= ca AND d_null_count@1 != row_count@2 AND d_min@3 <= cb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=8 total → 8 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=6, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=2, scan_efficiency_ratio=21.45% (172/802)] statement ok reset datafusion.explain.analyze_categories; @@ -605,8 +605,8 @@ LIMIT 2; Plan with Metrics 01)SortExec: TopK(fetch=2), expr=[e@0 ASC NULLS LAST], preserve_partitioning=[false], filter=[e@0 < bb], metrics=[output_rows=2, output_batches=1, row_replacements=2] 02)--HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, d@0)], projection=[e@2], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_join_build.parquet]]}, projection=[a], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=6.39% (64/1.00 K)] -04)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_join_probe.parquet]]}, projection=[d, e], file_type=parquet, predicate=DynamicFilter [ d@0 >= aa AND d@0 <= ab AND d@0 IN (SET) ([aa, ab]) ] AND DynamicFilter [ e@1 < bb ], dynamic_rg_pruning=eligible, pruning_predicate=d_null_count@1 != row_count@2 AND d_max@0 >= aa AND d_null_count@1 != row_count@2 AND d_min@3 <= ab AND (d_null_count@1 != row_count@2 AND d_min@3 <= aa AND aa <= d_max@0 OR d_null_count@1 != row_count@2 AND d_min@3 <= ab AND ab <= d_max@0) AND e_null_count@5 != row_count@2 AND e_min@4 < bb, required_guarantees=[d in (aa, ab)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=4 total → 4 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=14.89% (154/1.03 K)] +03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_join_build.parquet]]}, projection=[a], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=6.39% (64/1.00 K)] +04)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_join_probe.parquet]]}, projection=[d, e], file_type=parquet, predicate=DynamicFilter [ d@0 >= aa AND d@0 <= ab AND d@0 IN (SET) ([aa, ab]) ] AND DynamicFilter [ e@1 < bb ], dynamic_rg_pruning=eligible, pruning_predicate=d_null_count@1 != row_count@2 AND d_max@0 >= aa AND d_null_count@1 != row_count@2 AND d_min@3 <= ab AND (d_null_count@1 != row_count@2 AND d_min@3 <= aa AND aa <= d_max@0 OR d_null_count@1 != row_count@2 AND d_min@3 <= ab AND ab <= d_max@0) AND e_null_count@5 != row_count@2 AND e_min@4 < bb, required_guarantees=[d in (aa, ab)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=4 total → 4 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=14.89% (154/1.03 K)] statement ok reset datafusion.explain.analyze_categories; @@ -656,7 +656,7 @@ EXPLAIN ANALYZE SELECT b, a FROM topk_proj ORDER BY a LIMIT 2; Plan with Metrics 01)ProjectionExec: expr=[b@1 as b, a@0 as a], metrics=[output_rows=2, output_batches=1] 02)--SortExec: TopK(fetch=2), expr=[a@0 ASC NULLS LAST], preserve_partitioning=[false], filter=[a@0 < 2], metrics=[output_rows=2, output_batches=1, row_replacements=2] -03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[a, b], file_type=parquet, predicate=DynamicFilter [ a@0 < 2 ], sort_order_for_reorder=[a@0 ASC NULLS LAST], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_min@0 < 2, required_guarantees=[], metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=13.21% (141/1.07 K)] +03)----DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[a, b], file_type=parquet, predicate=DynamicFilter [ a@0 < 2 ], sort_order_for_reorder=[a@0 ASC NULLS LAST], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_min@0 < 2, required_guarantees=[], metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=13.21% (141/1.07 K)] # Case 2: prune — `SELECT a` — filter stays as `a < 2` on the scan. query TT @@ -664,7 +664,7 @@ EXPLAIN ANALYZE SELECT a FROM topk_proj ORDER BY a LIMIT 2; ---- Plan with Metrics 01)SortExec: TopK(fetch=2), expr=[a@0 ASC NULLS LAST], preserve_partitioning=[false], filter=[a@0 < 2], metrics=[output_rows=2, output_batches=1, row_replacements=2] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[a], file_type=parquet, predicate=DynamicFilter [ a@0 < 2 ], sort_order_for_reorder=[a@0 ASC NULLS LAST], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_min@0 < 2, required_guarantees=[], metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=6.84% (73/1.07 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[a], file_type=parquet, predicate=DynamicFilter [ a@0 < 2 ], sort_order_for_reorder=[a@0 ASC NULLS LAST], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_min@0 < 2, required_guarantees=[], metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=6.84% (73/1.07 K)] # Case 3: expression — `SELECT a+1 AS a_plus_1` — the TopK filter is on # `a_plus_1`, the scan predicate must read `a@0 + 1`. @@ -673,7 +673,7 @@ EXPLAIN ANALYZE SELECT a + 1 AS a_plus_1, b FROM topk_proj ORDER BY a_plus_1 LIM ---- Plan with Metrics 01)SortExec: TopK(fetch=2), expr=[a_plus_1@0 ASC NULLS LAST], preserve_partitioning=[false], filter=[a_plus_1@0 < 3], metrics=[output_rows=2, output_batches=1, row_replacements=2] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[CAST(a@0 AS Int64) + 1 as a_plus_1, b], file_type=parquet, predicate=DynamicFilter [ CAST(a@0 AS Int64) + 1 < 3 ], dynamic_rg_pruning=eligible, metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=13.21% (141/1.07 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[CAST(a@0 AS Int64) + 1 as a_plus_1, b], file_type=parquet, predicate=DynamicFilter [ CAST(a@0 AS Int64) + 1 < 3 ], dynamic_rg_pruning=eligible, metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=13.21% (141/1.07 K)] # Case 4: alias shadowing — `SELECT a+1 AS a` — the projection renames # `a+1` to `a`, so the TopK's `a < 3` must still be rewritten to @@ -683,7 +683,7 @@ EXPLAIN ANALYZE SELECT a + 1 AS a, b FROM topk_proj ORDER BY a LIMIT 2; ---- Plan with Metrics 01)SortExec: TopK(fetch=2), expr=[a@0 ASC NULLS LAST], preserve_partitioning=[false], filter=[a@0 < 3], metrics=[output_rows=2, output_batches=1, row_replacements=2] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[CAST(a@0 AS Int64) + 1 as a, b], file_type=parquet, predicate=DynamicFilter [ CAST(a@0 AS Int64) + 1 < 3 ], sort_order_for_reorder=[a@0 ASC NULLS LAST], dynamic_rg_pruning=eligible, metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=13.21% (141/1.07 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/topk_proj.parquet]]}, projection=[CAST(a@0 AS Int64) + 1 as a, b], file_type=parquet, predicate=DynamicFilter [ CAST(a@0 AS Int64) + 1 < 3 ], sort_order_for_reorder=[a@0 ASC NULLS LAST], dynamic_rg_pruning=eligible, metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=3, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=3, predicate_cache_records=3, scan_efficiency_ratio=13.21% (141/1.07 K)] statement ok reset datafusion.explain.analyze_categories; @@ -740,12 +740,12 @@ INNER JOIN ( ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, a@0)], projection=[a@0, min_value@2], metrics=[output_rows=2, output_batches=2, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=2, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_agg_build.parquet]]}, projection=[a], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=14.45% (64/443)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_agg_build.parquet]]}, projection=[a], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=14.45% (64/443)] 03)--ProjectionExec: expr=[a@0 as a, min(join_agg_probe.value)@1 as min_value], metrics=[output_rows=2, output_batches=2] 04)----AggregateExec: mode=FinalPartitioned, gby=[a@0 as a], aggr=[min(join_agg_probe.value)], metrics=[output_rows=2, output_batches=2, spill_count=0, spilled_rows=0] 05)------RepartitionExec: partitioning=Hash([a@0], 4), input_partitions=1, metrics=[output_rows=2, output_batches=2, spill_count=0, spilled_rows=0] 06)--------AggregateExec: mode=Partial, gby=[a@0 as a], aggr=[min(join_agg_probe.value)], metrics=[output_rows=2, output_batches=1, spill_count=0, spilled_rows=0, skipped_aggregation_rows=0, reduction_factor=100% (2/2)] -07)----------DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_agg_probe.parquet]]}, projection=[a, value], file_type=parquet, predicate=DynamicFilter [ a@0 >= h1 AND a@0 <= h2 AND a@0 IN (SET) ([h1, h2]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= h1 AND a_null_count@1 != row_count@2 AND a_min@3 <= h2 AND (a_null_count@1 != row_count@2 AND a_min@3 <= h1 AND h1 <= a_max@0 OR a_null_count@1 != row_count@2 AND a_min@3 <= h2 AND h2 <= a_max@0), required_guarantees=[a in (h1, h2)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=4 total → 4 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=4, predicate_cache_records=2, scan_efficiency_ratio=19.07% (151/792)] +07)----------DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/join_agg_probe.parquet]]}, projection=[a, value], file_type=parquet, predicate=DynamicFilter [ a@0 >= h1 AND a@0 <= h2 AND a@0 IN (SET) ([h1, h2]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= h1 AND a_null_count@1 != row_count@2 AND a_min@3 <= h2 AND (a_null_count@1 != row_count@2 AND a_min@3 <= h1 AND h1 <= a_max@0 OR a_null_count@1 != row_count@2 AND a_min@3 <= h2 AND h2 <= a_max@0), required_guarantees=[a in (h1, h2)], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=1 total → 1 matched, page_index_rows_pruned=4 total → 4 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=4, predicate_cache_records=2, scan_efficiency_ratio=19.07% (151/792)] statement ok reset datafusion.explain.analyze_categories; @@ -807,8 +807,8 @@ ON nulls_build.a = nulls_probe.a AND nulls_build.b = nulls_probe.b; ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, a@0), (b@1, b@1)], metrics=[output_rows=1, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=3, input_batches=1, input_rows=1, avg_fanout=100% (1/1), probe_hit_rate=100% (1/1)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nulls_build.parquet]]}, projection=[a, b], file_type=parquet, metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=18.6% (144/774)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nulls_probe.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= 1 AND b@1 <= 2 AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:1}, {c0:,c1:2}, {c0:ab,c1:}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= 1 AND b_null_count@5 != row_count@2 AND b_min@6 <= 2, required_guarantees=[], metrics=[output_rows=1, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=1, pushdown_rows_pruned=3, predicate_cache_inner_records=8, predicate_cache_records=2, scan_efficiency_ratio=20.18% (225/1.11 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nulls_build.parquet]]}, projection=[a, b], file_type=parquet, metrics=[output_rows=3, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=18.6% (144/774)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/nulls_probe.parquet]]}, projection=[a, b, c], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= 1 AND b@1 <= 2 AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:1}, {c0:,c1:2}, {c0:ab,c1:}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= 1 AND b_null_count@5 != row_count@2 AND b_min@6 <= 2, required_guarantees=[], metrics=[output_rows=1, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=1, pushdown_rows_pruned=3, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=2, scan_efficiency_ratio=20.18% (225/1.11 K)] statement ok reset datafusion.explain.analyze_categories; @@ -873,8 +873,8 @@ ON lj_build.a = lj_probe.a AND lj_build.b = lj_probe.b; ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Left, on=[(a@0, a@0), (b@1, b@1)], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=2, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] # LEFT SEMI JOIN: only matching build rows are returned; probe scan still # receives the dynamic filter. @@ -889,8 +889,8 @@ WHERE EXISTS ( ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=LeftSemi, on=[(a@0, a@0), (b@1, b@1)], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=2, input_rows=4, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=14.89% (154/1.03 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/lj_probe.parquet]]}, projection=[a, b], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND struct(a@0, b@1) IN (SET) ([{c0:aa,c1:ba}, {c0:ab,c1:bb}]) ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=14.89% (154/1.03 K)] statement ok reset datafusion.explain.analyze_categories; @@ -959,8 +959,8 @@ FROM hl_probe p INNER JOIN hl_build AS build ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(a@0, a@0), (b@1, b@1)], projection=[a@3, b@4, c@2, e@5], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/hl_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/hl_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/hl_build.parquet]]}, projection=[a, b, c], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=19.58% (196/1.00 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/hl_probe.parquet]]}, projection=[a, b, e], file_type=parquet, predicate=DynamicFilter [ a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=a_null_count@1 != row_count@2 AND a_max@0 >= aa AND a_null_count@1 != row_count@2 AND a_min@3 <= ab AND b_null_count@5 != row_count@2 AND b_max@4 >= ba AND b_null_count@5 != row_count@2 AND b_min@6 <= bb, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=22.05% (228/1.03 K)] statement ok drop table hl_build; @@ -1008,8 +1008,8 @@ FROM int_build b INNER JOIN int_probe p ---- Plan with Metrics 01)HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(id1@0, id1@0), (id2@1, id2@1)], projection=[id1@0, id2@1, value@2, data@5], metrics=[output_rows=2, output_batches=1, array_map_created_count=0, build_input_batches=1, build_input_rows=2, input_batches=1, input_rows=2, avg_fanout=100% (2/2), probe_hit_rate=100% (2/2)] -02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/int_build.parquet]]}, projection=[id1, id2, value], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=18.23% (204/1.12 K)] -03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/int_probe.parquet]]}, projection=[id1, id2, data], file_type=parquet, predicate=DynamicFilter [ id1@0 >= 1 AND id1@0 <= 2 AND id2@1 >= 10 AND id2@1 <= 20 AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=id1_null_count@1 != row_count@2 AND id1_max@0 >= 1 AND id1_null_count@1 != row_count@2 AND id1_min@3 <= 2 AND id2_null_count@5 != row_count@2 AND id2_max@4 >= 10 AND id2_null_count@5 != row_count@2 AND id2_min@6 <= 20, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=20.67% (221/1.07 K)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/int_build.parquet]]}, projection=[id1, id2, value], file_type=parquet, metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=0, predicate_cache_records=0, scan_efficiency_ratio=18.23% (204/1.12 K)] +03)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/push_down_filter_parquet/int_probe.parquet]]}, projection=[id1, id2, data], file_type=parquet, predicate=DynamicFilter [ id1@0 >= 1 AND id1@0 <= 2 AND id2@1 >= 10 AND id2@1 <= 20 AND hash_lookup ], dynamic_rg_pruning=eligible, pruning_predicate=id1_null_count@1 != row_count@2 AND id1_max@0 >= 1 AND id1_null_count@1 != row_count@2 AND id1_min@3 <= 2 AND id2_null_count@5 != row_count@2 AND id2_max@4 >= 10 AND id2_null_count@5 != row_count@2 AND id2_min@6 <= 20, required_guarantees=[], metrics=[output_rows=2, output_batches=1, files_ranges_pruned_statistics=1 total → 1 matched, row_groups_pruned_statistics=1 total → 1 matched, row_groups_pruned_bloom_filter=1 total → 1 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, file_open_errors=0, file_scan_errors=0, files_opened=1, files_processed=1, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=2, pushdown_rows_pruned=2, row_filter_skipped_fully_matched=0, predicate_cache_inner_records=8, predicate_cache_records=4, scan_efficiency_ratio=20.67% (221/1.07 K)] statement ok reset datafusion.explain.analyze_categories;