Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
2 changes: 1 addition & 1 deletion benchmarks/src/statistics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,7 +249,7 @@ enum QError {
}

fn capture_statistics(plan: &dyn ExecutionPlan) -> Result<Vec<CapturedStatistics>> {
let statistics_context = StatisticsRegistry::default_with_builtin_providers();
let statistics_context = StatisticsRegistry::new();
let mut result = vec![];
capture_statistics_inner(plan, &statistics_context, "0", &mut result)?;
Ok(result)
Expand Down
6 changes: 3 additions & 3 deletions datafusion-examples/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -220,9 +220,9 @@ cargo run --example dataframe -- dataframe

#### Category: Single Process

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

## UDF Examples

Expand Down
41 changes: 18 additions & 23 deletions datafusion-examples/examples/statistics/join_reorder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,21 +15,18 @@
// specific language governing permissions and limitations
// under the License.

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

use std::sync::Arc;

Expand Down Expand Up @@ -111,11 +108,9 @@ fn build_ctx(with_registry: bool) -> Result<SessionContext> {
.with_config(config)
.with_default_features();
if with_registry {
let mut registry = StatisticsRegistry::default_with_builtin_providers();
registry.register(Arc::new(ClosureStatisticsProvider::with_matches(
catalog_matches,
catalog_stats,
)));
let registry = StatisticsRegistry::with_providers(vec![Arc::new(
ClosureStatisticsProvider::with_matches(catalog_matches, catalog_stats),
)]);
builder = builder.with_statistics_registry(registry);
}
let ctx = SessionContext::new_with_state(builder.build());
Expand All @@ -135,8 +130,8 @@ fn build_ctx(with_registry: bool) -> Result<SessionContext> {
ctx.register_table(
"dims",
mem_table(&[
("user_id", int_col(&(0..48).collect::<Vec<_>>())),
("label", int_col(&(0..48).collect::<Vec<_>>())),
("user_id", int_col(&(0..60).collect::<Vec<_>>())),
("label", int_col(&(0..60).collect::<Vec<_>>())),
])?,
)?;
Ok(ctx)
Expand Down Expand Up @@ -167,13 +162,13 @@ pub async fn join_reorder() -> Result<()> {
"A hash join builds its in-memory hash table from one input and probes with\n\
the other, so the smaller input should be the build side. Default estimation\n\
sizes the grouped `events` at 1000 rows and builds from `dims`; the\n\
registry's refined ~32 estimate is below `dims` (48 rows), so it flips the\n\
build side to `events`. The ground-truth count above (also below 48)\n\
catalog statistics bring the estimate to 50, below `dims` (60 rows), so the\n\
build side flips to `events`. The ground-truth count above (also below 60)\n\
confirms `events` really is the smaller, cheaper side.\n"
);
println!("-- Without the registry (default estimation) --");
println!("{}\n", explain(&build_ctx(false)?).await?);
println!("-- With the registry (built-in refinement) --");
println!("-- With the registry (catalog statistics) --");
println!("{}", explain(&build_ctx(true)?).await?);
Ok(())
}
2 changes: 1 addition & 1 deletion datafusion-examples/examples/statistics/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@
//! - `all`: run all examples included in this module
//!
//! - `join_reorder`
//! (file: join_reorder.rs, desc: Supply and refine column statistics via a provider to flip a join order)
//! (file: join_reorder.rs, desc: Supply catalog column statistics via a provider to flip a join order)

mod join_reorder;

Expand Down
2 changes: 1 addition & 1 deletion datafusion/ffi/src/physical_optimizer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -586,7 +586,7 @@ mod tests {

let ctx = ContextWithRegistry {
config: ConfigOptions::new(),
registry: StatisticsRegistry::default_with_builtin_providers(),
registry: StatisticsRegistry::new(),
};

let plan = create_test_plan();
Expand Down
58 changes: 47 additions & 11 deletions datafusion/physical-plan/src/operator_statistics/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,17 +37,12 @@
//! - [`StatisticsRegistry`]: Chains providers, lives in SessionState
//! - [`ExtendedStatistics`]: Statistics with type-safe custom extensions
//!
//! # Built-in Providers
//! # Bundled Providers
//!
//! The following providers are included and can be registered in this order:
//!
//! 1. [`FilterStatisticsProvider`] - selectivity-based filter estimation
//! 2. [`ProjectionStatisticsProvider`] - column mapping through projections
//! 3. [`PassthroughStatisticsProvider`] - passthrough for cardinality-preserving operators
//! 4. [`AggregateStatisticsProvider`] - NDV-based GROUP BY cardinality estimation
//! 5. [`JoinStatisticsProvider`] - NDV-based join output estimation (hash, sort-merge, cross)
//! 6. [`LimitStatisticsProvider`] - caps output at the fetch limit (local and global)
//! 7. [`UnionStatisticsProvider`] - sums input row counts
//! The providers bundled in this module are deprecated: the default statistics
//! estimation is each operator's [`ExecutionPlan::statistics_from_inputs`], and
//! estimation improvements belong there. Register user-defined providers to
//! plug in other estimation.
//!
//! # Statistics walk
//!
Expand Down Expand Up @@ -392,7 +387,7 @@ impl StatisticsRegistry {
Self { providers }
}

/// Create a registry pre-loaded with the standard built-in providers.
/// Create a registry with the deprecated bundled providers.
///
/// Provider order (first match wins):
/// 1. [`FilterStatisticsProvider`]
Expand All @@ -402,6 +397,11 @@ impl StatisticsRegistry {
/// 5. [`JoinStatisticsProvider`]
/// 6. [`LimitStatisticsProvider`]
/// 7. [`UnionStatisticsProvider`]
#[deprecated(
since = "56.0.0",
note = "the bundled providers are deprecated; the statistics walk uses each operator's `statistics_from_inputs` when no provider matches"
)]
#[expect(deprecated)]
Comment on lines +400 to +404

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

👍 looks good. Not sure which tickets are already open for this, but what comes to mind is:

  1. Fold the improvements brought by these providers into the built-in partition_statistics method.
  2. Remove the now redundant providers

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The two follow-up issues are listed here, which follow your proposal 1. to fold improvements in the built-in estimation.

For 2. I have filed #26055 following https://datafusion.apache.org/contributor-guide/api-health.html#deprecation-guidelines so we won't forget.

pub fn default_with_builtin_providers() -> Self {
Self::with_providers(vec![
Arc::new(FilterStatisticsProvider),
Expand Down Expand Up @@ -590,9 +590,14 @@ fn computed_with_row_count(
/// estimation logic as `FilterExec::statistics_helper`, then additionally
/// adjusts each column's `distinct_count` using [`ndv_after_selectivity`] based
/// on the computed selectivity ratio.
#[deprecated(
since = "56.0.0",
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"
)]
#[derive(Debug, Default)]
pub struct FilterStatisticsProvider;

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

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

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

#[expect(deprecated)]
impl StatisticsProvider for AggregateStatisticsProvider {
fn matches(&self, plan: &dyn ExecutionPlan) -> bool {
plan.downcast_ref::<AggregateExec>().is_some()
Expand Down Expand Up @@ -835,9 +855,14 @@ impl StatisticsProvider for AggregateStatisticsProvider {
/// Delegates when:
/// - The plan is not a supported join type
/// - Either input lacks row count information
#[deprecated(
since = "56.0.0",
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"
)]
#[derive(Debug, Default)]
pub struct JoinStatisticsProvider;

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

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

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

#[cfg(test)]
#[expect(deprecated)]
mod tests {
use super::*;
use crate::filter::FilterExec;
Expand Down
33 changes: 30 additions & 3 deletions datafusion/sqllogictest/src/test_context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ use datafusion::catalog::{
CatalogProvider, MemoryCatalogProvider, MemorySchemaProvider, SchemaProvider, Session,
};
use datafusion::common::config::Dialect;
use datafusion::common::stats::Precision;
use datafusion::common::{DataFusionError, Result, not_impl_err};
use datafusion::functions::math::abs;
use datafusion::logical_expr::async_udf::{AsyncScalarUDF, AsyncScalarUDFImpl};
Expand All @@ -63,7 +64,11 @@ use datafusion::common::cast::as_float64_array;
use datafusion::execution::SessionStateBuilder;
use datafusion::execution::memory_pool::UnboundedMemoryPool;
use datafusion::execution::runtime_env::RuntimeEnvBuilder;
use datafusion::physical_plan::operator_statistics::StatisticsRegistry;
use datafusion::physical_plan::joins::HashJoinExec;
use datafusion::physical_plan::operator_statistics::{
ClosureStatisticsProvider, StatisticsRegistry, StatisticsResult,
};
use datafusion::physical_plan::statistics::StatisticsArgs;
use log::info;
use sqlparser::ast;
use tempfile::TempDir;
Expand Down Expand Up @@ -144,9 +149,31 @@ impl TestContext {
relative_path.file_name().and_then(|name| name.to_str()),
Some("statistics_registry.slt")
) {
state_builder = state_builder.with_statistics_registry(
StatisticsRegistry::default_with_builtin_providers(),
// Replaces the join estimate with the Cartesian product
let join_provider = ClosureStatisticsProvider::with_matches(
|plan| plan.downcast_ref::<HashJoinExec>().is_some(),
Comment on lines +152 to +154

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I see this is still using a ClosureStatisticsProvider. I see this is mainly just for testing that the registry still works and actually listens to what the ClosureStatisticsProvider has to say.

If it's just that, all good, but let me know if missed something.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, I trimmed the example down to exactly that, just testing the registry works, you got it right

|plan, child_stats| {
let (Some(&left_rows), Some(&right_rows)) = (
child_stats[0].base().num_rows.get_value(),
child_stats[1].base().num_rows.get_value(),
) else {
return Ok(StatisticsResult::Delegate);
};
let child_base = child_stats
.iter()
.map(|c| Arc::clone(c.base_arc()))
.collect::<Vec<_>>();
let mut stats = Arc::unwrap_or_clone(
plan.statistics_from_inputs(&child_base, &StatisticsArgs::new())?,
);
stats.num_rows =
Precision::Inexact(left_rows.saturating_mul(right_rows));
Ok(StatisticsResult::Computed(stats.into()))
},
);
let registry =
StatisticsRegistry::with_providers(vec![Arc::new(join_provider)]);
state_builder = state_builder.with_statistics_registry(registry);
}

let state = state_builder.build();
Expand Down
7 changes: 4 additions & 3 deletions datafusion/sqllogictest/test_files/statistics_registry.slt
Original file line number Diff line number Diff line change
Expand Up @@ -18,16 +18,17 @@
# StatisticsRegistry: demonstrates improved join ordering via more conservative
# cardinality estimates on a skewed dataset.
#
# The session registers a test-only provider in test_context.rs that replaces
# the join estimate with the Cartesian product of the input row counts.
#
# customers (10 rows, 3 distinct customer_ids, skewed 8:1:1)
# orders (10 rows, same distribution)
# dim_small (50 rows)
#
# Built-in operator stats give 10*10 / NDV(3) = 33 < 50, which would keep the
# The operators' own estimates give 10*10 / NDV(3) = 33 < 50, which would keep the
# inner join on the build side (wrong; actual output is 66). The registry's
# conservative estimate 10*10 = 100 > 50 swaps dim_small to the build side.
#
# Parquet files written by COPY TO carry min/max stats (NDV=3 via range) but no
# distinct_count, so the registry falls back to the cartesian product upper bound.
# Threshold settings force Partitioned mode so statistics alone drive the swap.

statement ok
Expand Down
23 changes: 20 additions & 3 deletions docs/source/library-user-guide/upgrading/56.0.0.md
Original file line number Diff line number Diff line change
Expand Up @@ -448,9 +448,7 @@ previous default (`use_statistics_registry = false`).

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

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

### The bundled statistics providers are deprecated

`StatisticsRegistry::default_with_builtin_providers()` and all the statistics
providers bundled in `datafusion_physical_plan::operator_statistics`
(`FilterStatisticsProvider`, `ProjectionStatisticsProvider`,
`PassthroughStatisticsProvider`, `AggregateStatisticsProvider`,
`JoinStatisticsProvider`, `LimitStatisticsProvider` and
`UnionStatisticsProvider`) are deprecated. They duplicate or replace the
estimation each operator already does in `statistics_from_inputs`, which is the
default. `StatisticsRegistry` remains the extension point for user-defined
providers.

**Migration guide:**

- Stop registering the bundled providers; with no matching provider, the
statistics walk uses the operator's own estimate.
- To keep a bundled estimation technique, implement it in a user-defined
`StatisticsProvider`.

### `StatisticsRegistry::compute` and `compute_base` are deprecated

Use the walk instead:
Expand Down
Loading