Repository navigation
Feat(parquet) : Introduce support for optionally writing Distinct Values to statistics - #25576
Conversation
|
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 |
6aac80c to
085ec41
Compare
|
@adriangb in case your interested |
adriangb
left a comment
There was a problem hiding this comment.
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.
|
A couple of questions:
|
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) |
In that case, it's probably never safe or even helpful to read it in as |
yes thats what the code currently does. in |
|
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 |
b1fead6 to
fd3fbe5
Compare
5977456 to
d3fc56a
Compare
|
hmm 🤔 seems the ci is failing because the AggregateStatistics optimizer only folds |
Codecov Report❌ Patch coverage is 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. 🚀 New features to boost your workflow:
|
We have been working so far under the assumption that |
|
I commented on that issue: apache/arrow-rs#8608 (comment) |
|
is there any update on this? @adriangb |
…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.
bacf5b4 to
3ed9076
Compare
|
took a second look at thread Adrian proposed and I think its fine for us to assume its |
jayzhan211
left a comment
There was a problem hiding this comment.
Thanks @Rich-T-kid , here are 2 suggestions
|
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.
@jayzhan211 this should be good now |
|
Merge the main since: Tag 40 is taken on $ 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 40Fix: rebase on - bool write_row_group_number_distinct_values = 40; // default = false
+ bool write_row_group_number_distinct_values = 41; // default = false |
a165d45 to
efeb978
Compare
|
@jayzhan211 should be good to go now |
|
Thanks @Rich-T-kid @adriangb ! |
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:
StatisticsConverter::row_group_distinct_counts(): a typed API for reading NDV from row group metadataWriterPropertiesBuilder::set_write_row_group_number_distinct_values(): a writer-side option to track and emit NDV during writeThis 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?
StatisticsConverter::row_group_distinct_counts()instead of manually walking column chunk statistics.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?
Are there any user-facing changes?
yes, new
write_row_group_number_distinct_valuesto parquet config options.