Skip to content

Feat(parquet) : Introduce support for optionally writing Distinct Values to statistics - #25576

Merged
jayzhan211 merged 10 commits into
apache:mainfrom
Rich-T-kid:rich-T-kid/Introduce-roundtrip-ndv
Oct 9, 2026
Merged

jayzhan211 merged 10 commits into
apache:mainfrom
Rich-T-kid:rich-T-kid/Introduce-roundtrip-ndv

Conversation

@Rich-T-kid

@Rich-T-kid Rich-T-kid commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

DataFusion can use NDV (number of distinct values) from Parquet row group statistics to power optimizations like COUNT(DISTINCT) aggregation pushdown and cardinality estimation (#24111) . However, currently the Parquet writers doesn't emit NDV by default, so these statistics are absent for roundtrip queries.

Two recent arrow-rs additions make it possible to close this gap end-to-end:

This PR wires both into DataFusion so that users can write parquet files with NDV and have DataFusion read them back correctly.

What changes are included in this PR?

  • summarize_distinct_counts now calls StatisticsConverter::row_group_distinct_counts() instead of manually walking column chunk statistics.
  • Adds datafusion.execution.parquet.write_row_group_number_distinct_values (bool, default false) to ParquetOptions. When enabled, the ArrowWriter hashes every non-null value per column per row group and writes the distinct count into the Parquet statistics footer.

What is the testing strategy for this PR?

  • The NDV read change is covered by the 8 existing unit tests in metadata.rs (test_distinct_count_*), which all continue to pass.
  • The write config is tested by sqllogictest/test_files/parquet_ndv_write.slt, which writes 4 rows with 3 distinct values and verifies:
    • Without the option: Distinct is absent from DataSourceExec statistics
    • With the option: Distinct=Exact(3) appears in DataSourceExec statistics

Are there any user-facing changes?

yes, new write_row_group_number_distinct_values to parquet config options.

@github-actions github-actions Bot added sqllogictest SQL Logic Tests (.slt) common Related to common crate proto Related to proto crate datasource Changes to the datasource crate labels Sep 21, 2026
@github-actions

github-actions Bot commented Sep 21, 2026 •

Copy link
Copy Markdown

Thank you for opening this pull request!

Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch).

Details
     Cloning apache/main
    Building datafusion-common v55.1.0 (current)
       Built [  57.916s] (current)
     Parsing datafusion-common v55.1.0 (current)
      Parsed [   0.065s] (current)
    Building datafusion-common v55.1.0 (baseline)
       Built [  35.862s] (baseline)
     Parsing datafusion-common v55.1.0 (baseline)
      Parsed [   0.067s] (baseline)
    Checking datafusion-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.788s] 223 checks: 222 pass, 1 fail, 0 warn, 31 skip

--- failure constructible_struct_adds_field: struct exhaustively constructible through public API adds field ---

Description:
A pub struct that could be exhaustively constructed with a literal using only public API has a new pub field, breaking existing exhaustive literals.
        ref: https://doc.rust-lang.org/reference/expressions/struct-expr.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/constructible_struct_adds_field.ron

Failed in:
  field ParquetOptions.write_row_group_number_distinct_values in /home/runner/work/datafusion/datafusion/datafusion/common/src/config.rs:1315

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  96.556s] datafusion-common
    Building datafusion-datasource-parquet v55.1.0 (current)
       Built [  52.862s] (current)
     Parsing datafusion-datasource-parquet v55.1.0 (current)
      Parsed [   0.033s] (current)
    Building datafusion-datasource-parquet v55.1.0 (baseline)
       Built [  53.340s] (baseline)
     Parsing datafusion-datasource-parquet v55.1.0 (baseline)
      Parsed [   0.034s] (baseline)
    Checking datafusion-datasource-parquet v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.152s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 108.112s] datafusion-datasource-parquet
    Building datafusion-proto-common v55.1.0 (current)
       Built [  24.324s] (current)
     Parsing datafusion-proto-common v55.1.0 (current)
      Parsed [   0.050s] (current)
    Building datafusion-proto-common v55.1.0 (baseline)
       Built [  25.200s] (baseline)
     Parsing datafusion-proto-common v55.1.0 (baseline)
      Parsed [   0.056s] (baseline)
    Checking datafusion-proto-common v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   1.270s] 223 checks: 222 pass, 1 fail, 0 warn, 31 skip

--- failure constructible_struct_adds_field: struct exhaustively constructible through public API adds field ---

Description:
A pub struct that could be exhaustively constructed with a literal using only public API has a new pub field, breaking existing exhaustive literals.
        ref: https://doc.rust-lang.org/reference/expressions/struct-expr.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/constructible_struct_adds_field.ron

