From 9e36f80fc7b21fbab6631bb74ffdd8c63d6c9b07 Mon Sep 17 00:00:00 2001 From: Denys Tsomenko Date: Sat, 12 Sep 2026 03:36:39 +0300 Subject: [PATCH 1/2] experiment(datafusion_iceberg): enable parquet row-filter pushdown only for wide scans (>= 8 projected columns) with a narrow predicate (<= 2 columns) --- datafusion_iceberg/src/table/mod.rs | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/datafusion_iceberg/src/table/mod.rs b/datafusion_iceberg/src/table/mod.rs index 591448d0..1d12b1d8 100644 --- a/datafusion_iceberg/src/table/mod.rs +++ b/datafusion_iceberg/src/table/mod.rs @@ -810,9 +810,22 @@ async fn table_scan( .iter() .map(|index| scan_schema.index_of(arrow_schema.field(*index).name())) .collect::, _>>()?; + // Row-filter pushdown (predicate evaluated inside the parquet decoder, late materialization + // of the remaining columns, TopK / join dynamic filters reaching the scan) pays off only when + // the scan is wide and the predicate narrow: `SELECT * ... WHERE url LIKE ... ORDER BY t LIMIT n` + // decodes 100+ columns for the few surviving rows. On narrow scans the same machinery costs + // more than the vectorized FilterExec it replaces, so it stays off there. + let filter_columns: std::collections::HashSet<_> = filters + .iter() + .flat_map(|f| f.column_refs().into_iter().cloned()) + .collect(); + let pushdown_filters = requested_projection.len() >= 8 + && !filter_columns.is_empty() + && filter_columns.len() <= 2; let file_source = Arc::new( ParquetSource::new(table_schema) - .with_parquet_file_reader_factory(parquet_reader_factory.clone()), + .with_parquet_file_reader_factory(parquet_reader_factory.clone()) + .with_pushdown_filters(pushdown_filters), ); // Create plan for every partition with delete files From f0d65fb11e0e536fd799eb774e94bbc45ae0535d Mon Sep 17 00:00:00 2001 From: Denys Tsomenko Date: Sat, 12 Sep 2026 14:18:26 +0300 Subject: [PATCH 2/2] style: rustfmt the pushdown-filter condition Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_012m5Yx6yEpZsacqAhZotgkT --- datafusion_iceberg/src/table/mod.rs | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/datafusion_iceberg/src/table/mod.rs b/datafusion_iceberg/src/table/mod.rs index 1d12b1d8..e45f057d 100644 --- a/datafusion_iceberg/src/table/mod.rs +++ b/datafusion_iceberg/src/table/mod.rs @@ -819,9 +819,8 @@ async fn table_scan( .iter() .flat_map(|f| f.column_refs().into_iter().cloned()) .collect(); - let pushdown_filters = requested_projection.len() >= 8 - && !filter_columns.is_empty() - && filter_columns.len() <= 2; + let pushdown_filters = + requested_projection.len() >= 8 && !filter_columns.is_empty() && filter_columns.len() <= 2; let file_source = Arc::new( ParquetSource::new(table_schema) .with_parquet_file_reader_factory(parquet_reader_factory.clone())