Repository navigation
Use exponential decay for multi-column join selectivity estimation #21583
Description
Activity
- changed the title
[-]feat: use exponential decay for multi-column join selectivity estimation[/-][+]Use exponential decay for multi-column join selectivity estimation[/+]on Apr 13, 2026 A reduced TPC-H SF1 reproducer for this issue, tracked in #25610.
Describe the bug
A partsupp–lineitem composite-key join overestimates output fourfold despite
exact input row counts. The estimator uses only the strongest equality key’s
selectivity.To Reproduce
From the repository root, with the CLI fix from
PR #25570 applied:cargo build --profile ci --locked -p datafusion-benchmarks --bin dfbench cargo install tpchgen-cli --version 1.1.1 --locked # if not already installed repro_dir=$(mktemp -d) tpchgen-cli --scale-factor 1 --format parquet \ --parquet-compression 'ZSTD(1)' --parts 1 --output-dir "$repro_dir/data" cat > "$repro_dir/repro.sql" <<'SQL' SET datafusion.execution.target_partitions = 1; SET datafusion.optimizer.enable_dynamic_filter_pushdown = false; SELECT l_orderkey FROM partsupp JOIN lineitem ON ps_partkey = l_partkey AND ps_suppkey = l_suppkey; SQL target/ci/dfbench statistics \ --path "$repro_dir/data" --query_path "$repro_dir/repro.sql"
Observed with
tpchgen-cli1.1.1 at
6c320561b5.
Inspect the SELECT reports; ignore the emptySETreports.Operator Node Estimated rows Actual rows HashJoinExec 024,004,860 6,001,215 Expected behavior
Account for additional keys using a correlation-aware rule, multi-column NDVs or
composite constraints. Blindly multiplying selectivities can underestimate
correlated keys.Additional context
Input row counts are 800,000 partsupp rows and 6,001,215 lineitems.
Reacted by Kumar Ujjawal- added a parent issue
on Sep 22, 2026 Apache Hive has a configuration option driving the use of either formula: the default is the current implementation that picks the max NDV value (for Hive it's
maxNdvForCorrelatedColumns), which works well for correlated columns.When columns are independent, Hive uses the same exponential backoff formula proposed here: HiveRelMdSelectivity.java#L217-L235.
So I am unsure we can just replace one formula with another. Until we have information about column independence, I think we will need a similar configuration option, as Apache Hive does. The
StatisticsRegistryframework can already carry such information through statistics extensions (ExtendedStatistics) from a custom provider, without changes to DataFusion (see this test), but the built-in join estimate does not read it yet.If in the future we standardize how to pass and consume optional information about correlated columns in DataFusion (for example a distinct count over the join key columns together, the kind of statistic PostgreSQL collects with
CREATE STATISTICS ... (ndistinct)), we can propagate this extended statistics by default, and the join cardinality estimate can use it when it is available, and fall back to the user-provided configuration otherwise.- added a commit that references this issue
on Oct 6, 2026
Is your feature request related to a problem or challenge?
Current join cardinality estimation for multi-column equi-joins is too conservative.
When a join has several equality keys, the estimate mostly follows only the single strongest key. That means extra join predicates do not reduce the estimated output rows enough. This can overestimate join size and lead to weaker join planning decisions.
Describe the solution you'd like
Use NDV-based exponential decay for multi-column join selectivity.
For an equi-join, compute one NDV factor per join key pair as:
max(NDV(left_key), NDV(right_key))Then sort these factors from largest to smallest and combine them with decay:
ndv0 * ndv1^(1/2) * ndv2^(1/4) * ndv3^(1/8) * ...Use that decayed value as the denominator for inner join cardinality estimation instead of using only the single largest NDV.
In practice, this would mean:
1,1/2,1/4, ...left_num_rows * right_num_rows / decayed_ndvThis should make multi-key join estimates tighter without being as aggressive as multiplying all key NDVs directly.
It would also be good to add tests for:
Describe alternatives you've considered
Keep the current largest-NDV-only rule.
This is simple, but it ignores useful information from the other join keys.
Multiply all join-key NDVs directly.
This is likely too aggressive and can under-estimate badly when keys are correlated.
Use histograms or stronger correlation models.
This could be more accurate, but it is a larger change and not needed for a first improvement.
Additional context
Part of #20766.
Generated with Codex