Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
293 changes: 291 additions & 2 deletions datafusion/core/tests/parquet/filter_pushdown.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,20 +28,32 @@

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,
};
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;
Expand All @@ -62,6 +74,283 @@ async fn read_parquet_test_data<T: Into<String>>(path: T) -> Vec<RecordBatch> {
.unwrap()
}

fn small_parquet_file(
tempdir: &TempDir,
schema: &SchemaRef,
name: &str,
row_groups: &[&[i32]],
row_group_size: Option<usize>,
) -> 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<usize>,
) -> Arc<dyn ExecutionPlan> {
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<dyn datafusion::physical_expr::PhysicalExpr> {
ctx.create_physical_expr(
col("a").eq(lit(value)),
&Arc::clone(schema).to_dfschema().unwrap(),
)
.unwrap()
}

async fn int_values(plan: Arc<dyn ExecutionPlan>, ctx: &SessionContext) -> Vec<i32> {
collect(plan, ctx.task_ctx())
.await
.unwrap()
.into_iter()
.flat_map(|batch| {
batch
.column(0)
.as_any()
.downcast_ref::<Int32Array>()
.unwrap()
.values()
.iter()
.copied()
.collect::<Vec<_>>()
})
.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<dyn ExecutionPlan> = 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::<i32>::new()
);

let zero_fetch: Arc<dyn ExecutionPlan> = 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::<i32>::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::<i32>::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::<i32>::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::<i32>::new());
let twice = planner
.optimize_physical_plan(once, &ctx.state(), |_, _| {})
.unwrap();
assert_eq!(int_values(twice, &ctx).await, Vec::<i32>::new());

let uncapped: Arc<dyn ExecutionPlan> = 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<dyn ExecutionPlan> = 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::<i32>::new()
);
let pushed = FilterPushdown::new()
.optimize(Arc::clone(&original), ctx.state().config_options())
.unwrap();
assert_eq!(
int_values(Arc::clone(&pushed), &ctx).await,
Vec::<i32>::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<dyn ExecutionPlan> = Arc::new(
FilterExec::try_new(physical_predicate(&ctx, &schema, 2), internal_scan).unwrap(),
);
assert_eq!(
int_values(Arc::clone(&original), &ctx).await,
Vec::<i32>::new()
);
let pushed = FilterPushdown::new()
.optimize(original, ctx.state().config_options())
.unwrap();
assert_eq!(
int_values(Arc::clone(&pushed), &ctx).await,
Vec::<i32>::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 =
Expand Down
9 changes: 8 additions & 1 deletion datafusion/datasource/src/file_scan_config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -1028,6 +1028,13 @@ impl DataSource for FileScanConfig {
filters: Vec<Arc<dyn PhysicalExpr>>,
config: &ConfigOptions,
) -> Result<FilterPushdownPropagation<Arc<dyn DataSource>>> {
// 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
Expand Down