Failed in:
  field ParquetOptions.write_row_group_number_distinct_values in /home/runner/work/datafusion/datafusion/datafusion/proto-common/src/generated/prost.rs:874
  field ParquetOptions.write_row_group_number_distinct_values in /home/runner/work/datafusion/datafusion/datafusion/proto-common/src/generated/prost.rs:874
  field ParquetOptions.write_row_group_number_distinct_values in /home/runner/work/datafusion/datafusion/datafusion/proto-common/src/generated/prost.rs:874

     Summary semver requires new major version: 1 major and 0 minor checks failed
    Finished [  52.305s] datafusion-proto-common
    Building datafusion-proto-models v55.1.0 (current)
       Built [  27.647s] (current)
     Parsing datafusion-proto-models v55.1.0 (current)
      Parsed [   0.144s] (current)
    Building datafusion-proto-models v55.1.0 (baseline)
       Built [  27.755s] (baseline)
     Parsing datafusion-proto-models v55.1.0 (baseline)
      Parsed [   0.142s] (baseline)
    Checking datafusion-proto-models v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   1.956s] 223 checks: 221 pass, 2 fail, 0 warn, 31 skip

--- failure constructible_struct_adds_field: struct exhaustively constructible through public API adds field ---

Description:
A pub struct that could be exhaustively constructed with a literal using only public API has a new pub field, breaking existing exhaustive literals.
        ref: https://doc.rust-lang.org/reference/expressions/struct-expr.html
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/constructible_struct_adds_field.ron

Failed in:
  field ParquetOptions.write_row_group_number_distinct_values in /home/runner/work/datafusion/datafusion/datafusion/proto-models/src/generated/datafusion_proto_common.rs:874
  field ParquetOptions.write_row_group_number_distinct_values in /home/runner/work/datafusion/datafusion/datafusion/proto-models/src/generated/datafusion_proto_common.rs:874

--- failure struct_pub_field_missing: pub struct's pub field removed or renamed ---

Description:
A publicly-visible struct has at least one public field that is no longer available under its prior name. It may have been renamed or removed entirely.
        ref: https://doc.rust-lang.org/cargo/reference/semver.html#item-remove
       impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.50.0/src/lints/struct_pub_field_missing.ron

Failed in:
  field schema of struct ProjectionExecNode, previously in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/ab3c7d1dd038ad81ee9955865aaf882140668cb7/datafusion/proto-models/src/generated/prost.rs:2310
  field schema of struct ProjectionExecNode, previously in file /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/ab3c7d1dd038ad81ee9955865aaf882140668cb7/datafusion/proto-models/src/generated/prost.rs:2310

     Summary semver requires new major version: 2 major and 0 minor checks failed
    Finished [  59.534s] datafusion-proto-models
    Building datafusion-sqllogictest v55.1.0 (current)
       Built [ 107.741s] (current)
     Parsing datafusion-sqllogictest v55.1.0 (current)
      Parsed [   0.017s] (current)
    Building datafusion-sqllogictest v55.1.0 (baseline)
       Built [ 102.403s] (baseline)
     Parsing datafusion-sqllogictest v55.1.0 (baseline)
      Parsed [   0.017s] (baseline)
    Checking datafusion-sqllogictest v55.1.0 -> v55.1.0 (no change; assume patch)
     Checked [   0.108s] 223 checks: 223 pass, 31 skip
     Summary no semver update required
    Finished [ 213.448s] datafusion-sqllogictest

@github-actions github-actions Bot added the auto detected api change Auto detected API change label Sep 21, 2026
@Rich-T-kid
Rich-T-kid force-pushed the rich-T-kid/Introduce-roundtrip-ndv branch from 6aac80c to 085ec41 Compare September 21, 2026 18:19
@github-actions github-actions Bot added the core Core DataFusion crate label Sep 21, 2026
@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

@adriangb in case your interested

@adriangb adriangb left a comment

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.

Very cool! I'm interested in what algorithm is being used to calculate distinct values, what cost/perf impact it has, etc. But I suppose that (1) that's in arrow-rs and (2) we can always change it later w/o any impact on these new options / public API.

Comment thread datafusion/proto-common/src/generated/pbjson.rs
@adriangb

adriangb commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

