From 99c04f2b54841cdfca411953790a105793af0b8b Mon Sep 17 00:00:00 2001 From: Xinyao Zhang <43081360+zhangxinyao88@users.noreply.github.com> Date: Mon, 24 Aug 2026 19:19:00 -0400 Subject: [PATCH] fix: validate parquet statistics config --- datafusion/common/src/config.rs | 30 +++++++- .../common/src/file_options/parquet_writer.rs | 13 ++-- datafusion/common/src/parquet_config.rs | 74 +++++++++++++++++++ .../datasource-parquet/src/file_format.rs | 2 +- datafusion/proto-common/src/from_proto/mod.rs | 40 +++++++++- datafusion/proto-common/src/to_proto/mod.rs | 2 +- datafusion/proto-models/src/from_proto.rs | 12 +-- .../sqllogictest/test_files/set_variable.slt | 19 +++++ 8 files changed, 173 insertions(+), 19 deletions(-) diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index b340ec86283d8..159c4b3b31e1c 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -23,7 +23,7 @@ use arrow_ipc::CompressionType; use crate::encryption::{FileDecryptionProperties, FileEncryptionProperties}; use crate::error::{_config_datafusion_err, _config_err}; use crate::format::{ExplainAnalyzeCategories, ExplainFormat, MetricType}; -use crate::parquet_config::DFParquetWriterVersion; +use crate::parquet_config::{DFParquetStatistics, DFParquetWriterVersion}; use crate::parsers::{CompressionTypeVariant, CsvQuoteStyle}; use crate::utils::get_available_parallelism; use crate::{DataFusionError, Result}; @@ -1419,7 +1419,7 @@ config_namespace! { /// Valid values are: "none", "chunk", and "page" /// These values are not case sensitive. If NULL, uses /// default parquet writer setting - pub statistics_enabled: Option, transform = str::to_lowercase, default = Some("page".into()) + pub statistics_enabled: Option, default = Some(DFParquetStatistics::Page) /// (writing) Target maximum number of rows in each row group (defaults to 1M /// rows). Writing larger row groups requires more memory to write, but @@ -4577,6 +4577,32 @@ mod tests { ); } + #[test] + fn test_parquet_statistics_validation() { + use crate::{config::ConfigOptions, parquet_config::DFParquetStatistics}; + + let mut config = ConfigOptions::default(); + + for (value, expected) in [ + ("none", DFParquetStatistics::None), + ("CHUNK", DFParquetStatistics::Chunk), + ("page", DFParquetStatistics::Page), + ] { + config + .set("datafusion.execution.parquet.statistics_enabled", value) + .unwrap(); + assert_eq!(config.execution.parquet.statistics_enabled, Some(expected)); + } + + let err = config + .set("datafusion.execution.parquet.statistics_enabled", "invalid") + .unwrap_err(); + assert_contains!( + err.to_string(), + "Invalid parquet statistics setting: invalid. Expected one of: none, chunk, page" + ); + } + #[cfg(feature = "parquet")] #[test] fn set_cdc_enabled_flag() { diff --git a/datafusion/common/src/file_options/parquet_writer.rs b/datafusion/common/src/file_options/parquet_writer.rs index c539245764d45..6064d7bc9d334 100644 --- a/datafusion/common/src/file_options/parquet_writer.rs +++ b/datafusion/common/src/file_options/parquet_writer.rs @@ -257,9 +257,8 @@ impl ParquetOptions { .set_writer_version((*writer_version).into()) .set_dictionary_page_size_limit(*dictionary_page_size_limit) .set_statistics_enabled( - statistics_enabled - .as_ref() - .and_then(|s| parse_statistics_string(s).ok()) + (*statistics_enabled) + .map(Into::into) .unwrap_or(DEFAULT_STATISTICS_ENABLED), ) .set_max_row_group_row_count(Some(*max_row_group_size)) @@ -434,7 +433,7 @@ mod tests { MaxRowGroupBytes, ParquetCdcOptions, ParquetColumnOptions, ParquetEncryptionOptions, ParquetOptions, }; - use crate::parquet_config::DFParquetWriterVersion; + use crate::parquet_config::{DFParquetStatistics, DFParquetWriterVersion}; use parquet::basic::Compression; use parquet::file::properties::{ BloomFilterProperties, DEFAULT_BLOOM_FILTER_FPP, DEFAULT_BLOOM_FILTER_NDV, @@ -475,7 +474,7 @@ mod tests { compression: Some("zstd(22)".into()), dictionary_enabled: Some(!defaults.dictionary_enabled.unwrap_or(false)), dictionary_page_size_limit: 43, - statistics_enabled: Some("chunk".into()), + statistics_enabled: Some(DFParquetStatistics::Chunk), max_row_group_size: 42, max_row_group_bytes: Some(MaxRowGroupBytes::try_new(42).unwrap()), created_by: "wordy".into(), @@ -545,7 +544,7 @@ mod tests { /// (use identity to confirm correct.) fn session_config_from_writer_props(props: &WriterProperties) -> TableParquetOptions { let default_col = ColumnPath::from("col doesn't have specific config"); - let default_col_props = extract_column_options(props, default_col); + let default_col_props = extract_column_options(props, default_col.clone()); let configured_col = ColumnPath::from(COL_NAME); let configured_col_props = extract_column_options(props, configured_col); @@ -600,7 +599,7 @@ mod tests { encoding: default_col_props.encoding, compression: default_col_props.compression, dictionary_enabled: default_col_props.dictionary_enabled, - statistics_enabled: default_col_props.statistics_enabled, + statistics_enabled: Some(props.statistics_enabled(&default_col).into()), bloom_filter_on_write: default_col_props .bloom_filter_enabled .unwrap_or_default(), diff --git a/datafusion/common/src/parquet_config.rs b/datafusion/common/src/parquet_config.rs index 9d6d7a88566a7..06f1e0af217ff 100644 --- a/datafusion/common/src/parquet_config.rs +++ b/datafusion/common/src/parquet_config.rs @@ -106,3 +106,77 @@ impl From for DFParquetWriterVersion { } } } + +/// Parquet statistics levels supported by the writer +/// +/// This enum validates statistics settings at configuration time, ensuring only +/// `none`, `chunk`, or `page` can be set via `SET` commands or deserialization. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum DFParquetStatistics { + /// Do not write statistics + None, + /// Write chunk-level statistics + Chunk, + /// Write page-level statistics + #[default] + Page, +} + +impl FromStr for DFParquetStatistics { + type Err = DataFusionError; + + fn from_str(s: &str) -> Result { + match s.to_lowercase().as_str() { + "none" => Ok(Self::None), + "chunk" => Ok(Self::Chunk), + "page" => Ok(Self::Page), + other => Err(DataFusionError::Configuration(format!( + "Invalid parquet statistics setting: {other}. Expected one of: none, chunk, page" + ))), + } + } +} + +impl Display for DFParquetStatistics { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let s = match self { + Self::None => "none", + Self::Chunk => "chunk", + Self::Page => "page", + }; + f.write_str(s) + } +} + +impl ConfigField for DFParquetStatistics { + fn visit(&self, v: &mut V, key: &str, description: &'static str) { + v.some(key, self, description) + } + + fn set(&mut self, _: &str, value: &str) -> Result<()> { + *self = Self::from_str(value)?; + Ok(()) + } +} + +#[cfg(feature = "parquet")] +impl From for parquet::file::properties::EnabledStatistics { + fn from(value: DFParquetStatistics) -> Self { + match value { + DFParquetStatistics::None => Self::None, + DFParquetStatistics::Chunk => Self::Chunk, + DFParquetStatistics::Page => Self::Page, + } + } +} + +#[cfg(feature = "parquet")] +impl From for DFParquetStatistics { + fn from(value: parquet::file::properties::EnabledStatistics) -> Self { + match value { + parquet::file::properties::EnabledStatistics::None => Self::None, + parquet::file::properties::EnabledStatistics::Chunk => Self::Chunk, + parquet::file::properties::EnabledStatistics::Page => Self::Page, + } + } +} diff --git a/datafusion/datasource-parquet/src/file_format.rs b/datafusion/datasource-parquet/src/file_format.rs index 1b25a2c632510..1c6c875f016a4 100644 --- a/datafusion/datasource-parquet/src/file_format.rs +++ b/datafusion/datasource-parquet/src/file_format.rs @@ -743,7 +743,7 @@ impl From<&ParquetFormatFactory> for protobuf::TableParquetOptions { }), dictionary_page_size_limit: global_options.global.dictionary_page_size_limit as u64, statistics_enabled_opt: global_options.global.statistics_enabled.map(|enabled| { - parquet_options::StatisticsEnabledOpt::StatisticsEnabled(enabled) + parquet_options::StatisticsEnabledOpt::StatisticsEnabled(enabled.to_string()) }), max_row_group_size: global_options.global.max_row_group_size as u64, max_in_list_size: global_options.global.max_in_list_size as u64, diff --git a/datafusion/proto-common/src/from_proto/mod.rs b/datafusion/proto-common/src/from_proto/mod.rs index 169ff7f3d9ff2..52570a6f3f759 100644 --- a/datafusion/proto-common/src/from_proto/mod.rs +++ b/datafusion/proto-common/src/from_proto/mod.rs @@ -1075,11 +1075,11 @@ impl TryFrom<&protobuf::ParquetOptions> for ParquetOptions { // Continuing from where we left off in the TryFrom implementation dictionary_page_size_limit: value.dictionary_page_size_limit as usize, statistics_enabled: value - .statistics_enabled_opt.clone() + .statistics_enabled_opt.as_ref() .map(|opt| match opt { - protobuf::parquet_options::StatisticsEnabledOpt::StatisticsEnabled(v) => Some(v), + protobuf::parquet_options::StatisticsEnabledOpt::StatisticsEnabled(v) => v.parse(), }) - .unwrap_or(None), + .transpose()?, max_row_group_size: value.max_row_group_size as usize, max_in_list_size: value.max_in_list_size as usize, created_by: value.created_by.clone(), @@ -1335,6 +1335,7 @@ mod tests { use datafusion_common::config::{ MaxRowGroupBytes, ParquetCdcOptions, ParquetOptions, TableParquetOptions, }; + use datafusion_common::parquet_config::DFParquetStatistics; fn parquet_options_proto_round_trip(opts: ParquetOptions) -> ParquetOptions { let proto: crate::protobuf_common::ParquetOptions = @@ -1378,6 +1379,39 @@ mod tests { assert_eq!(recovered.coerce_int96_tz, Some("UTC".to_string())); } + #[test] + fn test_parquet_statistics_round_trip() { + let opts = ParquetOptions { + statistics_enabled: Some(DFParquetStatistics::Chunk), + ..ParquetOptions::default() + }; + let recovered = parquet_options_proto_round_trip(opts); + assert_eq!( + recovered.statistics_enabled, + Some(DFParquetStatistics::Chunk) + ); + } + + #[test] + fn test_invalid_parquet_statistics_rejected_from_proto() { + let opts = ParquetOptions::default(); + let mut proto: crate::protobuf_common::ParquetOptions = + (&opts).try_into().expect("to_proto"); + proto.statistics_enabled_opt = Some( + crate::protobuf_common::parquet_options::StatisticsEnabledOpt::StatisticsEnabled( + "invalid".to_string(), + ), + ); + + let err = ParquetOptions::try_from(&proto).unwrap_err(); + assert!( + err.to_string().contains( + "Invalid parquet statistics setting: invalid. Expected one of: none, chunk, page" + ), + "unexpected error: {err}" + ); + } + #[test] fn test_parquet_options_max_row_group_bytes_round_trip() { let opts = ParquetOptions { diff --git a/datafusion/proto-common/src/to_proto/mod.rs b/datafusion/proto-common/src/to_proto/mod.rs index 360981746585b..9ce283b8d20c4 100644 --- a/datafusion/proto-common/src/to_proto/mod.rs +++ b/datafusion/proto-common/src/to_proto/mod.rs @@ -910,7 +910,7 @@ impl TryFrom<&ParquetOptions> for protobuf::ParquetOptions { compression_opt: value.compression.clone().map(protobuf::parquet_options::CompressionOpt::Compression), dictionary_enabled_opt: value.dictionary_enabled.map(protobuf::parquet_options::DictionaryEnabledOpt::DictionaryEnabled), dictionary_page_size_limit: value.dictionary_page_size_limit as u64, - statistics_enabled_opt: value.statistics_enabled.clone().map(protobuf::parquet_options::StatisticsEnabledOpt::StatisticsEnabled), + statistics_enabled_opt: value.statistics_enabled.map(|v| protobuf::parquet_options::StatisticsEnabledOpt::StatisticsEnabled(v.to_string())), max_row_group_size: value.max_row_group_size as u64, max_in_list_size: value.max_in_list_size as u64, created_by: value.created_by.clone(), diff --git a/datafusion/proto-models/src/from_proto.rs b/datafusion/proto-models/src/from_proto.rs index 74ead8c52049b..30ff73d7d4831 100644 --- a/datafusion/proto-models/src/from_proto.rs +++ b/datafusion/proto-models/src/from_proto.rs @@ -366,13 +366,15 @@ impl TryFrom<&ParquetOptionsProto> for ParquetOptions { } }), dictionary_page_size_limit: proto.dictionary_page_size_limit as usize, - statistics_enabled: proto.statistics_enabled_opt.as_ref().map( - |opt| match opt { + statistics_enabled: proto + .statistics_enabled_opt + .as_ref() + .map(|opt| match opt { parquet_options::StatisticsEnabledOpt::StatisticsEnabled( statistics, - ) => statistics.clone(), - }, - ), + ) => statistics.parse(), + }) + .transpose()?, max_row_group_size: proto.max_row_group_size as usize, max_in_list_size: proto.max_in_list_size as usize, created_by: proto.created_by.clone(), diff --git a/datafusion/sqllogictest/test_files/set_variable.slt b/datafusion/sqllogictest/test_files/set_variable.slt index e36e59bccb66b..b1dfc7163c035 100644 --- a/datafusion/sqllogictest/test_files/set_variable.slt +++ b/datafusion/sqllogictest/test_files/set_variable.slt @@ -732,6 +732,25 @@ statement error DataFusion error: Error during planning: Duration has overflowed SET datafusion.runtime.list_files_cache_ttl = '1m18446744073709551556s' # Set invalid value and ensures error +statement ok +SET datafusion.execution.parquet.statistics_enabled = 'CHUNK' + +query TT +SHOW datafusion.execution.parquet.statistics_enabled +---- +datafusion.execution.parquet.statistics_enabled chunk + +statement error +SET datafusion.execution.parquet.statistics_enabled = 'invalid' +---- +DataFusion error: Error setting config datafusion.execution.parquet.statistics_enabled +caused by +Invalid or Unsupported Configuration: Invalid parquet statistics setting: invalid. Expected one of: none, chunk, page + + +statement ok +RESET datafusion.execution.parquet.statistics_enabled + statement error DataFusion error: Error setting config datafusion\.execution\.batch_size\ncaused by\nInvalid or Unsupported Configuration: value must be greater than 0 SET datafusion.execution.batch_size = 0