Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions datafusion/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1514,6 +1514,12 @@ config_namespace! {
/// default parquet writer setting
pub bloom_filter_ndv: Option<u64>, default = None

/// (writing) Write the number of distinct values (NDV) for each column
Comment thread
Rich-T-kid marked this conversation as resolved.
/// 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
Expand Down
10 changes: 9 additions & 1 deletion datafusion/common/src/file_options/parquet_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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(),
Expand Down
3 changes: 3 additions & 0 deletions datafusion/datasource-parquet/src/file_format.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}),
Expand Down
51 changes: 23 additions & 28 deletions datafusion/datasource-parquet/src/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -930,29 +930,27 @@ where

/// Extract distinct counts from row group column statistics.
fn summarize_distinct_counts(
parquet_idx: Option<usize>,
stats_converter: &StatisticsConverter,
row_groups_metadata: &[RowGroupMetaData],
) -> Precision<usize> {
let Some(parquet_idx) = parquet_idx else {
return Precision::Absent;
};
) -> Result<Precision<usize>> {
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<u64> = 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),
Expand All @@ -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 {
Comment thread
Rich-T-kid marked this conversation as resolved.
Some(distinct_count) => Precision::Inexact(distinct_count as usize),
None => Precision::Absent,
}
})
}

/// Compute the Arrow in-memory size for a single column
Expand Down Expand Up @@ -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);

Expand All @@ -1560,7 +1555,7 @@ mod tests {

assert_eq!(
result.column_statistics[0].distinct_count,
Precision::Exact(42)
Precision::Inexact(42)
);
}

Expand Down Expand Up @@ -1747,15 +1742,15 @@ mod tests {

assert_eq!(
result.column_statistics[0].distinct_count,
Precision::Exact(5)
Precision::Inexact(5)
);
assert_eq!(
result.column_statistics[1].distinct_count,
Precision::Absent
);
assert_eq!(
result.column_statistics[2].distinct_count,
Precision::Exact(100)
Precision::Inexact(100)
);
}

Expand Down Expand Up @@ -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"
);
}
}
Expand Down
1 change: 1 addition & 0 deletions datafusion/proto-common/proto/datafusion_common.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
1 change: 1 addition & 0 deletions datafusion/proto-common/src/from_proto/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
24 changes: 21 additions & 3 deletions datafusion/proto-common/src/generated/pbjson.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
;
}
Expand All @@ -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)
;
}
Expand Down Expand Up @@ -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::<std::result::Result<Vec<_>, _>>()?;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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)]
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -6825,6 +6833,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions {
SchemaForceViewTypes,
BinaryAsString,
SkipArrowMetadata,
WriteRowGroupNumberDistinctValues,
DictionaryPageSizeLimit,
DataPageRowCountLimit,
MaxRowGroupSize,
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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"));
Expand Down Expand Up @@ -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(),
Expand Down
3 changes: 3 additions & 0 deletions datafusion/proto-common/src/generated/prost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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")]
Expand Down
1 change: 1 addition & 0 deletions datafusion/proto-common/src/to_proto/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)),
Expand Down
2 changes: 2 additions & 0 deletions datafusion/proto-models/src/from_proto.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")]
Expand Down
4 changes: 2 additions & 2 deletions datafusion/sqllogictest/test_files/clickbench.slt
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading