From ec51f9bbda54c364f64fda075bb0154a584363e5 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 1/3] fix: validate parquet statistics config --- datafusion/common/src/config.rs | 43 ++++++++- .../common/src/file_options/parquet_writer.rs | 13 ++- datafusion/common/src/parquet_config.rs | 94 +++++++++++++++++++ .../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 | 35 ++++++- .../sqllogictest/test_files/set_variable.slt | 19 ++++ 8 files changed, 229 insertions(+), 19 deletions(-) diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 7c7604e041368..cdbb6dc91226b 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}; @@ -1421,7 +1421,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 @@ -4578,6 +4578,45 @@ mod tests { ); } + #[cfg(feature = "parquet")] + #[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" + ); + + // An unset value can arise from deserialization. An invalid update must + // leave that state unchanged rather than inserting the default. + config.execution.parquet.statistics_enabled = None; + assert_eq!(config.execution.parquet.statistics_enabled, None); + + assert!( + config + .set("datafusion.execution.parquet.statistics_enabled", "invalid") + .is_err() + ); + assert_eq!(config.execution.parquet.statistics_enabled, None); + } + #[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..f6184682cab17 100644 --- a/datafusion/common/src/parquet_config.rs +++ b/datafusion/common/src/parquet_config.rs @@ -106,3 +106,97 @@ 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)] +pub enum DFParquetStatistics { + /// Do not write statistics + None, + /// Write chunk-level statistics + Chunk, + /// Write page-level statistics + 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(()) + } +} + +/// `ConfigField` for `Option` parses before assigning so +/// an invalid value does not turn an unset option into the default. +impl ConfigField for Option { + fn visit(&self, v: &mut V, key: &str, description: &'static str) { + match self { + Some(statistics) => statistics.visit(v, key, description), + None => v.none(key, description), + } + } + + fn set(&mut self, _key: &str, value: &str) -> Result<()> { + *self = Some(DFParquetStatistics::from_str(value)?); + Ok(()) + } + + fn reset(&mut self, _key: &str) -> Result<()> { + *self = None; + 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 cf3a3cb75a0f0..689d983c16c8a 100644 --- a/datafusion/proto-common/src/from_proto/mod.rs +++ b/datafusion/proto-common/src/from_proto/mod.rs @@ -1077,11 +1077,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(), @@ -1337,6 +1337,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 = @@ -1380,6 +1381,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..ca532eb303422 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(), @@ -517,3 +519,26 @@ impl TryFrom<&TableParquetOptionsProto> for TableParquetOptions { }) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn rejects_invalid_parquet_statistics() { + let proto = ParquetOptionsProto { + statistics_enabled_opt: Some( + parquet_options::StatisticsEnabledOpt::StatisticsEnabled( + "invalid".to_string(), + ), + ), + ..Default::default() + }; + + let err = ParquetOptions::try_from(&proto).unwrap_err(); + assert!( + err.to_string() + .contains("Invalid parquet statistics setting: invalid") + ); + } +} 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 From 661f7fc8edb667d79049edad3b7d991c8b1c8cb4 Mon Sep 17 00:00:00 2001 From: Xinyao Zhang <43081360+zhangxinyao88@users.noreply.github.com> Date: Tue, 1 Sep 2026 08:56:01 -0400 Subject: [PATCH 2/3] fix: reject nested parquet statistics reset --- datafusion/common/src/config.rs | 11 +++++++++++ datafusion/common/src/parquet_config.rs | 13 ++++++++++--- 2 files changed, 21 insertions(+), 3 deletions(-) diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index cdbb6dc91226b..36c5acf257068 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -4615,6 +4615,17 @@ mod tests { .is_err() ); assert_eq!(config.execution.parquet.statistics_enabled, None); + + config.execution.parquet.statistics_enabled = Some(DFParquetStatistics::Page); + assert!( + config + .reset("datafusion.execution.parquet.statistics_enabled.typo") + .is_err() + ); + assert_eq!( + config.execution.parquet.statistics_enabled, + Some(DFParquetStatistics::Page) + ); } #[cfg(feature = "parquet")] diff --git a/datafusion/common/src/parquet_config.rs b/datafusion/common/src/parquet_config.rs index f6184682cab17..20922b790aedf 100644 --- a/datafusion/common/src/parquet_config.rs +++ b/datafusion/common/src/parquet_config.rs @@ -173,9 +173,16 @@ impl ConfigField for Option { Ok(()) } - fn reset(&mut self, _key: &str) -> Result<()> { - *self = None; - Ok(()) + fn reset(&mut self, key: &str) -> Result<()> { + if key.is_empty() { + *self = None; + Ok(()) + } else { + crate::error::_config_err!( + "Config field parquet.statistics_enabled is a scalar Option and does not have nested field \"{}\"", + key + ) + } } } From cdc9debe2ccaddcae4a929ae77fafa4cfed17556 Mon Sep 17 00:00:00 2001 From: Xinyao Zhang <43081360+zhangxinyao88@users.noreply.github.com> Date: Tue, 1 Sep 2026 23:40:07 -0400 Subject: [PATCH 3/3] fix: reject nested parquet statistics set --- datafusion/common/src/config.rs | 17 +++++++++++++++++ datafusion/common/src/parquet_config.rs | 18 ++++++++++++++++-- 2 files changed, 33 insertions(+), 2 deletions(-) diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 36c5acf257068..999d5c7406d9f 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -4617,6 +4617,19 @@ mod tests { assert_eq!(config.execution.parquet.statistics_enabled, None); config.execution.parquet.statistics_enabled = Some(DFParquetStatistics::Page); + assert!( + config + .set( + "datafusion.execution.parquet.statistics_enabled.typo", + "none" + ) + .is_err() + ); + assert_eq!( + config.execution.parquet.statistics_enabled, + Some(DFParquetStatistics::Page) + ); + assert!( config .reset("datafusion.execution.parquet.statistics_enabled.typo") @@ -4626,6 +4639,10 @@ mod tests { config.execution.parquet.statistics_enabled, Some(DFParquetStatistics::Page) ); + + let mut scalar = DFParquetStatistics::Page; + assert!(ConfigField::set(&mut scalar, "typo", "none").is_err()); + assert_eq!(scalar, DFParquetStatistics::Page); } #[cfg(feature = "parquet")] diff --git a/datafusion/common/src/parquet_config.rs b/datafusion/common/src/parquet_config.rs index 20922b790aedf..b727664e77886 100644 --- a/datafusion/common/src/parquet_config.rs +++ b/datafusion/common/src/parquet_config.rs @@ -152,7 +152,14 @@ impl ConfigField for DFParquetStatistics { v.some(key, self, description) } - fn set(&mut self, _: &str, value: &str) -> Result<()> { + fn set(&mut self, key: &str, value: &str) -> Result<()> { + if !key.is_empty() { + return crate::error::_config_err!( + "Config field parquet.statistics_enabled is a scalar DFParquetStatistics and does not have nested field \"{}\"", + key + ); + } + *self = Self::from_str(value)?; Ok(()) } @@ -168,7 +175,14 @@ impl ConfigField for Option { } } - fn set(&mut self, _key: &str, value: &str) -> Result<()> { + fn set(&mut self, key: &str, value: &str) -> Result<()> { + if !key.is_empty() { + return crate::error::_config_err!( + "Config field parquet.statistics_enabled is a scalar Option and does not have nested field \"{}\"", + key + ); + } + *self = Some(DFParquetStatistics::from_str(value)?); Ok(()) }