diff --git a/datafusion/core/tests/parquet/filter_pushdown.rs b/datafusion/core/tests/parquet/filter_pushdown.rs index d337979e5fd00..2039e2a72b259 100644 --- a/datafusion/core/tests/parquet/filter_pushdown.rs +++ b/datafusion/core/tests/parquet/filter_pushdown.rs @@ -28,10 +28,19 @@ use arrow::array::{ArrayRef, Int32Array, StringArray}; use arrow::compute::concat_batches; +use arrow::datatypes::SchemaRef; use arrow::error::ArrowError; use arrow::record_batch::RecordBatch; +use datafusion::datasource::listing::PartitionedFile; +use datafusion::datasource::object_store::ObjectStoreUrl; +use datafusion::datasource::physical_plan::ParquetSource; +use datafusion::datasource::source::DataSourceExec; +use datafusion::physical_expr::expressions::col as physical_col; +use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr}; +use datafusion::physical_plan::filter::FilterExec; use datafusion::physical_plan::metrics::{MetricValue, MetricsSet}; -use datafusion::physical_plan::{collect, displayable}; +use datafusion::physical_plan::{ExecutionPlan, collect, displayable}; +use datafusion::physical_planner::DefaultPhysicalPlanner; use datafusion::prelude::{ Expr, ParquetReadOptions, SessionContext, col, lit, lit_timestamp_nano, }; @@ -39,9 +48,12 @@ use datafusion::test_util::parquet::{ParquetScanOptions, TestParquetFile}; use datafusion_expr::utils::{conjunction, disjunction, split_conjunction}; use std::path::Path; -use datafusion_common::DataFusionError; use datafusion_common::test_util::parquet_test_data; +use datafusion_common::{DataFusionError, ToDFSchema}; +use datafusion_datasource::file_scan_config::FileScanConfigBuilder; use datafusion_execution::config::SessionConfig; +use datafusion_physical_optimizer::PhysicalOptimizerRule; +use datafusion_physical_optimizer::filter_pushdown::FilterPushdown; use itertools::Itertools; use parquet::arrow::ArrowWriter; use parquet::file::properties::WriterProperties; @@ -62,6 +74,283 @@ async fn read_parquet_test_data>(path: T) -> Vec { .unwrap() } +fn small_parquet_file( + tempdir: &TempDir, + schema: &SchemaRef, + name: &str, + row_groups: &[&[i32]], + row_group_size: Option, +) -> TestParquetFile { + let mut properties = WriterProperties::builder(); + if let Some(row_group_size) = row_group_size { + properties = properties.set_max_row_group_row_count(Some(row_group_size)); + } + TestParquetFile::try_new( + tempdir.path().join(name), + properties.build(), + row_groups.iter().map(|values| { + RecordBatch::try_new( + Arc::clone(schema), + vec![Arc::new(Int32Array::from(values.to_vec())) as ArrayRef], + ) + .unwrap() + }), + ) + .unwrap() +} + +fn ordered_scan_with_fetch( + file: &TestParquetFile, + schema: &SchemaRef, + source: ParquetSource, + fetch: Option, +) -> Arc { + let path = file.path().canonicalize().unwrap(); + let ordering = LexOrdering::new(vec![ + PhysicalSortExpr::new_default(physical_col("a", schema).unwrap()).asc(), + ]) + .unwrap(); + let config = + FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), Arc::new(source)) + .with_file(PartitionedFile::new( + path.to_string_lossy(), + path.metadata().unwrap().len(), + )) + .with_limit(fetch) + .with_output_ordering(vec![ordering]) + .build(); + DataSourceExec::from_data_source(config) +} + +fn physical_predicate( + ctx: &SessionContext, + schema: &SchemaRef, + value: i32, +) -> Arc { + ctx.create_physical_expr( + col("a").eq(lit(value)), + &Arc::clone(schema).to_dfschema().unwrap(), + ) + .unwrap() +} + +async fn int_values(plan: Arc, ctx: &SessionContext) -> Vec { + collect(plan, ctx.task_ctx()) + .await + .unwrap() + .into_iter() + .flat_map(|batch| { + batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .iter() + .copied() + .collect::>() + }) + .collect() +} + +#[tokio::test] +async fn filter_pushdown_preserves_fetch_with_row_filtering() { + let schema = Arc::new(arrow::datatypes::Schema::new(vec![ + arrow::datatypes::Field::new("a", arrow::datatypes::DataType::Int32, false), + ])); + let tempdir = TempDir::new_in(Path::new(".")).unwrap(); + // One row group makes this a row-filtering regression, not a pruning one. + let file = + small_parquet_file(&tempdir, &schema, "one_group.parquet", &[&[0, 1]], None); + let mut config = SessionConfig::new() + .with_target_partitions(1) + .with_parquet_pruning(false); + config.options_mut().execution.parquet.pushdown_filters = true; + let ctx = SessionContext::new_with_config(config); + + let original: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&ctx, &schema, 1), + ordered_scan_with_fetch( + &file, + &schema, + ParquetSource::new(Arc::clone(&schema)).with_pushdown_filters(true), + Some(1), + ), + ) + .unwrap(), + ); + assert_eq!( + int_values(Arc::clone(&original), &ctx).await, + Vec::::new() + ); + + let zero_fetch: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&ctx, &schema, 1), + ordered_scan_with_fetch( + &file, + &schema, + ParquetSource::new(Arc::clone(&schema)).with_pushdown_filters(true), + Some(0), + ), + ) + .unwrap(), + ); + assert_eq!( + int_values(Arc::clone(&zero_fetch), &ctx).await, + Vec::::new() + ); + let zero_pushed = FilterPushdown::new() + .optimize(zero_fetch, ctx.state().config_options()) + .unwrap(); + assert_eq!( + int_values(Arc::clone(&zero_pushed), &ctx).await, + Vec::::new() + ); + let zero_plan = displayable(zero_pushed.as_ref()).indent(false).to_string(); + assert!( + zero_plan.contains("FilterExec") && !zero_plan.contains("predicate="), + "{zero_plan}" + ); + + let pushed = FilterPushdown::new() + .optimize(Arc::clone(&original), ctx.state().config_options()) + .unwrap(); + assert_eq!( + int_values(Arc::clone(&pushed), &ctx).await, + Vec::::new() + ); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!( + plan.contains("FilterExec") && !plan.contains("predicate="), + "{plan}" + ); + + let planner = DefaultPhysicalPlanner::default(); + let once = planner + .optimize_physical_plan(Arc::clone(&original), &ctx.state(), |_, _| {}) + .unwrap(); + assert_eq!(int_values(Arc::clone(&once), &ctx).await, Vec::::new()); + let twice = planner + .optimize_physical_plan(once, &ctx.state(), |_, _| {}) + .unwrap(); + assert_eq!(int_values(twice, &ctx).await, Vec::::new()); + + let uncapped: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&ctx, &schema, 1), + ordered_scan_with_fetch( + &file, + &schema, + ParquetSource::new(Arc::clone(&schema)), + None, + ), + ) + .unwrap(), + ); + let pushed = FilterPushdown::new() + .optimize(uncapped, ctx.state().config_options()) + .unwrap(); + assert_eq!(int_values(Arc::clone(&pushed), &ctx).await, vec![1]); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!( + !plan.contains("FilterExec") && plan.contains("predicate=a@0 = 1"), + "{plan}" + ); +} + +#[tokio::test] +async fn filter_pushdown_preserves_fetch_with_pruning() { + let schema = Arc::new(arrow::datatypes::Schema::new(vec![ + arrow::datatypes::Field::new("a", arrow::datatypes::DataType::Int32, false), + ])); + let tempdir = TempDir::new_in(Path::new(".")).unwrap(); + // Separate row groups let statistics pruning change the row admitted by fetch. + let pruning_file = small_parquet_file( + &tempdir, + &schema, + "two_groups.parquet", + &[&[0], &[1]], + Some(1), + ); + let ctx = SessionContext::new_with_config( + SessionConfig::new() + .with_target_partitions(1) + .with_parquet_pruning(true), + ); + + let original: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&ctx, &schema, 1), + ordered_scan_with_fetch( + &pruning_file, + &schema, + ParquetSource::new(Arc::clone(&schema)).with_pushdown_filters(false), + Some(1), + ), + ) + .unwrap(), + ); + assert_eq!( + int_values(Arc::clone(&original), &ctx).await, + Vec::::new() + ); + let pushed = FilterPushdown::new() + .optimize(Arc::clone(&original), ctx.state().config_options()) + .unwrap(); + assert_eq!( + int_values(Arc::clone(&pushed), &ctx).await, + Vec::::new() + ); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!( + plan.contains("FilterExec") && !plan.contains("predicate="), + "{plan}" + ); + + let internal_file = small_parquet_file( + &tempdir, + &schema, + "internal.parquet", + &[&[0], &[1], &[2]], + Some(1), + ); + let internal_predicate = ctx + .create_physical_expr( + col("a").gt_eq(lit(1_i32)), + &Arc::clone(&schema).to_dfschema().unwrap(), + ) + .unwrap(); + let internal_scan = ordered_scan_with_fetch( + &internal_file, + &schema, + ParquetSource::new(Arc::clone(&schema)) + .with_pushdown_filters(false) + .with_predicate(internal_predicate), + Some(1), + ); + assert_eq!(int_values(Arc::clone(&internal_scan), &ctx).await, vec![1]); + let original: Arc = Arc::new( + FilterExec::try_new(physical_predicate(&ctx, &schema, 2), internal_scan).unwrap(), + ); + assert_eq!( + int_values(Arc::clone(&original), &ctx).await, + Vec::::new() + ); + let pushed = FilterPushdown::new() + .optimize(original, ctx.state().config_options()) + .unwrap(); + assert_eq!( + int_values(Arc::clone(&pushed), &ctx).await, + Vec::::new() + ); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!(plan.contains("FilterExec"), "{plan}"); + assert!(plan.contains("predicate=a@0 >= 1"), "{plan}"); + assert!(!plan.contains("predicate=a@0 = 2"), "{plan}"); +} + #[tokio::test] async fn single_file() { let batches = diff --git a/datafusion/datasource/src/file_scan_config/mod.rs b/datafusion/datasource/src/file_scan_config/mod.rs index 4e72b5e83bd23..6fbc0c2192209 100644 --- a/datafusion/datasource/src/file_scan_config/mod.rs +++ b/datafusion/datasource/src/file_scan_config/mod.rs @@ -62,7 +62,7 @@ use datafusion_physical_plan::execution_plan::SchedulingType; use datafusion_physical_plan::{ DisplayAs, DisplayFormatType, display::{ProjectSchemaDisplay, display_orderings}, - filter_pushdown::FilterPushdownPropagation, + filter_pushdown::{FilterPushdownPropagation, PushedDown}, metrics::ExecutionPlanMetricsSet, }; use log::{debug, warn}; @@ -1028,6 +1028,13 @@ impl DataSource for FileScanConfig { filters: Vec>, config: &ConfigOptions, ) -> Result>> { + // A new filter or pruning predicate can change which rows reach the enforced scan cap. + if self.limit.is_some() { + return Ok(FilterPushdownPropagation::with_parent_pushdown_result( + vec![PushedDown::No; filters.len()], + )); + } + // Remap filter Column indices to match the table schema (file + partition columns). // This is necessary because filters refer to the output schema of this `DataSource` // (e.g., after projection pushdown has been applied) and need to be remapped to the table schema