From 65fdcca5646c30149fa4c95422df483baf3409b0 Mon Sep 17 00:00:00 2001 From: discord9 <55937128+discord9@users.noreply.github.com> Date: Tue, 8 Sep 2026 19:25:23 +0800 Subject: [PATCH 1/2] fix: preserve file scan fetch during filter pushdown Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --- .../core/tests/parquet/filter_pushdown.rs | 234 +++++++++++++++++- .../datasource/src/file_scan_config/mod.rs | 9 +- 2 files changed, 240 insertions(+), 3 deletions(-) diff --git a/datafusion/core/tests/parquet/filter_pushdown.rs b/datafusion/core/tests/parquet/filter_pushdown.rs index d337979e5fd00..4e69e37476342 100644 --- a/datafusion/core/tests/parquet/filter_pushdown.rs +++ b/datafusion/core/tests/parquet/filter_pushdown.rs @@ -28,10 +28,17 @@ 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_plan::filter::{FilterExec, FilterExecBuilder}; 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 +46,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 +72,226 @@ async fn read_parquet_test_data>(path: T) -> Vec { .unwrap() } +fn scan_with_fetch( + file: &TestParquetFile, + source: ParquetSource, + fetch: Option, +) -> Arc { + let path = file.path().canonicalize().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) + .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_scan_fetch_semantics() { + 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(); + let file = TestParquetFile::try_new( + tempdir.path().join("ordered.parquet"), + WriterProperties::builder() + .set_max_row_group_row_count(Some(1)) + .build(), + [RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from(vec![0, 1])) as ArrayRef], + ) + .unwrap()], + ) + .unwrap(); + + let mut pushdown_config = SessionConfig::new().with_target_partitions(1); + pushdown_config + .options_mut() + .execution + .parquet + .pushdown_filters = true; + let pushdown_ctx = SessionContext::new_with_config(pushdown_config); + + // A scan-owned fetch is an admission barrier: a parent filter must remain above it. + for (name, fetch, expected) in [ + ("fetch_one", Some(1), Vec::::new()), + ("fetch_zero", Some(0), Vec::::new()), + ] { + let original: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&pushdown_ctx, &schema, 1), + scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), fetch), + ) + .unwrap(), + ); + let pushed = FilterPushdown::new() + .optimize(Arc::clone(&original), pushdown_ctx.state().config_options()) + .unwrap(); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!(plan.contains("FilterExec"), "{name}: {plan}"); + assert!(!plan.contains("predicate="), "{name}: {plan}"); + assert_eq!(int_values(pushed, &pushdown_ctx).await, expected, "{name}"); + } + + // The full optimizer must preserve the same result when run once or twice. + let original: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&pushdown_ctx, &schema, 1), + scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), Some(1)), + ) + .unwrap(), + ); + let planner = DefaultPhysicalPlanner::default(); + let once = planner + .optimize_physical_plan(Arc::clone(&original), &pushdown_ctx.state(), |_, _| {}) + .unwrap(); + let twice = planner + .optimize_physical_plan(Arc::clone(&once), &pushdown_ctx.state(), |_, _| {}) + .unwrap(); + assert_eq!(int_values(original, &pushdown_ctx).await, Vec::::new()); + assert_eq!(int_values(once, &pushdown_ctx).await, Vec::::new()); + assert_eq!(int_values(twice, &pushdown_ctx).await, Vec::::new()); + + // An uncapped scan still accepts an exact pushed-down filter. + let uncapped: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&pushdown_ctx, &schema, 1), + scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), None), + ) + .unwrap(), + ); + let pushed = FilterPushdown::new() + .optimize(uncapped, pushdown_ctx.state().config_options()) + .unwrap(); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!(!plan.contains("FilterExec"), "{plan}"); + assert!(plan.contains("predicate=a@0 = 1"), "{plan}"); + assert_eq!(int_values(pushed, &pushdown_ctx).await, vec![1]); + + // With row filtering disabled, the filter is only eligible for pruning; it still + // cannot cross a scan fetch, even with statistics pruning enabled. + let pruning_ctx = SessionContext::new_with_config( + SessionConfig::new() + .with_target_partitions(1) + .with_parquet_pruning(true), + ); + let pruning_only: Arc = Arc::new( + FilterExec::try_new( + physical_predicate(&pruning_ctx, &schema, 1), + scan_with_fetch( + &file, + ParquetSource::new(Arc::clone(&schema)).with_pushdown_filters(false), + Some(1), + ), + ) + .unwrap(), + ); + let pushed = FilterPushdown::new() + .optimize(pruning_only, pruning_ctx.state().config_options()) + .unwrap(); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!( + plan.contains("FilterExec") && !plan.contains("predicate="), + "{plan}" + ); + assert_eq!(int_values(pushed, &pruning_ctx).await, Vec::::new()); + + // A predicate already owned by the scan still controls its fetch, while a + // new parent predicate remains above that capped scan. + let internal_file = TestParquetFile::try_new( + tempdir.path().join("internal.parquet"), + WriterProperties::builder() + .set_max_row_group_row_count(Some(1)) + .build(), + [RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from(vec![0, 1, 2])) as ArrayRef], + ) + .unwrap()], + ) + .unwrap(); + let internal_predicate = pruning_ctx + .create_physical_expr( + col("a").gt_eq(lit(1_i32)), + &Arc::clone(&schema).to_dfschema().unwrap(), + ) + .unwrap(); + let internal_scan = scan_with_fetch( + &internal_file, + ParquetSource::new(Arc::clone(&schema)) + .with_pushdown_filters(false) + .with_predicate(internal_predicate), + Some(1), + ); + assert_eq!( + int_values(Arc::clone(&internal_scan), &pruning_ctx).await, + vec![1] + ); + let parent: Arc = Arc::new( + FilterExec::try_new(physical_predicate(&pruning_ctx, &schema, 2), internal_scan) + .unwrap(), + ); + let pushed = FilterPushdown::new() + .optimize(parent, pruning_ctx.state().config_options()) + .unwrap(); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!(plan.contains("predicate=a@0 >= 1"), "{plan}"); + assert!(!plan.contains("predicate=a@0 = 2"), "{plan}"); + assert_eq!(int_values(pushed, &pruning_ctx).await, Vec::::new()); + + // A fetch owned by FilterExec, unlike one owned by the scan, may still push. + let filter_fetch: Arc = Arc::new( + FilterExecBuilder::new( + physical_predicate(&pushdown_ctx, &schema, 1), + scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), None), + ) + .with_fetch(Some(1)) + .build() + .unwrap(), + ); + let pushed = FilterPushdown::new() + .optimize(filter_fetch, pushdown_ctx.state().config_options()) + .unwrap(); + let plan = displayable(pushed.as_ref()).indent(false).to_string(); + assert!(plan.contains("predicate=a@0 = 1"), "{plan}"); + assert_eq!(int_values(pushed, &pushdown_ctx).await, vec![1]); +} + #[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 From 67f3e6543e7e6f5efe33346e1b9b2f13b01bd9dd Mon Sep 17 00:00:00 2001 From: discord9 <55937128+discord9@users.noreply.github.com> Date: Tue, 8 Sep 2026 19:44:58 +0800 Subject: [PATCH 2/2] test: separate capped scan filter admission coverage Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --- .../core/tests/parquet/filter_pushdown.rs | 279 +++++++++++------- 1 file changed, 169 insertions(+), 110 deletions(-) diff --git a/datafusion/core/tests/parquet/filter_pushdown.rs b/datafusion/core/tests/parquet/filter_pushdown.rs index 4e69e37476342..2039e2a72b259 100644 --- a/datafusion/core/tests/parquet/filter_pushdown.rs +++ b/datafusion/core/tests/parquet/filter_pushdown.rs @@ -35,7 +35,9 @@ 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_plan::filter::{FilterExec, FilterExecBuilder}; +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::{ExecutionPlan, collect, displayable}; use datafusion::physical_planner::DefaultPhysicalPlanner; @@ -72,12 +74,42 @@ async fn read_parquet_test_data>(path: T) -> Vec { .unwrap() } -fn scan_with_fetch( +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( @@ -85,6 +117,7 @@ fn scan_with_fetch( path.metadata().unwrap().len(), )) .with_limit(fetch) + .with_output_ordering(vec![ordering]) .build(); DataSourceExec::from_data_source(config) } @@ -121,175 +154,201 @@ async fn int_values(plan: Arc, ctx: &SessionContext) -> Vec::new()), - ("fetch_zero", Some(0), Vec::::new()), - ] { - let original: Arc = Arc::new( - FilterExec::try_new( - physical_predicate(&pushdown_ctx, &schema, 1), - scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), fetch), - ) - .unwrap(), - ); - let pushed = FilterPushdown::new() - .optimize(Arc::clone(&original), pushdown_ctx.state().config_options()) - .unwrap(); - let plan = displayable(pushed.as_ref()).indent(false).to_string(); - assert!(plan.contains("FilterExec"), "{name}: {plan}"); - assert!(!plan.contains("predicate="), "{name}: {plan}"); - assert_eq!(int_values(pushed, &pushdown_ctx).await, expected, "{name}"); - } + // 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); - // The full optimizer must preserve the same result when run once or twice. let original: Arc = Arc::new( FilterExec::try_new( - physical_predicate(&pushdown_ctx, &schema, 1), - scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), Some(1)), + 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), &pushdown_ctx.state(), |_, _| {}) + .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(Arc::clone(&once), &pushdown_ctx.state(), |_, _| {}) + .optimize_physical_plan(once, &ctx.state(), |_, _| {}) .unwrap(); - assert_eq!(int_values(original, &pushdown_ctx).await, Vec::::new()); - assert_eq!(int_values(once, &pushdown_ctx).await, Vec::::new()); - assert_eq!(int_values(twice, &pushdown_ctx).await, Vec::::new()); + assert_eq!(int_values(twice, &ctx).await, Vec::::new()); - // An uncapped scan still accepts an exact pushed-down filter. let uncapped: Arc = Arc::new( FilterExec::try_new( - physical_predicate(&pushdown_ctx, &schema, 1), - scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), None), + physical_predicate(&ctx, &schema, 1), + ordered_scan_with_fetch( + &file, + &schema, + ParquetSource::new(Arc::clone(&schema)), + None, + ), ) .unwrap(), ); let pushed = FilterPushdown::new() - .optimize(uncapped, pushdown_ctx.state().config_options()) + .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}"); - assert!(plan.contains("predicate=a@0 = 1"), "{plan}"); - assert_eq!(int_values(pushed, &pushdown_ctx).await, vec![1]); + assert!( + !plan.contains("FilterExec") && plan.contains("predicate=a@0 = 1"), + "{plan}" + ); +} - // With row filtering disabled, the filter is only eligible for pruning; it still - // cannot cross a scan fetch, even with statistics pruning enabled. - let pruning_ctx = SessionContext::new_with_config( +#[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 pruning_only: Arc = Arc::new( + + let original: Arc = Arc::new( FilterExec::try_new( - physical_predicate(&pruning_ctx, &schema, 1), - scan_with_fetch( - &file, + 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(pruning_only, pruning_ctx.state().config_options()) + .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}" ); - assert_eq!(int_values(pushed, &pruning_ctx).await, Vec::::new()); - - // A predicate already owned by the scan still controls its fetch, while a - // new parent predicate remains above that capped scan. - let internal_file = TestParquetFile::try_new( - tempdir.path().join("internal.parquet"), - WriterProperties::builder() - .set_max_row_group_row_count(Some(1)) - .build(), - [RecordBatch::try_new( - Arc::clone(&schema), - vec![Arc::new(Int32Array::from(vec![0, 1, 2])) as ArrayRef], - ) - .unwrap()], - ) - .unwrap(); - let internal_predicate = pruning_ctx + + 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 = scan_with_fetch( + 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), &pruning_ctx).await, - vec![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(), ); - let parent: Arc = Arc::new( - FilterExec::try_new(physical_predicate(&pruning_ctx, &schema, 2), internal_scan) - .unwrap(), + assert_eq!( + int_values(Arc::clone(&original), &ctx).await, + Vec::::new() ); let pushed = FilterPushdown::new() - .optimize(parent, pruning_ctx.state().config_options()) + .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}"); - assert_eq!(int_values(pushed, &pruning_ctx).await, Vec::::new()); - - // A fetch owned by FilterExec, unlike one owned by the scan, may still push. - let filter_fetch: Arc = Arc::new( - FilterExecBuilder::new( - physical_predicate(&pushdown_ctx, &schema, 1), - scan_with_fetch(&file, ParquetSource::new(Arc::clone(&schema)), None), - ) - .with_fetch(Some(1)) - .build() - .unwrap(), - ); - let pushed = FilterPushdown::new() - .optimize(filter_fetch, pushdown_ctx.state().config_options()) - .unwrap(); - let plan = displayable(pushed.as_ref()).indent(false).to_string(); - assert!(plan.contains("predicate=a@0 = 1"), "{plan}"); - assert_eq!(int_values(pushed, &pushdown_ctx).await, vec![1]); } #[tokio::test]