Skip to content

Commit bb24345

Browse files
refactor: deprecate the bundled statistics providers (#25969)
## Which issue does this PR close? - Closes #25571. ## Rationale for this change `StatisticsRegistry::default_with_builtin_providers()` registers providers that are not the default estimation and that duplicate or replace what the operators already do in `statistics_from_inputs`. They can make estimates worse (TPC-H Q14: 14.73 billion rows instead of 73,650, see #25570), and fixes in the operators do not reach their users. Estimation improvements belong in the operators; `StatisticsRegistry` stays as the extension point for user-defined providers. ## What changes are included in this PR? - `default_with_builtin_providers()` and all seven bundled providers are deprecated. Their code is unchanged. - `statistics_registry.slt` and the `join_reorder` example use user-defined providers instead of bundled ones. - `dfbench statistics` reports the operators' own estimates instead of the bundled providers' estimates. Follow-ups: moving the Filter distinct count survival model into `FilterExec` (#26052); multi-key join estimation is tracked in #21583. ## What is the testing strategy for this PR? `statistics_registry.slt` passes unchanged; the `join_reorder` example still flips the build side. ## Are there any user-facing changes? Deprecations only, see the 56.0.0 upgrade guide. ---- Disclaimer: I used AI to assist in the code generation, I have manually reviewed the output and it matches my intention and understanding. --------- Co-authored-by: Gabriel <45515538+gabotechs@users.noreply.github.com>
1 parent 8f4c51f commit bb24345

9 files changed

Lines changed: 125 additions & 49 deletions

File tree

‎benchmarks/src/statistics.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -249,7 +249,7 @@ enum QError {
249249
}
250250

251251
fn capture_statistics(plan: &dyn ExecutionPlan) -> Result<Vec<CapturedStatistics>> {
252-
let statistics_context = StatisticsRegistry::default_with_builtin_providers();
252+
let statistics_context = StatisticsRegistry::new();
253253
let mut result = vec![];
254254
capture_statistics_inner(plan, &statistics_context, "0", &mut result)?;
255255
Ok(result)

‎datafusion-examples/README.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -220,9 +220,9 @@ cargo run --example dataframe -- dataframe
220220

221221
#### Category: Single Process
222222

223-
| Subcommand | File Path | Description |
224-
| ------------ | ------------------------------------------------------------------- | ----------------------------------------------------------------------- |
225-
| join_reorder | [`statistics/join_reorder.rs`](examples/statistics/join_reorder.rs) | Supply and refine column statistics via a provider to flip a join order |
223+
| Subcommand | File Path | Description |
224+
| ------------ | ------------------------------------------------------------------- | -------------------------------------------------------------------- |
225+
| join_reorder | [`statistics/join_reorder.rs`](examples/statistics/join_reorder.rs) | Supply catalog column statistics via a provider to flip a join order |
226226

227227
## UDF Examples
228228

‎datafusion-examples/examples/statistics/join_reorder.rs‎

Lines changed: 18 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -15,21 +15,18 @@
1515
// specific language governing permissions and limitations
1616
// under the License.
1717

18-
//! Plug refined cardinality estimation into the optimizer via the `StatisticsRegistry`.
18+
//! Supply catalog statistics to the optimizer via the `StatisticsRegistry`.
1919
//!
20-
//! DataFusion's built-in statistics are intentionally simple defaults; the registry
21-
//! is the seam for plugging in more refined estimation (formulas or heuristics
22-
//! suited to your system or data) without changing the core.
20+
//! The registry is the extension point for statistics the core cannot derive,
21+
//! such as column statistics a catalog knows.
2322
//!
24-
//! `(SELECT user_id FROM events WHERE amount < 50 GROUP BY user_id) JOIN dims`: a
25-
//! provider supplies the column stats a catalog knows (`amount` range, `user_id`
26-
//! distinct count); the built-in `FilterStatisticsProvider` then refines the
27-
//! post-filter distinct count with the survival formula
28-
//! `NDV * (1 - (1 - selectivity)^(rows / NDV))` (Yao/Cardenas) to ~32, below `dims`
29-
//! (48), flipping the join build side. Core's simpler `min(NDV, rows)` cap would
30-
//! give 50 (> 48) and keep the other order; the refinement is the point. The
31-
//! ground-truth query prints the true surviving distinct count (below 48),
32-
//! confirming the flip.
23+
//! `(SELECT user_id FROM events WHERE amount < 50 GROUP BY user_id) JOIN dims`:
24+
//! the in-memory `events` table has no column statistics, so the grouped side is
25+
//! estimated at 1000 rows and `dims` (60 rows) becomes the build side. A provider
26+
//! supplies the `amount` range and the `user_id` distinct count a catalog knows;
27+
//! with them the operators estimate the grouped side at 50 rows, below `dims`,
28+
//! which flips the join build side. The ground-truth query prints the true
29+
//! distinct count (below 60), confirming the flip.
3330
3431
use std::sync::Arc;
3532

@@ -111,11 +108,9 @@ fn build_ctx(with_registry: bool) -> Result<SessionContext> {
111108
.with_config(config)
112109
.with_default_features();
113110
if with_registry {
114-
let mut registry = StatisticsRegistry::default_with_builtin_providers();
115-
registry.register(Arc::new(ClosureStatisticsProvider::with_matches(
116-
catalog_matches,
117-
catalog_stats,
118-
)));
111+
let registry = StatisticsRegistry::with_providers(vec![Arc::new(
112+
ClosureStatisticsProvider::with_matches(catalog_matches, catalog_stats),
113+
)]);
119114
builder = builder.with_statistics_registry(registry);
120115
}
121116
let ctx = SessionContext::new_with_state(builder.build());
@@ -135,8 +130,8 @@ fn build_ctx(with_registry: bool) -> Result<SessionContext> {
135130
ctx.register_table(
136131
"dims",
137132
mem_table(&[
138-
("user_id", int_col(&(0..48).collect::<Vec<_>>())),
139-
("label", int_col(&(0..48).collect::<Vec<_>>())),
133+
("user_id", int_col(&(0..60).collect::<Vec<_>>())),
134+
("label", int_col(&(0..60).collect::<Vec<_>>())),
140135
])?,
141136
)?;
142137
Ok(ctx)
@@ -167,13 +162,13 @@ pub async fn join_reorder() -> Result<()> {
167162
"A hash join builds its in-memory hash table from one input and probes with\n\
168163
the other, so the smaller input should be the build side. Default estimation\n\
169164
sizes the grouped `events` at 1000 rows and builds from `dims`; the\n\
170-
registry's refined ~32 estimate is below `dims` (48 rows), so it flips the\n\
171-
build side to `events`. The ground-truth count above (also below 48)\n\
165+
catalog statistics bring the estimate to 50, below `dims` (60 rows), so the\n\
166+
build side flips to `events`. The ground-truth count above (also below 60)\n\
172167
confirms `events` really is the smaller, cheaper side.\n"
173168
);
174169
println!("-- Without the registry (default estimation) --");
175170
println!("{}\n", explain(&build_ctx(false)?).await?);
176-
println!("-- With the registry (built-in refinement) --");
171+
println!("-- With the registry (catalog statistics) --");
177172
println!("{}", explain(&build_ctx(true)?).await?);
178173
Ok(())
179174
}

‎datafusion-examples/examples/statistics/main.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@
3030
//! - `all`: run all examples included in this module
3131
//!
3232
//! - `join_reorder`
33-
//! (file: join_reorder.rs, desc: Supply and refine column statistics via a provider to flip a join order)
33+
//! (file: join_reorder.rs, desc: Supply catalog column statistics via a provider to flip a join order)
3434
3535
mod join_reorder;
3636

‎datafusion/ffi/src/physical_optimizer.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -586,7 +586,7 @@ mod tests {
586586

587587
let ctx = ContextWithRegistry {
588588
config: ConfigOptions::new(),
589-
registry: StatisticsRegistry::default_with_builtin_providers(),
589+
registry: StatisticsRegistry::new(),
590590
};
591591

592592
let plan = create_test_plan();

‎datafusion/physical-plan/src/operator_statistics/mod.rs‎

Lines changed: 47 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -37,17 +37,12 @@
3737
//! - [`StatisticsRegistry`]: Chains providers, lives in SessionState
3838
//! - [`ExtendedStatistics`]: Statistics with type-safe custom extensions
3939
//!
40-
//! # Built-in Providers
40+
//! # Bundled Providers
4141
//!
42-
//! The following providers are included and can be registered in this order:
43-
//!
44-
//! 1. [`FilterStatisticsProvider`] - selectivity-based filter estimation
45-
//! 2. [`ProjectionStatisticsProvider`] - column mapping through projections
46-
//! 3. [`PassthroughStatisticsProvider`] - passthrough for cardinality-preserving operators
47-
//! 4. [`AggregateStatisticsProvider`] - NDV-based GROUP BY cardinality estimation
48-
//! 5. [`JoinStatisticsProvider`] - NDV-based join output estimation (hash, sort-merge, cross)
49-
//! 6. [`LimitStatisticsProvider`] - caps output at the fetch limit (local and global)
50-
//! 7. [`UnionStatisticsProvider`] - sums input row counts
42+
//! The providers bundled in this module are deprecated: the default statistics
43+
//! estimation is each operator's [`ExecutionPlan::statistics_from_inputs`], and
44+
//! estimation improvements belong there. Register user-defined providers to
45+
//! plug in other estimation.
5146
//!
5247
//! # Statistics walk
5348
//!
@@ -392,7 +387,7 @@ impl StatisticsRegistry {
392387
Self { providers }
393388
}
394389

395-
/// Create a registry pre-loaded with the standard built-in providers.
390+
/// Create a registry with the deprecated bundled providers.
396391
///
397392
/// Provider order (first match wins):
398393
/// 1. [`FilterStatisticsProvider`]
@@ -402,6 +397,11 @@ impl StatisticsRegistry {
402397
/// 5. [`JoinStatisticsProvider`]
403398
/// 6. [`LimitStatisticsProvider`]
404399
/// 7. [`UnionStatisticsProvider`]
400+
#[deprecated(
401+
since = "56.0.0",
402+
note = "the bundled providers are deprecated; the statistics walk uses each operator's `statistics_from_inputs` when no provider matches"
403+
)]
404+
#[expect(deprecated)]
405405
pub fn default_with_builtin_providers() -> Self {
406406
Self::with_providers(vec![
407407
Arc::new(FilterStatisticsProvider),
@@ -590,9 +590,14 @@ fn computed_with_row_count(
590590
/// estimation logic as `FilterExec::statistics_helper`, then additionally
591591
/// adjusts each column's `distinct_count` using [`ndv_after_selectivity`] based
592592
/// on the computed selectivity ratio.
593+
#[deprecated(
594+
since = "56.0.0",
595+
note = "duplicates `FilterExec::statistics_from_inputs` and adds a distinct count adjustment that belongs in the operator, see https://github.com/apache/datafusion/issues/26052; without a matching provider the statistics walk uses the operator's own estimate"
596+
)]
593597
#[derive(Debug, Default)]
594598
pub struct FilterStatisticsProvider;
595599

600+
#[expect(deprecated)]
596601
impl StatisticsProvider for FilterStatisticsProvider {
597602
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
598603
plan.downcast_ref::<FilterExec>().is_some()
@@ -649,9 +654,14 @@ impl StatisticsProvider for FilterStatisticsProvider {
649654
/// Maps enhanced child column statistics to output columns based on the
650655
/// projection expressions, preserving NDV and other statistics through
651656
/// column references.
657+
#[deprecated(
658+
since = "56.0.0",
659+
note = "duplicates the `statistics_from_inputs` of `ProjectionExec`; without a matching provider the statistics walk uses the operator's own estimate"
660+
)]
652661
#[derive(Debug, Default)]
653662
pub struct ProjectionStatisticsProvider;
654663

664+
#[expect(deprecated)]
655665
impl StatisticsProvider for ProjectionStatisticsProvider {
656666
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
657667
plan.downcast_ref::<ProjectionExec>().is_some()
@@ -691,9 +701,14 @@ impl StatisticsProvider for ProjectionStatisticsProvider {
691701
/// transform statistics, so we pass through the enhanced child stats directly.
692702
/// This avoids the fallback calling `statistics_from_inputs` (overall) which would
693703
/// trigger a redundant internal recursion with raw (non-enhanced) stats.
704+
#[deprecated(
705+
since = "56.0.0",
706+
note = "duplicates the `statistics_from_inputs` of the cardinality-preserving operators; without a matching provider the statistics walk uses the operator's own estimate"
707+
)]
694708
#[derive(Debug, Default)]
695709
pub struct PassthroughStatisticsProvider;
696710

711+
#[expect(deprecated)]
697712
impl StatisticsProvider for PassthroughStatisticsProvider {
698713
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
699714
plan.children().len() == 1
@@ -742,9 +757,14 @@ impl StatisticsProvider for PassthroughStatisticsProvider {
742757
/// - GROUP BY is empty (scalar aggregate)
743758
/// - Any GROUP BY expression is not a simple column reference
744759
/// - Any GROUP BY column lacks NDV information
760+
#[deprecated(
761+
since = "56.0.0",
762+
note = "duplicates the `statistics_from_inputs` of `AggregateExec`; without a matching provider the statistics walk uses the operator's own estimate"
763+
)]
745764
#[derive(Debug, Default)]
746765
pub struct AggregateStatisticsProvider;
747766

767+
#[expect(deprecated)]
748768
impl StatisticsProvider for AggregateStatisticsProvider {
749769
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
750770
plan.downcast_ref::<AggregateExec>().is_some()
@@ -835,9 +855,14 @@ impl StatisticsProvider for AggregateStatisticsProvider {
835855
/// Delegates when:
836856
/// - The plan is not a supported join type
837857
/// - Either input lacks row count information
858+
#[deprecated(
859+
since = "56.0.0",
860+
note = "replaces the estimates of `HashJoinExec`, `SortMergeJoinExec` and `CrossJoinExec` with separate estimation logic; without a matching provider the statistics walk uses the operator's own estimate"
861+
)]
838862
#[derive(Debug, Default)]
839863
pub struct JoinStatisticsProvider;
840864

865+
#[expect(deprecated)]
841866
impl StatisticsProvider for JoinStatisticsProvider {
842867
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
843868
plan.downcast_ref::<HashJoinExec>().is_some()
@@ -960,9 +985,14 @@ impl StatisticsProvider for JoinStatisticsProvider {
960985
///
961986
/// Caps output row count at the limit value, accounting for any leading skip offset
962987
/// in `GlobalLimitExec`.
988+
#[deprecated(
989+
since = "56.0.0",
990+
note = "duplicates the `statistics_from_inputs` of the limit operators; without a matching provider the statistics walk uses the operator's own estimate"
991+
)]
963992
#[derive(Debug, Default)]
964993
pub struct LimitStatisticsProvider;
965994

995+
#[expect(deprecated)]
966996
impl StatisticsProvider for LimitStatisticsProvider {
967997
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
968998
plan.downcast_ref::<LocalLimitExec>().is_some()
@@ -1011,9 +1041,14 @@ impl StatisticsProvider for LimitStatisticsProvider {
10111041
/// Statistics provider for [`UnionExec`].
10121042
///
10131043
/// Sums row counts across all inputs.
1044+
#[deprecated(
1045+
since = "56.0.0",
1046+
note = "duplicates the `statistics_from_inputs` of `UnionExec`; without a matching provider the statistics walk uses the operator's own estimate"
1047+
)]
10141048
#[derive(Debug, Default)]
10151049
pub struct UnionStatisticsProvider;
10161050

1051+
#[expect(deprecated)]
10171052
impl StatisticsProvider for UnionStatisticsProvider {
10181053
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
10191054
plan.downcast_ref::<UnionExec>().is_some()
@@ -1134,6 +1169,7 @@ impl StatisticsProvider for ClosureStatisticsProvider {
11341169
}
11351170

11361171
#[cfg(test)]
1172+
#[expect(deprecated)]
11371173
mod tests {
11381174
use super::*;
11391175
use crate::filter::FilterExec;

‎datafusion/sqllogictest/src/test_context.rs‎

Lines changed: 30 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ use datafusion::catalog::{
3737
CatalogProvider, MemoryCatalogProvider, MemorySchemaProvider, SchemaProvider, Session,
3838
};
3939
use datafusion::common::config::Dialect;
40+
use datafusion::common::stats::Precision;
4041
use datafusion::common::{DataFusionError, Result, not_impl_err};
4142
use datafusion::functions::math::abs;
4243
use datafusion::logical_expr::async_udf::{AsyncScalarUDF, AsyncScalarUDFImpl};
@@ -63,7 +64,11 @@ use datafusion::common::cast::as_float64_array;
6364
use datafusion::execution::SessionStateBuilder;
6465
use datafusion::execution::memory_pool::UnboundedMemoryPool;
6566
use datafusion::execution::runtime_env::RuntimeEnvBuilder;
66-
use datafusion::physical_plan::operator_statistics::StatisticsRegistry;
67+
use datafusion::physical_plan::joins::HashJoinExec;
68+
use datafusion::physical_plan::operator_statistics::{
69+
ClosureStatisticsProvider, StatisticsRegistry, StatisticsResult,
70+
};
71+
use datafusion::physical_plan::statistics::StatisticsArgs;
6772
use log::info;
6873
use sqlparser::ast;
6974
use tempfile::TempDir;
@@ -144,9 +149,31 @@ impl TestContext {
144149
relative_path.file_name().and_then(|name| name.to_str()),
145150
Some("statistics_registry.slt")
146151
) {
147-
state_builder = state_builder.with_statistics_registry(
148-
StatisticsRegistry::default_with_builtin_providers(),
152+
// Replaces the join estimate with the Cartesian product
153+
let join_provider = ClosureStatisticsProvider::with_matches(
154+
|plan| plan.downcast_ref::<HashJoinExec>().is_some(),
155+
|plan, child_stats| {
156+
let (Some(&left_rows), Some(&right_rows)) = (
157+
child_stats[0].base().num_rows.get_value(),
158+
child_stats[1].base().num_rows.get_value(),
159+
) else {
160+
return Ok(StatisticsResult::Delegate);
161+
};
162+
let child_base = child_stats
163+
.iter()
164+
.map(|c| Arc::clone(c.base_arc()))
165+
.collect::<Vec<_>>();
166+
let mut stats = Arc::unwrap_or_clone(
167+
plan.statistics_from_inputs(&child_base, &StatisticsArgs::new())?,
168+
);
169+
stats.num_rows =
170+
Precision::Inexact(left_rows.saturating_mul(right_rows));
171+
Ok(StatisticsResult::Computed(stats.into()))
172+
},
149173
);
174+
let registry =
175+
StatisticsRegistry::with_providers(vec![Arc::new(join_provider)]);
176+
state_builder = state_builder.with_statistics_registry(registry);
150177
}
151178

152179
let state = state_builder.build();

‎datafusion/sqllogictest/test_files/statistics_registry.slt‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,16 +18,17 @@
1818
# StatisticsRegistry: demonstrates improved join ordering via more conservative
1919
# cardinality estimates on a skewed dataset.
2020
#
21+
# The session registers a test-only provider in test_context.rs that replaces
22+
# the join estimate with the Cartesian product of the input row counts.
23+
#
2124
# customers (10 rows, 3 distinct customer_ids, skewed 8:1:1)
2225
# orders (10 rows, same distribution)
2326
# dim_small (50 rows)
2427
#
25-
# Built-in operator stats give 10*10 / NDV(3) = 33 < 50, which would keep the
28+
# The operators' own estimates give 10*10 / NDV(3) = 33 < 50, which would keep the
2629
# inner join on the build side (wrong; actual output is 66). The registry's
2730
# conservative estimate 10*10 = 100 > 50 swaps dim_small to the build side.
2831
#
29-
# Parquet files written by COPY TO carry min/max stats (NDV=3 via range) but no
30-
# distinct_count, so the registry falls back to the cartesian product upper bound.
3132
# Threshold settings force Partitioned mode so statistics alone drive the swap.
3233

3334
statement ok

‎docs/source/library-user-guide/upgrading/56.0.0.md‎

Lines changed: 20 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -448,9 +448,7 @@ previous default (`use_statistics_registry = false`).
448448

449449
- Drop any `use_statistics_registry` setting.
450450
- To use statistics providers, register them on the session with
451-
`SessionStateBuilder::with_statistics_registry(...)`. To keep the built-in
452-
NDV-aware providers the flag previously enabled, register
453-
`StatisticsRegistry::default_with_builtin_providers()`.
451+
`SessionStateBuilder::with_statistics_registry(...)`.
454452
- To opt out, register nothing (the default).
455453

456454
### Generated protobuf `JsonWriterOptions` gained a `compression_level` field
@@ -559,6 +557,25 @@ part of `StatisticsRegistry::default_with_builtin_providers()`.
559557
- Remove `DefaultStatisticsProvider` from any custom provider chain; register no
560558
terminal provider instead (the walk falls back on its own).
561559

560+
### The bundled statistics providers are deprecated
561+
562+
`StatisticsRegistry::default_with_builtin_providers()` and all the statistics
563+
providers bundled in `datafusion_physical_plan::operator_statistics`
564+
(`FilterStatisticsProvider`, `ProjectionStatisticsProvider`,
565+
`PassthroughStatisticsProvider`, `AggregateStatisticsProvider`,
566+
`JoinStatisticsProvider`, `LimitStatisticsProvider` and
567+
`UnionStatisticsProvider`) are deprecated. They duplicate or replace the
568+
estimation each operator already does in `statistics_from_inputs`, which is the
569+
default. `StatisticsRegistry` remains the extension point for user-defined
570+
providers.
571+
572+
**Migration guide:**
573+
574+
- Stop registering the bundled providers; with no matching provider, the
575+
statistics walk uses the operator's own estimate.
576+
- To keep a bundled estimation technique, implement it in a user-defined
577+
`StatisticsProvider`.
578+
562579
### `StatisticsRegistry::compute` and `compute_base` are deprecated
563580

564581
Use the walk instead:

0 commit comments

Comments
 (0)