Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
2f007cd
added ability to eagerly prune parquet files for better statistics
athlcode Sep 15, 2026
dc44eb9
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 15, 2026
ead1d34
fixed license header
athlcode Sep 15, 2026
8619b96
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 16, 2026
01fa26a
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 17, 2026
e11e013
build fix and tests
athlcode Sep 17, 2026
6865b01
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 19, 2026
fad5f9c
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 22, 2026
6c6951c
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 23, 2026
45dc4ed
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 24, 2026
4cd61f0
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 24, 2026
1628174
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 26, 2026
2f04e6a
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 27, 2026
84dc1a5
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 27, 2026
5a2d7d4
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 28, 2026
40aaaab
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 28, 2026
fa3eef2
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 30, 2026
99acad3
Skip bloom filters of fully matched row groups in eager pruning; upda…
athlcode Sep 30, 2026
e45177d
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Sep 30, 2026
d928123
Merge branch 'main' of https://github.com/apache/datafusion into prun…
athlcode Oct 1, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ use datafusion::{
error::Result,
execution::session_state::SessionStateBuilder,
physical_expr_common::sort_expr::LexRequirement,
physical_plan::ExecutionPlan,
physical_plan::{ExecutionPlan, PhysicalExpr},
prelude::SessionContext,
};

Expand Down Expand Up @@ -147,8 +147,11 @@ impl FileFormat for TSVFileFormat {
&self,
state: &dyn Session,
conf: FileScanConfig,
filters: &[Arc<dyn PhysicalExpr>],
) -> Result<Arc<dyn ExecutionPlan>> {
self.csv_file_format.create_physical_plan(state, conf).await
self.csv_file_format
.create_physical_plan(state, conf, filters)
.await
}

async fn create_writer_physical_plan(
Expand Down
30 changes: 29 additions & 1 deletion datafusion/catalog-listing/src/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -769,11 +769,39 @@ impl ListingTable {
.with_expr_adapter(self.expr_adapter_factory.clone())
.build();

// Pass the (non partition) filters to the file format, which may use
// them e.g. to prune files while planning. Filters are still evaluated
// above the scan (see `supports_filters_pushdown`), so filters that
// cannot be converted are simply not passed and never fail the scan.
let physical_filters = if filters.is_empty() {
vec![]
} else {
match DFSchema::try_from(Arc::clone(&self.table_schema)) {
Ok(df_schema) => filters
.iter()
.filter_map(|filter| {
state
.create_physical_expr(filter.clone(), &df_schema)
.inspect_err(|e| {
log::debug!(
"Not passing filter {filter} to the file format: {e}"
)
})
.ok()
})
.collect::<Vec<_>>(),
Err(e) => {
log::debug!("Not passing filters to the file format: {e}");
vec![]
}
}
};

// create the execution plan
let plan = self
.options
.format
.create_physical_plan(state, scan_config)
.create_physical_plan(state, scan_config, &physical_filters)
.await?;

Ok(ScanResult::new(plan))
Expand Down
86 changes: 86 additions & 0 deletions datafusion/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -585,6 +585,75 @@ impl Display for SpillCompression {
}
}

/// How much pruning work the Parquet scan performs during planning to refine
/// the scan's statistics.
///
/// Pruning normally happens only when a `DataSourceExec` is executed, so the
/// optimizer sees row counts for whole files even when a filter can skip most
/// of them. When eager pruning is enabled, pruning is performed while the scan
/// is planned and its results are used to tighten the (inexact) row count and
/// byte size estimates of the scan, e.g. for join ordering.
///
/// The levels form a ladder: each level includes the work of all lower levels.
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum EagerParquetPruning {
/// No pruning during planning
#[default]
Disabled,
/// Prune row groups using row group statistics (min/max)
RowGroups,
/// Additionally prune data pages using the page index
PageIndex,
/// Additionally prune row groups using bloom filters
BloomFilters,
}

impl EagerParquetPruning {
/// Returns true if this level performs the work of `level`
pub fn includes(self, level: EagerParquetPruning) -> bool {
self >= level
}
}

impl FromStr for EagerParquetPruning {
type Err = DataFusionError;

fn from_str(s: &str) -> Result<Self, Self::Err> {
match s.to_ascii_lowercase().as_str() {
"disabled" | "" => Ok(Self::Disabled),
"row_groups" => Ok(Self::RowGroups),
"page_index" => Ok(Self::PageIndex),
"bloom_filters" => Ok(Self::BloomFilters),
other => Err(DataFusionError::Configuration(format!(
"Invalid eager parquet pruning level: {other}. Expected one of: disabled, row_groups, page_index, bloom_filters"
))),
}
}
}

impl ConfigField for EagerParquetPruning {
fn visit<V: Visit>(&self, v: &mut V, key: &str, description: &'static str) {
v.some(key, self, description)
}

fn set(&mut self, _: &str, value: &str) -> Result<()> {
*self = EagerParquetPruning::from_str(value)?;
Ok(())
}
}

impl Display for EagerParquetPruning {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let str = match self {
Self::Disabled => "disabled",
Self::RowGroups => "row_groups",
Self::PageIndex => "page_index",
Self::BloomFilters => "bloom_filters",
};
write!(f, "{str}")
}
}