A couple of questions:

  1. Does the spec say this value is / must be Exact or could it be an estimate? An estimate sounds much better to me overall, it's just as useful for the planner and I don't think an exact value can be used for e.g. select count(...) ... if there is a filter or more than 1 row group or more than 1 file (which is a lot of queries in practice).
  2. Could these values be derived from the dictionary size or bloom filters at read time? There's some interesting things that could be done e.g. you can even take the bloom filters, OR them and then derive an approximate distinct count for a single file with multiple row groups or even across files. You can't combine dictionary counts, but you could sum them or if willing to pay some cost upfront combine the dictionary hashes. This would be a totally different feature, and would be hard tradeoff between planning cost and execution cost (a whole can of worms).
  3. At write time could we do something similar and derive the value from the dictionary or bloom filters we are already writing, or prepare a bloom filter just to compute an estimated NDV? My estimate is that the current implementation uses ~9MB per column at the default row group size, a bloom filter would use ~1MB per column, would be cheaper to write and would have <1% error (this is all at the default NDV/FPP configuration) (not to mention it's free if you're already writing a bloom filter).

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author
  1. No, from my understanding it doesn't need to be exact. Even if it were, it wouldn't be any more useful than an estimation, for the reasons you mentioned. A parquet file with two row groups with an NDV of 100 and 101 respectively can have between 101 and 201 unique values.
  2. for bloom filters I think this is a good idea but from the arrow spec, dictionary arrays aren't guaranteed to be unique.
dictionary = keys:[0,1,2,3] ,values=["cat","cat","cat","cat"]

