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..e31be60cca4f7 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: 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/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/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index 51dc98e4f2f02..0fe3d001e4d1e 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,14 @@ 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 { - Some(distinct_count) if num_row_groups == 1 => { - Precision::Exact(distinct_count as usize) - } + Ok(match max_distinct_count { Some(distinct_count) => Precision::Inexact(distinct_count as usize), None => Precision::Absent, - } + }) } /// Compute the Arrow in-memory size for a single column @@ -1535,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); @@ -1560,7 +1555,7 @@ mod tests { assert_eq!( result.column_statistics[0].distinct_count, - Precision::Exact(42) + Precision::Inexact(42) ); } @@ -1747,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, @@ -1755,7 +1750,7 @@ mod tests { ); assert_eq!( result.column_statistics[2].distinct_count, - Precision::Exact(100) + Precision::Inexact(100) ); } @@ -1924,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" ); } } diff --git a/datafusion/proto-common/proto/datafusion_common.proto b/datafusion/proto-common/proto/datafusion_common.proto index bdabb0777ef81..880b5368ee511 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 = 41; // 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..e801f6f71b4f3 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 = "41")] + 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..e801f6f71b4f3 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 = "41")] + 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/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; 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 new file mode 100644 index 0000000000000..675b44e6e5c01 --- /dev/null +++ b/datafusion/sqllogictest/test_files/parquet_ndv_write.slt @@ -0,0 +1,132 @@ +# 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 Inexact (NDV from parquet statistics is never treated as exact). +###### + +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 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=Inexact(3) ScanBytes=Exact(32))]] + +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; + +statement ok +RESET datafusion.explain.show_statistics; + +statement ok +RESET datafusion.explain.physical_plan_only; + +statement ok +RESET datafusion.execution.collect_statistics; 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. |