/// A `usize` configuration value that rejects zero when set from strings.
///
/// Use this for options where zero is never a meaningful runtime value.
Expand Down Expand Up @@ -1423,6 +1492,23 @@ config_namespace! {
/// Defaults to 20.
pub max_in_list_size: usize, default = 20

/// (reading) How much pruning the parquet scan performs while the query
/// is planned, to tighten the scan's row count and byte size estimates.
/// When the scan executes, it starts from the pruned files and row
/// groups, but evaluates its own predicate again.
///
/// Valid values are `disabled`, `row_groups` (row group statistics),
/// `page_index` (additionally the page index) and `bloom_filters`
/// (additionally bloom filters). Each level includes the work of the
/// levels before it and may add I/O and latency to planning. Statistics
/// derived from eager pruning are always inexact.
pub eager_pruning: EagerParquetPruning, default = EagerParquetPruning::Disabled

/// (reading) Eager pruning is skipped for scans with more files than
/// this limit, bounding the planning time spent by
/// `eager_pruning`. A limit of 0 skips eager pruning for every scan.
pub eager_pruning_file_limit: usize, default = 256

// The following options affect writing to parquet files
// and map to parquet::file::properties::WriterProperties

Expand Down
7 changes: 7 additions & 0 deletions datafusion/common/src/file_options/parquet_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,8 @@ impl ParquetOptions {
skip_arrow_metadata: _,
max_predicate_cache_size: _,
max_in_list_size: _,
eager_pruning: _, // reads not used for writer props
eager_pruning_file_limit: _, // reads not used for writer props
} = self;

let mut builder = WriterProperties::builder()
Expand Down Expand Up @@ -409,6 +411,8 @@ mod tests {
enable_page_index: defaults.enable_page_index,
pruning: defaults.pruning,
max_in_list_size: defaults.max_in_list_size,
eager_pruning: defaults.eager_pruning,
eager_pruning_file_limit: defaults.eager_pruning_file_limit,
skip_metadata: defaults.skip_metadata,
metadata_size_hint: defaults.metadata_size_hint,
pushdown_filters: defaults.pushdown_filters,
Expand Down Expand Up @@ -534,6 +538,9 @@ mod tests {
enable_page_index: global_options_defaults.enable_page_index,
pruning: global_options_defaults.pruning,
max_in_list_size: global_options_defaults.max_in_list_size,
eager_pruning: global_options_defaults.eager_pruning,
eager_pruning_file_limit: global_options_defaults
.eager_pruning_file_limit,
skip_metadata: global_options_defaults.skip_metadata,
metadata_size_hint: global_options_defaults.metadata_size_hint,
pushdown_filters: global_options_defaults.pushdown_filters,
Expand Down
1 change: 1 addition & 0 deletions datafusion/core/src/datasource/file_format/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ pub(crate) mod test_util {
.with_projection_indices(projection)?
.with_limit(limit)
.build(),
&[],
)
.await?;
Ok(exec)
Expand Down
Loading
Loading