is a valid dictionary array. I think the bloom filter idea is the best approach to derive the ndv as its pretty expensive to get an exact count. With that being said, many work loads are write once read thousands of times so the trade off in practice is generally worth it. especially if we can take advantage of low ndv columns (👀 #24111)
3. +1 with this. i'm actually about to open up a PR that adds benchmarks for calculating the ndv. it may be interesting to open an issue with possible optimization ideas.

@adriangb

Copy link
Copy Markdown
Contributor

No, from my understanding it doesn't need to be exact. Even if it were, it wouldn't be any more useful than an estimation, for the reasons you mentioned. A parquet file with two row groups with an NDV of 100 and 101 respectively can have between 101 and 201 unique values.

In that case, it's probably never safe or even helpful to read it in as Exact, should we always read it in as Inexact?

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

In that case, it's probably never safe or even helpful to read it in as Exact, should we always read it in as Inexact?

yes thats what the code currently does. in datafusion/datasource-parquet/src/metadata.rs if there is more than one row-group the column ndv is Inexact due to the reason we went over above. in the case where theres only one row-group then the NDV is actually exactly correct.

@adriangb

Copy link
Copy Markdown
Contributor

my point is that even in the single row group case maybe it should be inexact? if the write used a bloom filter as discussed above it could be 1% off (in either direction i believe)

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

my point is that even in the single row group case maybe it should be inexact? if the write used a bloom filter as discussed above it could be 1% off (in either direction i believe)

@adriangb ahh okay I see what you mean now. makes sense to me, updating the PR

@Rich-T-kid
Rich-T-kid force-pushed the rich-T-kid/Introduce-roundtrip-ndv branch from b1fead6 to fd3fbe5 Compare September 22, 2026 00:25
@github-actions github-actions Bot added documentation Improvements or additions to documentation and removed core Core DataFusion crate labels Sep 22, 2026
@Rich-T-kid
Rich-T-kid force-pushed the rich-T-kid/Introduce-roundtrip-ndv branch from 5977456 to d3fc56a Compare September 22, 2026 04:53
@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

hmm 🤔 seems the ci is failing because the AggregateStatistics optimizer only folds COUNT(DISTINCT) to a constant when distinct_count is Precision::Exact. changing the ndv to Inexact may be correct but it may affect other piece of datafusion. will take a closer look tomorrow

@codecov-commenter

codecov-commenter commented Sep 22, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 70.00000% with 12 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.74%. Comparing base (4d167a1) to head (233a1c4).
⚠️ Report is 86 commits behind head on main.

Files with missing lines Patch % Lines
datafusion/proto-common/src/generated/pbjson.rs 41.66% 5 Missing and 2 partials ⚠️
datafusion/datasource-parquet/src/file_format.rs 0.00% 2 Missing and 1 partial ⚠️
datafusion/datasource-parquet/src/metadata.rs 91.66% 1 Missing ⚠️
datafusion/proto-models/src/from_proto.rs 50.00% 0 Missing and 1 partial ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25576      +/-   ##
==========================================
+ Coverage   82.60%   82.74%   +0.14%     
==========================================
  Files        1145     1147       +2     
  Lines      444057   449791    +5734     
  Branches   444057   449791    +5734     
==========================================
+ Hits       366796   372167    +5371     
+ Misses      54996    54950      -46     
- Partials    22265    22674     +409     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@asolimando

Copy link
Copy Markdown
Member

my point is that even in the single row group case maybe it should be inexact? if the write used a bloom filter as discussed above it could be 1% off (in either direction i believe)

We have been working so far under the assumption that distinct_count is exact in Parquet (see apache/arrow-rs#8608 (comment) from @alamb), has this changed in the meantime? If so, I agree that we should consume it as Inexact everywhere, as we have rules relying on that for correctness as @Rich-T-kid has noticed.

@adriangb

Copy link
Copy Markdown
Contributor

I commented on that issue: apache/arrow-rs#8608 (comment)

@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

is there any update on this? @adriangb

@alamb alamb changed the title Feat(parquet) : Introduce roundtrip number Distinct Valeus Feat(parquet) : Introduce roundtrip number Distinct Values Sep 24, 2026
@alamb alamb changed the title Feat(parquet) : Introduce roundtrip number Distinct Values Feat(parquet) : Introduce support for optionally writing Distinct Values to statistics Sep 24, 2026
…action

Replace manual row group column stat traversal in summarize_distinct_counts
with the new arrow-rs StatisticsConverter::row_group_distinct_counts API,
keeping the existing threshold and exactness logic intact.
Add a new ParquetOptions field that enables writing NDV (number of distinct
values) statistics per column when writing parquet files, backed by the
new arrow-rs WriterPropertiesBuilder::set_write_row_group_number_distinct_values
API. Includes proto serialization and a SQL logic roundtrip test.
- `summarize_distinct_counts`: always returns `Precision::Inexact` — NDV
  from parquet statistics is never exact across queries with filters or
  multiple files, even for a single row group
- Update `parquet_ndv_write.slt` to assert `Distinct=Inexact(3)`
- Add `write_row_group_number_distinct_values` to `information_schema.slt`
  and regenerate `configs.md` to fix the failing CI checks
…write_row_group_number_distinct_values

Two CI fixes:
- parquet_writer.rs: read write_row_group_number_distinct_values back from
  WriterProperties in session_config_from_writer_props instead of using the default
- metadata.rs: restore Precision::Exact for single-row-group NDV so the
  aggregate_statistics optimizer can fold COUNT(DISTINCT) to constants
NDV from parquet statistics is always Inexact, so the AggregateStatistics
optimizer cannot fold COUNT(DISTINCT) to a constant. Update the EXPLAIN
plan to match the actual AggregateExec output.
@Rich-T-kid
Rich-T-kid force-pushed the rich-T-kid/Introduce-roundtrip-ndv branch from bacf5b4 to 3ed9076 Compare October 2, 2026 09:15
@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

took a second look at thread Adrian proposed and I think its fine for us to assume its Inexact for the reasons mentioned above. there was a mention of possibly adding another parquet field but if/when that comes we can update the code then.

@jayzhan211 jayzhan211 left a comment

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.

Thanks @Rich-T-kid , here are 2 suggestions

Comment thread datafusion/datasource-parquet/src/metadata.rs
Comment thread datafusion/common/src/config.rs
@Rich-T-kid
Rich-T-kid requested a review from jayzhan211 October 5, 2026 05:00
@jayzhan211

Copy link
Copy Markdown
Contributor

tests are not passed

Tag 39 collides with row_group_range_assignment added in main after
this branch was created. Move to the next free tag and regenerate
prost bindings in both proto-common and proto-models.
@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

tests are not passed

@jayzhan211 this should be good now

@jayzhan211

Copy link
Copy Markdown
Contributor

Merge the main since:

Tag 40 is taken on main since #24227 (bool enable_rle_to_dictionary = 40;), which merged after this branch's last push. The merge is textually clean, so GitHub shows it as mergeable and CI is green, but the merged tree doesn't compile:

$ cargo check -p datafusion-proto-common   # this branch merged with current main
error: proc-macro derive panicked
  = help: message: called `Result::unwrap()` on an `Err` value: message ParquetOptions has multiple fields with tag 40

Fix: rebase on main, take the next free tag, and regenerate both generated copies:

-  bool write_row_group_number_distinct_values = 40; // default = false
+  bool write_row_group_number_distinct_values = 41; // default = false

@Rich-T-kid
Rich-T-kid force-pushed the rich-T-kid/Introduce-roundtrip-ndv branch from a165d45 to efeb978 Compare October 7, 2026 15:41
@Rich-T-kid

Copy link
Copy Markdown
Contributor Author

@jayzhan211 should be good to go now

@jayzhan211
jayzhan211 added this pull request to the merge queue Oct 9, 2026
@jayzhan211

Copy link
Copy Markdown
Contributor

Thanks @Rich-T-kid @adriangb !

Merged via the queue into apache:main with commit f84b82c Oct 9, 2026
43 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

auto detected api change Auto detected API change common Related to common crate datasource Changes to the datasource crate documentation Improvements or additions to documentation proto Related to proto crate sqllogictest SQL Logic Tests (.slt) v56.0.0

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Improve NDV statistics for Parquet columns (parquet RLE & Dictionary array focus)

5 participants