From b9df4d6266b882212698c2aa4c1e67f22d315861 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 21 Sep 2026 11:37:42 -0400 Subject: [PATCH 01/10] feat: use StatisticsConverter::row_group_distinct_counts for NDV extraction 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. --- datafusion/datasource-parquet/src/metadata.rs | 32 +++++++++---------- 1 file changed, 15 insertions(+), 17 deletions(-) diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index 51dc98e4f2f02..a418deccad035 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -779,7 +779,7 @@ fn summarize_column_statistics( summarize_null_counts(stats_converter, row_groups_metadata)?; accumulators.distinct_counts_array[logical_schema_index] = - summarize_distinct_counts(parquet_index, row_groups_metadata); + summarize_distinct_counts(stats_converter, row_groups_metadata)?; let arrow_field = logical_file_schema.field(logical_schema_index); accumulators.column_byte_sizes[logical_schema_index] = compute_arrow_column_size( @@ -930,29 +930,27 @@ where /// Extract distinct counts from row group column statistics. fn summarize_distinct_counts( - parquet_idx: Option, + stats_converter: &StatisticsConverter, row_groups_metadata: &[RowGroupMetaData], -) -> Precision { - let Some(parquet_idx) = parquet_idx else { - return Precision::Absent; - }; +) -> Result> { + if stats_converter.parquet_column_index().is_none() { + return Ok(Precision::Absent); + } let num_row_groups = row_groups_metadata.len(); if num_row_groups == 0 { - return Precision::Absent; + return Ok(Precision::Absent); } let required_count = (num_row_groups as f64 * PARTIAL_NDV_THRESHOLD).ceil() as usize; + let distinct_counts = + stats_converter.row_group_distinct_counts(row_groups_metadata)?; + let mut ndv_count = 0; let mut max_distinct_count: Option = None; - for (row_group_idx, row_group) in row_groups_metadata.iter().enumerate() { - if let Some(distinct_count) = row_group - .columns() - .get(parquet_idx) - .and_then(|col| col.statistics()) - .and_then(|stats| stats.distinct_count_opt()) - { + for (row_group_idx, value) in distinct_counts.iter().enumerate() { + if let Some(distinct_count) = value { ndv_count += 1; max_distinct_count = Some(match max_distinct_count { Some(max) => max.max(distinct_count), @@ -963,17 +961,17 @@ fn summarize_distinct_counts( // Return early if there's no chance to reach the required coverage. let remaining = num_row_groups - row_group_idx - 1; if ndv_count + remaining < required_count { - return Precision::Absent; + return Ok(Precision::Absent); } } - match max_distinct_count { + Ok(match max_distinct_count { Some(distinct_count) if num_row_groups == 1 => { Precision::Exact(distinct_count as usize) } Some(distinct_count) => Precision::Inexact(distinct_count as usize), None => Precision::Absent, - } + }) } /// Compute the Arrow in-memory size for a single column From 3ad4cd80e8d9bfda0de712ff3e6f8c7d08676dc4 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 21 Sep 2026 12:01:17 -0400 Subject: [PATCH 02/10] feat: add write_row_group_number_distinct_values parquet config option 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. --- datafusion/common/src/config.rs | 6 ++ .../common/src/file_options/parquet_writer.rs | 10 +- .../datasource-parquet/src/file_format.rs | 3 + .../proto/datafusion_common.proto | 1 + datafusion/proto-common/src/from_proto/mod.rs | 1 + .../proto-common/src/generated/pbjson.rs | 24 ++++- .../proto-common/src/generated/prost.rs | 3 + datafusion/proto-common/src/to_proto/mod.rs | 1 + datafusion/proto-models/src/from_proto.rs | 2 + .../src/generated/datafusion_proto_common.rs | 3 + .../test_files/parquet_ndv_write.slt | 97 +++++++++++++++++++ 11 files changed, 147 insertions(+), 4 deletions(-) create mode 100644 datafusion/sqllogictest/test_files/parquet_ndv_write.slt diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 4246dcc451d45..80e08c5630e22 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -1514,6 +1514,12 @@ config_namespace! { /// default parquet writer setting pub bloom_filter_ndv: Option, default = None + /// (writing) Write the number of distinct values (NDV) for each column + /// in row group statistics when creating parquet files. Enabling this + /// improves NDV-based query optimizations at the cost of hashing every + /// non-null value during write. + pub write_row_group_number_distinct_values: bool, default = false + /// (writing) Controls whether DataFusion will attempt to speed up writing /// parquet files by serializing them in parallel. Each column /// in each row group in each output file are serialized in parallel diff --git a/datafusion/common/src/file_options/parquet_writer.rs b/datafusion/common/src/file_options/parquet_writer.rs index af4554c4eef84..4bdea56586dd8 100644 --- a/datafusion/common/src/file_options/parquet_writer.rs +++ b/datafusion/common/src/file_options/parquet_writer.rs @@ -228,6 +228,7 @@ impl ParquetOptions { bloom_filter_on_write, bloom_filter_fpp, bloom_filter_ndv, + write_row_group_number_distinct_values, content_defined_chunking, // not in WriterProperties @@ -267,7 +268,10 @@ impl ParquetOptions { .set_column_index_truncate_length(*column_index_truncate_length) .set_statistics_truncate_length(*statistics_truncate_length) .set_data_page_row_count_limit(*data_page_row_count_limit) - .set_bloom_filter_enabled(*bloom_filter_on_write); + .set_bloom_filter_enabled(*bloom_filter_on_write) + .set_write_row_group_number_distinct_values( + *write_row_group_number_distinct_values, + ); if let Some(bloom_filter_fpp) = bloom_filter_fpp { builder = builder.set_bloom_filter_fpp(*bloom_filter_fpp); @@ -404,6 +408,8 @@ mod tests { bloom_filter_on_write: !defaults.bloom_filter_on_write, bloom_filter_fpp: Some(0.42), bloom_filter_ndv: Some(42), + write_row_group_number_distinct_values: !defaults + .write_row_group_number_distinct_values, // not in WriterProperties, but itemizing here to not skip newly added props enable_page_index: defaults.enable_page_index, @@ -551,6 +557,8 @@ mod tests { schema_force_view_types: global_options_defaults.schema_force_view_types, binary_as_string: global_options_defaults.binary_as_string, skip_arrow_metadata: global_options_defaults.skip_arrow_metadata, + write_row_group_number_distinct_values: global_options_defaults + .write_row_group_number_distinct_values, coerce_int96: None, coerce_int96_tz: None, content_defined_chunking: props.content_defined_chunking().into(), diff --git a/datafusion/datasource-parquet/src/file_format.rs b/datafusion/datasource-parquet/src/file_format.rs index dc0f9a7be434e..3ec1196d0c6c8 100644 --- a/datafusion/datasource-parquet/src/file_format.rs +++ b/datafusion/datasource-parquet/src/file_format.rs @@ -706,6 +706,9 @@ impl From<&ParquetFormatFactory> for protobuf::TableParquetOptions { schema_force_view_types: global_options.global.schema_force_view_types, binary_as_string: global_options.global.binary_as_string, skip_arrow_metadata: global_options.global.skip_arrow_metadata, + write_row_group_number_distinct_values: global_options + .global + .write_row_group_number_distinct_values, coerce_int96_opt: global_options.global.coerce_int96.map(|compression| { parquet_options::CoerceInt96Opt::CoerceInt96(compression) }), diff --git a/datafusion/proto-common/proto/datafusion_common.proto b/datafusion/proto-common/proto/datafusion_common.proto index bdabb0777ef81..e74211d01edc5 100644 --- a/datafusion/proto-common/proto/datafusion_common.proto +++ b/datafusion/proto-common/proto/datafusion_common.proto @@ -579,6 +579,7 @@ message ParquetOptions { bool schema_force_view_types = 28; // default = false bool binary_as_string = 29; // default = false bool skip_arrow_metadata = 30; // default = false + bool write_row_group_number_distinct_values = 39; // default = false oneof metadata_size_hint_opt { uint64 metadata_size_hint = 4; diff --git a/datafusion/proto-common/src/from_proto/mod.rs b/datafusion/proto-common/src/from_proto/mod.rs index 0bde696c7c46b..03a38fe4ca9ea 100644 --- a/datafusion/proto-common/src/from_proto/mod.rs +++ b/datafusion/proto-common/src/from_proto/mod.rs @@ -1278,6 +1278,7 @@ impl TryFrom<&protobuf::ParquetOptions> for ParquetOptions { protobuf::parquet_options::CoerceInt96TzOpt::CoerceInt96Tz(v) => Some(v), }).unwrap_or(None), skip_arrow_metadata: value.skip_arrow_metadata, + write_row_group_number_distinct_values: value.write_row_group_number_distinct_values, max_predicate_cache_size: value.max_predicate_cache_size_opt.map(|opt| match opt { protobuf::parquet_options::MaxPredicateCacheSizeOpt::MaxPredicateCacheSize(v) => { to_usize(v, "max_predicate_cache_size") diff --git a/datafusion/proto-common/src/generated/pbjson.rs b/datafusion/proto-common/src/generated/pbjson.rs index a1d86c5fe3187..dfa8e93883f17 100644 --- a/datafusion/proto-common/src/generated/pbjson.rs +++ b/datafusion/proto-common/src/generated/pbjson.rs @@ -2600,7 +2600,7 @@ impl<'de> serde::Deserialize<'de> for CsvWriterOptions { if compression_level__.is_some() { return Err(serde::de::Error::duplicate_field("compressionLevel")); } - compression_level__ = + compression_level__ = map_.next_value::<::std::option::Option<::pbjson::private::NumberDeserialize<_>>>()?.map(|x| x.0) ; } @@ -2614,7 +2614,7 @@ impl<'de> serde::Deserialize<'de> for CsvWriterOptions { if terminator__.is_some() { return Err(serde::de::Error::duplicate_field("terminator")); } - terminator__ = + terminator__ = Some(map_.next_value::<::pbjson::private::BytesDeserialize<_>>()?.0) ; } @@ -4056,7 +4056,7 @@ impl serde::Serialize for ExplainAnalyzeCategoriesNode { struct_ser.serialize_field("all", &self.all)?; } if !self.only.is_empty() { - let v = self.only.iter().copied().map(|v| { + let v = self.only.iter().cloned().map(|v| { MetricCategory::try_from(v) .map_err(|_| serde::ser::Error::custom(format!("Invalid variant {}", v))) }).collect::, _>>()?; @@ -6479,6 +6479,9 @@ impl serde::Serialize for ParquetOptions { if self.skip_arrow_metadata { len += 1; } + if self.write_row_group_number_distinct_values { + len += 1; + } if self.dictionary_page_size_limit != 0 { len += 1; } @@ -6596,6 +6599,9 @@ impl serde::Serialize for ParquetOptions { if self.skip_arrow_metadata { struct_ser.serialize_field("skipArrowMetadata", &self.skip_arrow_metadata)?; } + if self.write_row_group_number_distinct_values { + struct_ser.serialize_field("writeRowGroupNumberDistinctValues", &self.write_row_group_number_distinct_values)?; + } if self.dictionary_page_size_limit != 0 { #[allow(clippy::needless_borrow)] #[allow(clippy::needless_borrows_for_generic_args)] @@ -6768,6 +6774,8 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { "binaryAsString", "skip_arrow_metadata", "skipArrowMetadata", + "write_row_group_number_distinct_values", + "writeRowGroupNumberDistinctValues", "dictionary_page_size_limit", "dictionaryPageSizeLimit", "data_page_row_count_limit", @@ -6825,6 +6833,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { SchemaForceViewTypes, BinaryAsString, SkipArrowMetadata, + WriteRowGroupNumberDistinctValues, DictionaryPageSizeLimit, DataPageRowCountLimit, MaxRowGroupSize, @@ -6882,6 +6891,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { "schemaForceViewTypes" | "schema_force_view_types" => Ok(GeneratedField::SchemaForceViewTypes), "binaryAsString" | "binary_as_string" => Ok(GeneratedField::BinaryAsString), "skipArrowMetadata" | "skip_arrow_metadata" => Ok(GeneratedField::SkipArrowMetadata), + "writeRowGroupNumberDistinctValues" | "write_row_group_number_distinct_values" => Ok(GeneratedField::WriteRowGroupNumberDistinctValues), "dictionaryPageSizeLimit" | "dictionary_page_size_limit" => Ok(GeneratedField::DictionaryPageSizeLimit), "dataPageRowCountLimit" | "data_page_row_count_limit" => Ok(GeneratedField::DataPageRowCountLimit), "maxRowGroupSize" | "max_row_group_size" => Ok(GeneratedField::MaxRowGroupSize), @@ -6937,6 +6947,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { let mut schema_force_view_types__ = None; let mut binary_as_string__ = None; let mut skip_arrow_metadata__ = None; + let mut write_row_group_number_distinct_values__ = None; let mut dictionary_page_size_limit__ = None; let mut data_page_row_count_limit__ = None; let mut max_row_group_size__ = None; @@ -7068,6 +7079,12 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { } skip_arrow_metadata__ = Some(map_.next_value()?); } + GeneratedField::WriteRowGroupNumberDistinctValues => { + if write_row_group_number_distinct_values__.is_some() { + return Err(serde::de::Error::duplicate_field("writeRowGroupNumberDistinctValues")); + } + write_row_group_number_distinct_values__ = Some(map_.next_value()?); + } GeneratedField::DictionaryPageSizeLimit => { if dictionary_page_size_limit__.is_some() { return Err(serde::de::Error::duplicate_field("dictionaryPageSizeLimit")); @@ -7210,6 +7227,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { schema_force_view_types: schema_force_view_types__.unwrap_or_default(), binary_as_string: binary_as_string__.unwrap_or_default(), skip_arrow_metadata: skip_arrow_metadata__.unwrap_or_default(), + write_row_group_number_distinct_values: write_row_group_number_distinct_values__.unwrap_or_default(), dictionary_page_size_limit: dictionary_page_size_limit__.unwrap_or_default(), data_page_row_count_limit: data_page_row_count_limit__.unwrap_or_default(), max_row_group_size: max_row_group_size__.unwrap_or_default(), diff --git a/datafusion/proto-common/src/generated/prost.rs b/datafusion/proto-common/src/generated/prost.rs index f856c1aa3959c..0a47feae8c659 100644 --- a/datafusion/proto-common/src/generated/prost.rs +++ b/datafusion/proto-common/src/generated/prost.rs @@ -867,6 +867,9 @@ pub struct ParquetOptions { /// default = false #[prost(bool, tag = "30")] pub skip_arrow_metadata: bool, + /// default = false + #[prost(bool, tag = "39")] + pub write_row_group_number_distinct_values: bool, #[prost(uint64, tag = "12")] pub dictionary_page_size_limit: u64, #[prost(uint64, tag = "18")] diff --git a/datafusion/proto-common/src/to_proto/mod.rs b/datafusion/proto-common/src/to_proto/mod.rs index b8eb46f2c8da8..7470db0adbc4b 100644 --- a/datafusion/proto-common/src/to_proto/mod.rs +++ b/datafusion/proto-common/src/to_proto/mod.rs @@ -963,6 +963,7 @@ impl TryFrom<&ParquetOptions> for protobuf::ParquetOptions { schema_force_view_types: value.schema_force_view_types, binary_as_string: value.binary_as_string, skip_arrow_metadata: value.skip_arrow_metadata, + write_row_group_number_distinct_values: value.write_row_group_number_distinct_values, coerce_int96_opt: value.coerce_int96.clone().map(protobuf::parquet_options::CoerceInt96Opt::CoerceInt96), coerce_int96_tz_opt: value.coerce_int96_tz.clone().map(protobuf::parquet_options::CoerceInt96TzOpt::CoerceInt96Tz), max_predicate_cache_size_opt: value.max_predicate_cache_size.map(|v| protobuf::parquet_options::MaxPredicateCacheSizeOpt::MaxPredicateCacheSize(v as u64)), diff --git a/datafusion/proto-models/src/from_proto.rs b/datafusion/proto-models/src/from_proto.rs index fdf32396475b7..378625cd9b77d 100644 --- a/datafusion/proto-models/src/from_proto.rs +++ b/datafusion/proto-models/src/from_proto.rs @@ -465,6 +465,8 @@ impl TryFrom<&ParquetOptionsProto> for ParquetOptions { schema_force_view_types: proto.schema_force_view_types, binary_as_string: proto.binary_as_string, skip_arrow_metadata: proto.skip_arrow_metadata, + write_row_group_number_distinct_values: proto + .write_row_group_number_distinct_values, coerce_int96: proto.coerce_int96_opt.as_ref().map(|opt| match opt { parquet_options::CoerceInt96Opt::CoerceInt96(coerce_int96) => { coerce_int96.clone() diff --git a/datafusion/proto-models/src/generated/datafusion_proto_common.rs b/datafusion/proto-models/src/generated/datafusion_proto_common.rs index f856c1aa3959c..0a47feae8c659 100644 --- a/datafusion/proto-models/src/generated/datafusion_proto_common.rs +++ b/datafusion/proto-models/src/generated/datafusion_proto_common.rs @@ -867,6 +867,9 @@ pub struct ParquetOptions { /// default = false #[prost(bool, tag = "30")] pub skip_arrow_metadata: bool, + /// default = false + #[prost(bool, tag = "39")] + pub write_row_group_number_distinct_values: bool, #[prost(uint64, tag = "12")] pub dictionary_page_size_limit: u64, #[prost(uint64, tag = "18")] diff --git a/datafusion/sqllogictest/test_files/parquet_ndv_write.slt b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt new file mode 100644 index 0000000000000..8318b3799a098 --- /dev/null +++ b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt @@ -0,0 +1,97 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +# Roundtrip test for write_row_group_number_distinct_values. +# +# Writes a parquet file with NDV tracking enabled and verifies that +# DataFusion reads back the distinct count from the parquet statistics. + +statement ok +set datafusion.explain.physical_plan_only = true; + +statement ok +set datafusion.explain.show_statistics = true; + +statement ok +set datafusion.execution.collect_statistics = true; + +###### +# Without the option (default), distinct counts are absent. +###### + +query I +COPY (VALUES (1), (2), (3), (1)) +TO 'test_files/scratch/parquet_ndv_write/no_ndv.parquet' +STORED AS PARQUET; +---- +4 + +statement ok +CREATE EXTERNAL TABLE no_ndv +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_ndv_write/no_ndv.parquet'; + +# Distinct count should be absent (no NDV written to the file) +query TT +EXPLAIN SELECT * FROM no_ndv; +---- +physical_plan DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/parquet_ndv_write/no_ndv.parquet]]}, projection=[column1], file_type=parquet, statistics=[Rows=Exact(4), Bytes=Exact(32), [(Col[0]: Min=Exact(Int64(1)) Max=Exact(Int64(3)) Null=Exact(0) ScanBytes=Exact(32))]] + +statement ok +DROP TABLE no_ndv; + +###### +# With write_row_group_number_distinct_values = true, distinct count is written +# and read back as Exact (single row group). +###### + +statement ok +set datafusion.execution.parquet.write_row_group_number_distinct_values = true; + +query I +COPY (VALUES (1), (2), (3), (1)) +TO 'test_files/scratch/parquet_ndv_write/with_ndv.parquet' +STORED AS PARQUET; +---- +4 + +statement ok +CREATE EXTERNAL TABLE with_ndv +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_ndv_write/with_ndv.parquet'; + +# Distinct count should now be Exact(3) — 3 distinct values among 4 rows +query TT +EXPLAIN SELECT * FROM with_ndv; +---- +physical_plan DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/parquet_ndv_write/with_ndv.parquet]]}, projection=[column1], file_type=parquet, statistics=[Rows=Exact(4), Bytes=Exact(32), [(Col[0]: Min=Exact(Int64(1)) Max=Exact(Int64(3)) Null=Exact(0) Distinct=Exact(3) ScanBytes=Exact(32))]] + +statement ok +DROP TABLE with_ndv; + +# Reset settings +statement ok +RESET datafusion.execution.parquet.write_row_group_number_distinct_values; + +statement ok +RESET datafusion.explain.show_statistics; + +statement ok +RESET datafusion.explain.physical_plan_only; + +statement ok +RESET datafusion.execution.collect_statistics; From 0dbb380b12e2f914a753b8865bad7e37124463f0 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 21 Sep 2026 19:25:39 -0400 Subject: [PATCH 03/10] fix: change NDV row-group precision from Exact to Inexact and fix CI MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - `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 --- datafusion/datasource-parquet/src/metadata.rs | 3 --- datafusion/sqllogictest/test_files/information_schema.slt | 2 ++ datafusion/sqllogictest/test_files/parquet_ndv_write.slt | 6 +++--- docs/source/user-guide/configs.md | 1 + 4 files changed, 6 insertions(+), 6 deletions(-) diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index a418deccad035..abf711c335cbf 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -966,9 +966,6 @@ fn summarize_distinct_counts( } Ok(match max_distinct_count { - Some(distinct_count) if num_row_groups == 1 => { - Precision::Exact(distinct_count as usize) - } Some(distinct_count) => Precision::Inexact(distinct_count as usize), None => Precision::Absent, }) diff --git a/datafusion/sqllogictest/test_files/information_schema.slt b/datafusion/sqllogictest/test_files/information_schema.slt index f6633a165df18..e7e31a9b68a55 100644 --- a/datafusion/sqllogictest/test_files/information_schema.slt +++ b/datafusion/sqllogictest/test_files/information_schema.slt @@ -269,6 +269,7 @@ datafusion.execution.parquet.skip_metadata true datafusion.execution.parquet.statistics_enabled page datafusion.execution.parquet.statistics_truncate_length 64 datafusion.execution.parquet.write_batch_size 1024 +datafusion.execution.parquet.write_row_group_number_distinct_values false datafusion.execution.parquet.writer_version 1.0 datafusion.execution.perfect_hash_join_min_key_density 0.15 datafusion.execution.perfect_hash_join_small_build_threshold 1024 @@ -431,6 +432,7 @@ datafusion.execution.parquet.skip_metadata true (reading) If true, the parquet r datafusion.execution.parquet.statistics_enabled page (writing) Sets if statistics are enabled for any column Valid values are: "none", "chunk", and "page" These values are not case sensitive. If NULL, uses default parquet writer setting datafusion.execution.parquet.statistics_truncate_length 64 (writing) Sets statistics truncate length. If NULL, uses default parquet writer setting datafusion.execution.parquet.write_batch_size 1024 (writing) Sets write_batch_size in rows +datafusion.execution.parquet.write_row_group_number_distinct_values false (writing) Write the number of distinct values (NDV) for each column in row group statistics when creating parquet files. Enabling this improves NDV-based query optimizations at the cost of hashing every non-null value during write. datafusion.execution.parquet.writer_version 1.0 (writing) Sets parquet writer version valid values are "1.0" and "2.0" datafusion.execution.perfect_hash_join_min_key_density 0.15 The minimum required density of join keys on the build side to consider a perfect hash join (see `HashJoinExec` for more details). Density is calculated as: `(number of rows) / (max_key - min_key + 1)`. A perfect hash join may be used if the actual key density > this value. Currently only supports cases where build_side.num_rows() < u32::MAX. Support for build_side.num_rows() >= u32::MAX will be added in the future. datafusion.execution.perfect_hash_join_small_build_threshold 1024 A perfect hash join (see `HashJoinExec` for more details) will be considered if the range of keys (max - min) on the build side is < this threshold. This provides a fast path for joins with very small key ranges, bypassing the density check. Currently only supports cases where build_side.num_rows() < u32::MAX. Support for build_side.num_rows() >= u32::MAX will be added in the future. diff --git a/datafusion/sqllogictest/test_files/parquet_ndv_write.slt b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt index 8318b3799a098..4c8fa7e356405 100644 --- a/datafusion/sqllogictest/test_files/parquet_ndv_write.slt +++ b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt @@ -56,7 +56,7 @@ DROP TABLE no_ndv; ###### # With write_row_group_number_distinct_values = true, distinct count is written -# and read back as Exact (single row group). +# and read back as Inexact (NDV from parquet statistics is never treated as exact). ###### statement ok @@ -74,11 +74,11 @@ CREATE EXTERNAL TABLE with_ndv STORED AS PARQUET LOCATION 'test_files/scratch/parquet_ndv_write/with_ndv.parquet'; -# Distinct count should now be Exact(3) — 3 distinct values among 4 rows +# Distinct count should now be Inexact(3) — 3 distinct values among 4 rows query TT EXPLAIN SELECT * FROM with_ndv; ---- -physical_plan DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/parquet_ndv_write/with_ndv.parquet]]}, projection=[column1], file_type=parquet, statistics=[Rows=Exact(4), Bytes=Exact(32), [(Col[0]: Min=Exact(Int64(1)) Max=Exact(Int64(3)) Null=Exact(0) Distinct=Exact(3) ScanBytes=Exact(32))]] +physical_plan DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/parquet_ndv_write/with_ndv.parquet]]}, projection=[column1], file_type=parquet, statistics=[Rows=Exact(4), Bytes=Exact(32), [(Col[0]: Min=Exact(Int64(1)) Max=Exact(Int64(3)) Null=Exact(0) Distinct=Inexact(3) ScanBytes=Exact(32))]] statement ok DROP TABLE with_ndv; diff --git a/docs/source/user-guide/configs.md b/docs/source/user-guide/configs.md index ba8ea52bb0850..022e7e8a940df 100644 --- a/docs/source/user-guide/configs.md +++ b/docs/source/user-guide/configs.md @@ -112,6 +112,7 @@ The following configuration settings are available: | datafusion.execution.parquet.bloom_filter_on_write | false | (writing) Write bloom filters for all columns when creating parquet files | | datafusion.execution.parquet.bloom_filter_fpp | NULL | (writing) Sets bloom filter false positive probability. If NULL, uses default parquet writer setting | | datafusion.execution.parquet.bloom_filter_ndv | NULL | (writing) Sets bloom filter number of distinct values. If NULL, uses default parquet writer setting | +| datafusion.execution.parquet.write_row_group_number_distinct_values | false | (writing) Write the number of distinct values (NDV) for each column in row group statistics when creating parquet files. Enabling this improves NDV-based query optimizations at the cost of hashing every non-null value during write. | | datafusion.execution.parquet.allow_single_file_parallelism | true | (writing) Controls whether DataFusion will attempt to speed up writing parquet files by serializing them in parallel. Each column in each row group in each output file are serialized in parallel leveraging a maximum possible core count of n_files\*n_row_groups\*n_columns. | | datafusion.execution.parquet.maximum_parallel_row_group_writers | 1 | (writing) By default parallel parquet writer is tuned for minimum memory usage in a streaming execution plan. You may see a performance benefit when writing large parquet files by increasing maximum_parallel_row_group_writers and maximum_buffered_record_batches_per_stream if your system has idle cores and can tolerate additional memory usage. Boosting these values is likely worthwhile when writing out already in-memory data, such as from a cached data frame. | | datafusion.execution.parquet.maximum_buffered_record_batches_per_stream | 2 | (writing) By default parallel parquet writer is tuned for minimum memory usage in a streaming execution plan. You may see a performance benefit when writing large parquet files by increasing maximum_parallel_row_group_writers and maximum_buffered_record_batches_per_stream if your system has idle cores and can tolerate additional memory usage. Boosting these values is likely worthwhile when writing out already in-memory data, such as from a cached data frame. | From 568a166e93fd610ac25dd6d39c7a229d3de94ba5 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 22 Sep 2026 00:43:14 -0400 Subject: [PATCH 04/10] fix: restore Exact precision for single-row-group NDV and round-trip 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 --- datafusion/common/src/file_options/parquet_writer.rs | 4 ++-- datafusion/datasource-parquet/src/metadata.rs | 3 +++ 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/datafusion/common/src/file_options/parquet_writer.rs b/datafusion/common/src/file_options/parquet_writer.rs index 4bdea56586dd8..e31be60cca4f7 100644 --- a/datafusion/common/src/file_options/parquet_writer.rs +++ b/datafusion/common/src/file_options/parquet_writer.rs @@ -557,8 +557,8 @@ mod tests { schema_force_view_types: global_options_defaults.schema_force_view_types, binary_as_string: global_options_defaults.binary_as_string, skip_arrow_metadata: global_options_defaults.skip_arrow_metadata, - write_row_group_number_distinct_values: global_options_defaults - .write_row_group_number_distinct_values, + write_row_group_number_distinct_values: props + .write_row_group_number_distinct_values(), coerce_int96: None, coerce_int96_tz: None, content_defined_chunking: props.content_defined_chunking().into(), diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index abf711c335cbf..a418deccad035 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -966,6 +966,9 @@ fn summarize_distinct_counts( } Ok(match max_distinct_count { + Some(distinct_count) if num_row_groups == 1 => { + Precision::Exact(distinct_count as usize) + } Some(distinct_count) => Precision::Inexact(distinct_count as usize), None => Precision::Absent, }) From aa3c063217092918284f751675b0ace473fe7347 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 22 Sep 2026 00:53:18 -0400 Subject: [PATCH 05/10] fix --- datafusion/datasource-parquet/src/metadata.rs | 19 ++++++++----------- 1 file changed, 8 insertions(+), 11 deletions(-) diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index a418deccad035..0fe3d001e4d1e 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -966,9 +966,6 @@ fn summarize_distinct_counts( } Ok(match max_distinct_count { - Some(distinct_count) if num_row_groups == 1 => { - Precision::Exact(distinct_count as usize) - } Some(distinct_count) => Precision::Inexact(distinct_count as usize), None => Precision::Absent, }) @@ -1533,7 +1530,7 @@ mod tests { #[test] fn test_distinct_count_single_row_group_with_ndv() { - // Single row group with distinct count should return Exact + // Single row group with distinct count should return Inexact let schema_descr = create_schema_descr(1); let arrow_schema = create_arrow_schema(1); @@ -1558,7 +1555,7 @@ mod tests { assert_eq!( result.column_statistics[0].distinct_count, - Precision::Exact(42) + Precision::Inexact(42) ); } @@ -1745,7 +1742,7 @@ mod tests { assert_eq!( result.column_statistics[0].distinct_count, - Precision::Exact(5) + Precision::Inexact(5) ); assert_eq!( result.column_statistics[1].distinct_count, @@ -1753,7 +1750,7 @@ mod tests { ); assert_eq!( result.column_statistics[2].distinct_count, - Precision::Exact(100) + Precision::Inexact(100) ); } @@ -1922,15 +1919,15 @@ mod tests { // category: 10 distinct values assert_eq!( result.column_statistics[1].distinct_count, - Precision::Exact(10), - "category column should have Exact(10) distinct_count" + Precision::Inexact(10), + "category column should have Inexact(10) distinct_count" ); // name: 5 distinct values assert_eq!( result.column_statistics[2].distinct_count, - Precision::Exact(5), - "name column should have Exact(5) distinct_count" + Precision::Inexact(5), + "name column should have Inexact(5) distinct_count" ); } } From 3ed90761e0805fd6bde5e14d4f16375592b2ddfc Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 22 Sep 2026 02:05:06 -0400 Subject: [PATCH 06/10] fix: update clickbench.slt to reflect Inexact NDV physical plan 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. --- datafusion/sqllogictest/test_files/clickbench.slt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/datafusion/sqllogictest/test_files/clickbench.slt b/datafusion/sqllogictest/test_files/clickbench.slt index 7cb5547383c38..07cb828d65a14 100644 --- a/datafusion/sqllogictest/test_files/clickbench.slt +++ b/datafusion/sqllogictest/test_files/clickbench.slt @@ -1203,8 +1203,8 @@ logical_plan 02)--SubqueryAlias: hits 03)----TableScan: hits_raw projection=[HitColor, BrowserLanguage, BrowserCountry] physical_plan -01)ProjectionExec: expr=[1 as count(DISTINCT hits.HitColor), 1 as count(DISTINCT hits.BrowserCountry), 1 as count(DISTINCT hits.BrowserLanguage)] -02)--PlaceholderRowExec +01)AggregateExec: mode=Single, gby=[], aggr=[count(DISTINCT hits.HitColor), count(DISTINCT hits.BrowserCountry), count(DISTINCT hits.BrowserLanguage)] +02)--DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/core/tests/data/clickbench_hits_10.parquet]]}, projection=[HitColor, BrowserLanguage, BrowserCountry], file_type=parquet query III SELECT COUNT(DISTINCT "HitColor"), COUNT(DISTINCT "BrowserCountry"), COUNT(DISTINCT "BrowserLanguage") FROM hits; From e61aeb25d1f0cae9cc0becaafaa3cedd37682fd7 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 5 Oct 2026 00:57:44 -0400 Subject: [PATCH 07/10] address review --- .../test_files/parquet_ndv_write.slt | 35 +++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/datafusion/sqllogictest/test_files/parquet_ndv_write.slt b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt index 4c8fa7e356405..675b44e6e5c01 100644 --- a/datafusion/sqllogictest/test_files/parquet_ndv_write.slt +++ b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt @@ -83,6 +83,41 @@ physical_plan DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/ statement ok DROP TABLE with_ndv; +###### +# Dictionary column NDV is unreliable (arrow-rs hashes dictionary keys, not +# logical values), so stats-based `COUNT(DISTINCT)` folding must not use it. +# Pin this: the written NDV for `d` is wrong (2), but the query still returns +# the correct 6 because the Inexact precision prevents the stats fold. +###### + +statement ok +set datafusion.execution.batch_size = 2; + +query I +COPY (SELECT arrow_cast(column1, 'Dictionary(Int32, Utf8)') AS d, column2 AS x + FROM (VALUES ('a', 1), ('b', 2), ('c', 3), ('d', 4), ('e', 5), ('f', 6))) +TO 'test_files/scratch/parquet_ndv_write/dict.parquet' +STORED AS PARQUET; +---- +6 + +statement ok +CREATE EXTERNAL TABLE dict_ndv +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_ndv_write/dict.parquet'; + +# NDV written for `d` is wrong (2), so it must not be used as Exact +query II +SELECT count(DISTINCT d), count(DISTINCT x) FROM dict_ndv; +---- +6 6 + +statement ok +DROP TABLE dict_ndv; + +statement ok +RESET datafusion.execution.batch_size; + # Reset settings statement ok RESET datafusion.execution.parquet.write_row_group_number_distinct_values; From 7450603f29f0d109b8c91d03db8b272731bc62ae Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 5 Oct 2026 08:19:58 -0400 Subject: [PATCH 08/10] fix: use proto tag 40 for write_row_group_number_distinct_values 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. --- datafusion/proto-common/proto/datafusion_common.proto | 2 +- datafusion/proto-common/src/generated/prost.rs | 2 +- .../proto-models/src/generated/datafusion_proto_common.rs | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/datafusion/proto-common/proto/datafusion_common.proto b/datafusion/proto-common/proto/datafusion_common.proto index e74211d01edc5..19d9f85bf6900 100644 --- a/datafusion/proto-common/proto/datafusion_common.proto +++ b/datafusion/proto-common/proto/datafusion_common.proto @@ -579,7 +579,7 @@ message ParquetOptions { bool schema_force_view_types = 28; // default = false bool binary_as_string = 29; // default = false bool skip_arrow_metadata = 30; // default = false - bool write_row_group_number_distinct_values = 39; // default = false + bool write_row_group_number_distinct_values = 40; // default = false oneof metadata_size_hint_opt { uint64 metadata_size_hint = 4; diff --git a/datafusion/proto-common/src/generated/prost.rs b/datafusion/proto-common/src/generated/prost.rs index 0a47feae8c659..527415293f8b3 100644 --- a/datafusion/proto-common/src/generated/prost.rs +++ b/datafusion/proto-common/src/generated/prost.rs @@ -868,7 +868,7 @@ pub struct ParquetOptions { #[prost(bool, tag = "30")] pub skip_arrow_metadata: bool, /// default = false - #[prost(bool, tag = "39")] + #[prost(bool, tag = "40")] pub write_row_group_number_distinct_values: bool, #[prost(uint64, tag = "12")] pub dictionary_page_size_limit: u64, diff --git a/datafusion/proto-models/src/generated/datafusion_proto_common.rs b/datafusion/proto-models/src/generated/datafusion_proto_common.rs index 0a47feae8c659..527415293f8b3 100644 --- a/datafusion/proto-models/src/generated/datafusion_proto_common.rs +++ b/datafusion/proto-models/src/generated/datafusion_proto_common.rs @@ -868,7 +868,7 @@ pub struct ParquetOptions { #[prost(bool, tag = "30")] pub skip_arrow_metadata: bool, /// default = false - #[prost(bool, tag = "39")] + #[prost(bool, tag = "40")] pub write_row_group_number_distinct_values: bool, #[prost(uint64, tag = "12")] pub dictionary_page_size_limit: u64, From efeb9784dbe7a5e3baf1776d4863cff4f2095569 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Wed, 7 Oct 2026 11:41:42 -0400 Subject: [PATCH 09/10] fix proto issue --- datafusion/proto-common/proto/datafusion_common.proto | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datafusion/proto-common/proto/datafusion_common.proto b/datafusion/proto-common/proto/datafusion_common.proto index 19d9f85bf6900..880b5368ee511 100644 --- a/datafusion/proto-common/proto/datafusion_common.proto +++ b/datafusion/proto-common/proto/datafusion_common.proto @@ -579,7 +579,7 @@ message ParquetOptions { bool schema_force_view_types = 28; // default = false bool binary_as_string = 29; // default = false bool skip_arrow_metadata = 30; // default = false - bool write_row_group_number_distinct_values = 40; // default = false + bool write_row_group_number_distinct_values = 41; // default = false oneof metadata_size_hint_opt { uint64 metadata_size_hint = 4; From 233a1c4c9f9ae5e368d98c799253e80a939cb505 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Thu, 8 Oct 2026 09:41:21 -0400 Subject: [PATCH 10/10] fix: regenerate proto to sync write_row_group_number_distinct_values tag to 41 --- datafusion/proto-common/src/generated/prost.rs | 2 +- .../proto-models/src/generated/datafusion_proto_common.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/datafusion/proto-common/src/generated/prost.rs b/datafusion/proto-common/src/generated/prost.rs index 527415293f8b3..e801f6f71b4f3 100644 --- a/datafusion/proto-common/src/generated/prost.rs +++ b/datafusion/proto-common/src/generated/prost.rs @@ -868,7 +868,7 @@ pub struct ParquetOptions { #[prost(bool, tag = "30")] pub skip_arrow_metadata: bool, /// default = false - #[prost(bool, tag = "40")] + #[prost(bool, tag = "41")] pub write_row_group_number_distinct_values: bool, #[prost(uint64, tag = "12")] pub dictionary_page_size_limit: u64, diff --git a/datafusion/proto-models/src/generated/datafusion_proto_common.rs b/datafusion/proto-models/src/generated/datafusion_proto_common.rs index 527415293f8b3..e801f6f71b4f3 100644 --- a/datafusion/proto-models/src/generated/datafusion_proto_common.rs +++ b/datafusion/proto-models/src/generated/datafusion_proto_common.rs @@ -868,7 +868,7 @@ pub struct ParquetOptions { #[prost(bool, tag = "30")] pub skip_arrow_metadata: bool, /// default = false - #[prost(bool, tag = "40")] + #[prost(bool, tag = "41")] pub write_row_group_number_distinct_values: bool, #[prost(uint64, tag = "12")] pub dictionary_page_size_limit: u64,