From 6c320561b5b1aef7a235a12435c3c96b62956c67 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Mon, 21 Sep 2026 17:26:56 +0200 Subject: [PATCH 1/2] fix: use session statistics registry in dfbench reports --- benchmarks/src/statistics.rs | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/benchmarks/src/statistics.rs b/benchmarks/src/statistics.rs index 5c0be9984eceb..2787639f1e508 100644 --- a/benchmarks/src/statistics.rs +++ b/benchmarks/src/statistics.rs @@ -170,7 +170,10 @@ impl RunOpt { let logical_plan = state.optimize(&logical_plan)?; let physical_plan = state.create_physical_plan(&logical_plan).await?; - let statistics = capture_statistics(physical_plan.as_ref())?; + let statistics = capture_statistics( + physical_plan.as_ref(), + &state.statistics_registry().cloned().unwrap_or_default(), + )?; collect(Arc::clone(&physical_plan), state.task_ctx()).await?; let mut report = Vec::with_capacity(statistics.len()); @@ -234,10 +237,12 @@ enum QError { ExactZero, } -fn capture_statistics(plan: &dyn ExecutionPlan) -> Result> { - let statistics_context = StatisticsRegistry::default_with_builtin_providers(); +fn capture_statistics( + plan: &dyn ExecutionPlan, + statistics_context: &StatisticsRegistry, +) -> Result> { let mut result = vec![]; - capture_statistics_inner(plan, &statistics_context, "0", &mut result)?; + capture_statistics_inner(plan, statistics_context, "0", &mut result)?; Ok(result) } From d988cdf870280c3ac98daadd9d722ff3e52db373 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Mon, 5 Oct 2026 13:10:44 +0200 Subject: [PATCH 2/2] test: verify dfbench uses session statistics registry --- benchmarks/src/statistics.rs | 42 ++++++++++++++++++++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/benchmarks/src/statistics.rs b/benchmarks/src/statistics.rs index d94ddfc6dfbb1..fdc3c193ce02b 100644 --- a/benchmarks/src/statistics.rs +++ b/benchmarks/src/statistics.rs @@ -728,6 +728,11 @@ fn collect_parquet_files(path: &Path, files: &mut Vec) -> Result<()> { #[cfg(test)] mod tests { use super::*; + use datafusion::execution::session_state::SessionStateBuilder; + use datafusion::physical_plan::Statistics; + use datafusion::physical_plan::operator_statistics::{ + ClosureStatisticsProvider, StatisticsResult, + }; use tempfile::tempdir; #[test] @@ -791,6 +796,43 @@ mod tests { assert!(!reports.is_empty()); } + #[tokio::test] + async fn reports_estimates_from_session_statistics_registry() { + let directory = tempdir().unwrap(); + let options = RunOpt { + query: None, + compare: None, + path: directory.path().to_path_buf(), + query_path: directory.path().to_path_buf(), + }; + let mut registry = StatisticsRegistry::new(); + registry.register(Arc::new(ClosureStatisticsProvider::new( + |plan, _child_stats| { + Ok(StatisticsResult::Computed( + Statistics { + num_rows: Precision::Inexact(42), + ..Statistics::new_unknown(plan.schema().as_ref()) + } + .into(), + )) + }, + ))); + let ctx = SessionContext::from( + SessionStateBuilder::new() + .with_default_features() + .with_statistics_registry(registry) + .build(), + ); + let statement = + sql_statements("SELECT 1", &ctx.state().config_options().sql_parser) + .unwrap() + .pop_front() + .unwrap(); + + let reports = options.report_statement(&ctx, statement).await.unwrap(); + assert_eq!(reports[0].estimated_rows, StatisticValue::Inexact(42)); + } + async fn report_query_files( options: &RunOpt, ctx: &SessionContext,