apache / apache/datafusion

Improve NDV merge to be associative across multiple inputs

Open
#20,966 0 comments 1 reaction 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

The overlap-based NDV merge formula (`estimate_ndv_with_overlap`) introduced in #19957 (and shared with Union from #20846) is not associative: merging `(A+B)+C` produces a different result than `A+(B+C)` because the intermediate merge updates min/max/NDV, which "smears" the column statistics before the next pairwise merge.

Example:
```
A = [0,100], NDV=80
B = [50,150], NDV=60
C = [100,200], NDV=50

(A+B)+C = 135
A+(B+C) = 137
```

Current usage in `try_merge`, as addressed in #19957, folds left-to-right over row groups, so the result depends on row group ordering. Practically, row group order within a Parquet file is stable, so results are deterministic for a given file. The difference is typically small (bounded by the uniform-distribution assumption).

### Describe the solution you'd like

A multi-way merge that computes the overlap formula over all inputs at once rather than pairwise, using the original (unsmeared) min/max/NDV from each input.

### Describe alternatives you've considered

Store HLL sketches per row group/column in Parquet and merge them for exact NDV computation, bypassing the scalar heuristic entirely. Parquet does not natively support this, but [apache/arrow-rs#8608 (comment)](https://github.com/apache/arrow-rs/issues/8608#issuecomment-4041324952) proposes encoding distinct-count estimates in user-defined key=value metadata, which could serve as a vehicle for HLL sketches. Citing as an alternative for completeness, not currently viable as not implemented yet.

### Additional context

- Review comment: https://github.com/apache/datafusion/pull/19957#discussion_r2936704265 (@gene-bordegaray)
- Related PR: #20846 (Union NDV overlap formula)
- Related EPIC: https://github.com/apache/datafusion/issues/20766

Contributor guide

Open the contributing guide

Research direction

Start at the try_merge path and the estimate_ndv_with_overlap implementation, then trace how row-group statistics are folded. Design the multi-way merge around the original min/max/NDV inputs rather than intermediate results; done means equivalent results regardless of merge grouping and preserved overlap-based behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
databases
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.