From ed3d6201b69f1aa58b381f1fc86466516ed7ed53 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Wed, 26 Aug 2026 00:42:38 -0400 Subject: [PATCH 1/9] enable rle_to_dictionary flag --- datafusion/common/src/config.rs | 11 + .../common/src/file_options/parquet_writer.rs | 4 + .../datasource-parquet/src/file_format.rs | 48 +- datafusion/datasource-parquet/src/metadata.rs | 101 ++- .../datasource-parquet/src/opener/mod.rs | 72 ++- .../datasource-parquet/src/schema_coercion.rs | 574 +++++++++++++++++- datafusion/datasource-parquet/src/source.rs | 4 + .../proto/datafusion_common.proto | 2 + datafusion/proto-common/src/from_proto/mod.rs | 1 + .../proto-common/src/generated/pbjson.rs | 20 +- .../proto-common/src/generated/prost.rs | 2 + datafusion/proto-common/src/to_proto/mod.rs | 1 + datafusion/proto-models/src/from_proto.rs | 1 + .../src/generated/datafusion_proto_common.rs | 2 + .../proto/src/logical_plan/file_formats.rs | 27 + .../test_files/information_schema.slt | 2 + .../test_files/parquet_rle_to_dictionary.slt | 250 ++++++++ docs/source/user-guide/configs.md | 1 + 18 files changed, 1111 insertions(+), 12 deletions(-) create mode 100644 datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 4246dcc451d45..179eb550fc231 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -1423,6 +1423,17 @@ config_namespace! { /// Defaults to 20. pub max_in_list_size: usize, default = 20 + /// (reading) If true, top-level string and binary Parquet columns with + /// dictionary pages are inferred and scanned as + /// `Dictionary` / `Dictionary` instead of + /// their plain value type. + /// + /// This applies only when DataFusion infers the table schema. Tables with + /// a user-supplied schema are not promoted because the Parquet footer is + /// not read at DDL time, so dictionary pages cannot be detected per column. + /// See + pub enable_rle_to_dictionary: bool, default = false + // The following options affect writing to parquet files // and map to parquet::file::properties::WriterProperties diff --git a/datafusion/common/src/file_options/parquet_writer.rs b/datafusion/common/src/file_options/parquet_writer.rs index af4554c4eef84..2de10509396e0 100644 --- a/datafusion/common/src/file_options/parquet_writer.rs +++ b/datafusion/common/src/file_options/parquet_writer.rs @@ -249,6 +249,7 @@ impl ParquetOptions { skip_arrow_metadata: _, max_predicate_cache_size: _, max_in_list_size: _, + enable_rle_to_dictionary: _, } = self; let mut builder = WriterProperties::builder() @@ -427,6 +428,7 @@ mod tests { coerce_int96_tz: None, max_predicate_cache_size: defaults.max_predicate_cache_size, content_defined_chunking: defaults.content_defined_chunking.clone(), + enable_rle_to_dictionary: defaults.enable_rle_to_dictionary, } } @@ -554,6 +556,8 @@ mod tests { coerce_int96: None, coerce_int96_tz: None, content_defined_chunking: props.content_defined_chunking().into(), + enable_rle_to_dictionary: global_options_defaults + .enable_rle_to_dictionary, }, column_specific_options, key_value_metadata, diff --git a/datafusion/datasource-parquet/src/file_format.rs b/datafusion/datasource-parquet/src/file_format.rs index dc0f9a7be434e..8bb8a475a5f89 100644 --- a/datafusion/datasource-parquet/src/file_format.rs +++ b/datafusion/datasource-parquet/src/file_format.rs @@ -145,6 +145,9 @@ impl Debug for ParquetFormatFactory { #[derive(Debug, Default)] pub struct ParquetFormat { options: TableParquetOptions, + /// When set, restricts RLE->Dictionary promotion to only these columns. + /// Overrides `enable_rle_to_dictionary`; see [`Self::with_rle_column_allowlist`]. + rle_column_allowlist: Option>, } impl ParquetFormat { @@ -209,6 +212,22 @@ impl ParquetFormat { &self.options } + /// Restrict RLE->Dictionary promotion to a named subset of columns. + /// + /// Only columns in `columns` that also have dictionary pages in the file + /// are promoted to `Dictionary(Int32, …)` in the inferred schema. + /// Overrides `enable_rle_to_dictionary` when set; an empty set disables + /// promotion entirely. See . + pub fn with_rle_column_allowlist(mut self, columns: HashSet) -> Self { + self.rle_column_allowlist = Some(columns); + self + } + + /// Returns the RLE->Dictionary column allowlist, if set. + pub fn rle_column_allowlist(&self) -> Option<&HashSet> { + self.rle_column_allowlist.as_ref() + } + /// Get [`schema_force_view_types`] /// /// [`schema_force_view_types`]: datafusion_common::config::ParquetOptions::schema_force_view_types @@ -356,14 +375,19 @@ impl FileFormat for ParquetFormat { &object.location, ) .await?; - let result = DFParquetMetadata::new(store.as_ref(), object) + let mut meta = DFParquetMetadata::new(store.as_ref(), object) .with_metadata_size_hint(self.metadata_size_hint()) .with_decryption_properties(file_decryption_properties) .with_file_metadata_cache(Some(Arc::clone(&file_metadata_cache))) .with_coerce_int96(coerce_int96) .with_coerce_int96_tz(coerce_int96_tz.clone()) - .fetch_schema_with_location() - .await?; + .with_enable_rle_to_dictionary( + self.options.global.enable_rle_to_dictionary, + ); + if let Some(allowlist) = &self.rle_column_allowlist { + meta = meta.with_rle_column_allowlist(allowlist.clone()); + } + let result = meta.fetch_schema_with_location().await?; Ok::<_, DataFusionError>(result) }) .boxed() // Workaround https://github.com/rust-lang/rust/issues/64552 @@ -398,7 +422,12 @@ impl FileFormat for ParquetFormat { } drop(seen); - let schemas = schemas.into_iter().map(|(_, schema)| schema); + // Normalize dict-promoted schemas before merging so mixed dict/plain files merge cleanly. + let mut schemas: Vec = + schemas.into_iter().map(|(_, schema)| schema).collect(); + if self.options.global.enable_rle_to_dictionary { + schemas = crate::schema_coercion::uniform_dict_schemas(schemas); + } let schema = if self.skip_metadata() { Schema::try_merge(clear_metadata(schemas)) @@ -508,7 +537,13 @@ impl FileFormat for ParquetFormat { .downcast_ref::() .cloned() .ok_or_else(|| internal_datafusion_err!("Expected ParquetSource"))?; - source = source.with_table_parquet_options(self.options.clone()); + let mut source_options = self.options.clone(); + // An allowlist implies dict-typed fields in the table schema, so the reader + // must also have the flag set to coerce those columns at scan time. + if self.rle_column_allowlist.is_some() { + source_options.global.enable_rle_to_dictionary = true; + } + source = source.with_table_parquet_options(source_options); // Use the CachedParquetFileReaderFactory let metadata_cache = state.runtime_env().cache_manager.get_file_metadata_cache(); @@ -520,7 +555,7 @@ impl FileFormat for ParquetFormat { source = source.with_parquet_file_reader_factory(cached_parquet_read_factory); if let Some(metadata_size_hint) = metadata_size_hint { - source = source.with_metadata_size_hint(metadata_size_hint) + source = source.with_metadata_size_hint(metadata_size_hint); } source = self.set_source_encryption_factory(source, state)?; @@ -672,6 +707,7 @@ impl From<&ParquetFormatFactory> for protobuf::TableParquetOptions { compression_opt: global_options.global.compression.map(|compression| { parquet_options::CompressionOpt::Compression(compression.to_string()) }), + enable_rle_to_dictionary: global_options.global.enable_rle_to_dictionary, dictionary_enabled_opt: global_options.global.dictionary_enabled.map(|enabled| { parquet_options::DictionaryEnabledOpt::DictionaryEnabled(enabled) }), diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index 51dc98e4f2f02..18f2eaefb6b01 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -23,11 +23,11 @@ use crate::{Int96Coercer, apply_file_schema_type_coercions}; use arrow::array::{Array, ArrayRef, BooleanArray}; use arrow::compute::kernels::cmp::eq; use arrow::compute::{and, sum}; -use arrow::datatypes::{DataType, Schema, SchemaRef, TimeUnit}; +use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit}; use datafusion_common::encryption::FileDecryptionProperties; use datafusion_common::stats::Precision; use datafusion_common::{ - ColumnStatistics, DataFusionError, HashMap, Result, ScalarValue, Statistics, + ColumnStatistics, DataFusionError, HashMap, HashSet, Result, ScalarValue, Statistics, internal_datafusion_err, }; use datafusion_execution::cache::cache_manager::{ @@ -152,6 +152,11 @@ pub struct DFParquetMetadata<'a> { pub coerce_int96: Option, /// Optional timezone applied to INT96-coerced timestamps. pub coerce_int96_tz: Option>, + /// If true, promote string/binary columns with dictionary pages to `Dictionary(Int32, ...)`. + enable_rle_to_dictionary: bool, + /// When set, restricts RLE->Dictionary promotion to only these columns. + /// Takes priority over `enable_rle_to_dictionary`; `None` defers to the flag. + rle_column_allowlist: Option>, } impl<'a> DFParquetMetadata<'a> { @@ -169,9 +174,27 @@ impl<'a> DFParquetMetadata<'a> { page_index_policy: None, coerce_int96: None, coerce_int96_tz: None, + enable_rle_to_dictionary: false, + rle_column_allowlist: None, } } + /// Promote string/binary columns with dictionary pages to `Dictionary(Int32, ...)`. + pub fn with_enable_rle_to_dictionary(mut self, enable: bool) -> Self { + self.enable_rle_to_dictionary = enable; + self + } + + /// Restrict RLE->Dictionary promotion to a named subset of columns. + /// Overrides `enable_rle_to_dictionary`; see [`crate::file_format::ParquetFormat::with_rle_column_allowlist`]. + pub(crate) fn with_rle_column_allowlist( + mut self, + columns: impl IntoIterator, + ) -> Self { + self.rle_column_allowlist = Some(columns.into_iter().collect()); + self + } + /// Set a hint for the number of trailing bytes to prefetch from the end /// of the file, equivalent to /// [`ParquetMetaDataReader::with_prefetch_hint`]. @@ -449,6 +472,80 @@ impl<'a> DFParquetMetadata<'a> { .coerce() }) .unwrap_or(schema); + + let should_promote = + self.enable_rle_to_dictionary || self.rle_column_allowlist.is_some(); + let schema = if should_promote { + let schema_descr = file_metadata.schema_descr(); + // Top-level columns that have a dictionary page in at least one row group. + let dict_cols: HashSet = metadata + .row_groups() + .iter() + .flat_map(|rg| { + rg.columns() + .iter() + .enumerate() + .filter_map(|(col_idx, col)| { + col.dictionary_page_offset()?; + let col_desc = schema_descr.column(col_idx); + let parts = col_desc.path().parts(); + // Skip nested columns: their leaf name doesn't match the + // Arrow top-level field name. + (parts.len() == 1).then(|| parts[0].clone()) + }) + }) + .collect(); + let dict_cols = if let Some(allowlist) = &self.rle_column_allowlist { + dict_cols + .into_iter() + .filter(|name| allowlist.contains(name)) + .collect::>() + } else { + dict_cols + }; + if dict_cols.is_empty() { + schema + } else { + let promoted: Vec<_> = schema + .fields() + .iter() + .map(|field| { + if !dict_cols.contains(field.name()) { + return Arc::clone(field); + } + let dict_value_type = match field.data_type() { + DataType::Utf8 => Some(DataType::Utf8), + DataType::LargeUtf8 => Some(DataType::LargeUtf8), + DataType::Utf8View => Some(DataType::Utf8View), + DataType::Binary => Some(DataType::Binary), + DataType::LargeBinary => Some(DataType::LargeBinary), + DataType::BinaryView => Some(DataType::BinaryView), + _ => None, + }; + dict_value_type.map_or_else( + || Arc::clone(field), + |value_type| { + Arc::new( + Field::new( + field.name(), + DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(value_type), + ), + field.is_nullable(), + ) + .with_metadata(field.metadata().clone()), + ) + }, + ) + }) + .collect(); + Schema::new_with_metadata(promoted, schema.metadata().clone()) + } + } else { + schema + }; + Ok(schema) } diff --git a/datafusion/datasource-parquet/src/opener/mod.rs b/datafusion/datasource-parquet/src/opener/mod.rs index 1f4a09a58b2d6..54e4fd09b36a6 100644 --- a/datafusion/datasource-parquet/src/opener/mod.rs +++ b/datafusion/datasource-parquet/src/opener/mod.rs @@ -35,7 +35,7 @@ use crate::row_group_filter::{RowGroupAccessPlanFilter, row_group_in_range}; use crate::{ BloomFilterStatistics, Int96Coercer, ParquetAccessPlan, ParquetFileMetrics, ParquetFileReaderFactory, ParquetRowSelection, ParquetVirtualColumn, RowGroupAccess, - apply_file_schema_type_coercions, + schema_coercion::apply_file_schema_type_coercions_with_rle, }; use arrow::array::RecordBatch; use arrow::datatypes::DataType; @@ -298,6 +298,8 @@ pub(super) struct ParquetMorselizer { /// lists skip container-level pruning. Sourced from /// `datafusion.execution.parquet.max_in_list_size`. pub max_in_list_size: usize, + /// Whether to ask arrow-rs to read promoted dictionary columns directly. + pub enable_rle_to_dictionary: bool, /// Whether to read row groups in reverse order pub reverse_row_groups: bool, /// Optional sort order used to reorder row groups by their min/max statistics. @@ -493,6 +495,7 @@ struct PreparedParquetOpen { predicate_creation_errors: Count, max_predicate_cache_size: Option, max_in_list_size: usize, + enable_rle_to_dictionary: bool, reverse_row_groups: bool, sort_order_for_reorder: Option, preserve_order: bool, @@ -994,6 +997,7 @@ impl ParquetMorselizer { predicate_creation_errors, max_predicate_cache_size: self.max_predicate_cache_size, max_in_list_size: self.max_in_list_size, + enable_rle_to_dictionary: self.enable_rle_to_dictionary, reverse_row_groups: self.reverse_row_groups, sort_order_for_reorder: self.sort_order_for_reorder.clone(), preserve_order: self.preserve_order, @@ -1109,9 +1113,10 @@ impl MetadataLoadedParquetOpen { // desired schema (for example if we want to instruct the parquet // reader to read strings using Utf8View instead). Update if necessary let mut metadata_dirty = false; - if let Some(merged) = apply_file_schema_type_coercions( + if let Some(merged) = apply_file_schema_type_coercions_with_rle( &prepared.logical_file_schema, &physical_file_schema, + prepared.enable_rle_to_dictionary, ) { physical_file_schema = Arc::new(merged); options = options.with_schema(Arc::clone(&physical_file_schema)); @@ -2187,6 +2192,7 @@ mod test { coerce_int96: Option, max_predicate_cache_size: Option, max_in_list_size: usize, + enable_rle_to_dictionary: bool, reverse_row_groups: bool, preserve_order: bool, } @@ -2410,6 +2416,7 @@ mod test { coerce_int96: None, max_predicate_cache_size: None, max_in_list_size: MAX_IN_LIST_SIZE, + enable_rle_to_dictionary: false, reverse_row_groups: false, preserve_order: false, } @@ -2485,6 +2492,11 @@ mod test { self } + fn with_enable_rle_to_dictionary(mut self, enable: bool) -> Self { + self.enable_rle_to_dictionary = enable; + self + } + fn with_metrics(mut self, metrics: ExecutionPlanMetricsSet) -> Self { self.metrics = metrics; self @@ -2595,6 +2607,7 @@ mod test { encryption_factory: None, max_predicate_cache_size: self.max_predicate_cache_size, max_in_list_size: self.max_in_list_size, + enable_rle_to_dictionary: self.enable_rle_to_dictionary, reverse_row_groups: self.reverse_row_groups, sort_order_for_reorder: None, virtual_state, @@ -5709,4 +5722,59 @@ mod test { assert_eq!(rows, 5); } } + + async fn collect_batches( + morselizer: &ParquetMorselizer, + file: PartitionedFile, + ) -> Vec { + let mut stream = open_file(morselizer, file).await.unwrap(); + let mut batches = Vec::new(); + while let Some(batch) = stream.next().await { + batches.push(batch.unwrap()); + } + batches + } + + // Proves the opener passes a promoted binary Dictionary schema to arrow-rs. + #[tokio::test] + async fn test_rle_binary_column_promotion() { + let store = Arc::new(InMemory::new()) as Arc; + let bin_schema = Arc::new(Schema::new(vec![Field::new( + "payload", + DataType::Binary, + true, + )])); + let values = + Arc::new(arrow::array::BinaryArray::from_vec(vec![b"a", b"b", b"a"])); + let batch = RecordBatch::try_new(Arc::clone(&bin_schema), vec![values]).unwrap(); + let props = WriterProperties::builder() + .set_dictionary_enabled(true) + .build(); + let bin_size = write_parquet_batches( + Arc::clone(&store), + "binary.parquet", + vec![batch], + Some(props), + ) + .await; + let dict_bin_schema = Arc::new(Schema::new(vec![Field::new( + "payload", + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Binary)), + true, + )])); + let morselizer = ParquetMorselizerBuilder::new() + .with_store(Arc::clone(&store)) + .with_schema(Arc::clone(&dict_bin_schema)) + .with_enable_rle_to_dictionary(true) + .build(); + let batches = collect_batches( + &morselizer, + PartitionedFile::new("binary.parquet".to_string(), bin_size as u64), + ) + .await; + assert_eq!( + batches[0].schema().field(0).data_type(), + &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Binary)) + ); + } } diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index 828070659b707..807c04cc515ae 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -17,7 +17,7 @@ //! Arrow-schema coercion utilities used by the Parquet reader to make a //! file schema match the table schema (binary→string, regular→view, -//! INT96→Timestamp). +//! INT96→Timestamp, plain->Dictionary). //! //! These helpers are independent of the [`ParquetFormat`](crate::file_format::ParquetFormat) //! type and several have been re-exported at the crate root for use by @@ -84,6 +84,183 @@ pub fn apply_file_schema_type_coercions( )) } +/// Like [`apply_file_schema_type_coercions`], but also coerces compatible +/// string/binary file fields to dictionary types already present in the table +/// schema. +pub(crate) fn apply_file_schema_type_coercions_with_rle( + table_schema: &Schema, + file_schema: &Schema, + enable_rle_to_dictionary: bool, +) -> Option { + let mut needs_view_transform = false; + let mut needs_string_transform = false; + let mut needs_nested_transform = false; + let mut needs_dict_transform = false; + + // Create a mapping of table field names to their data types for fast lookup + // and simultaneously check if we need any transformations + let table_fields: HashMap<_, _> = table_schema + .fields() + .iter() + .map(|field| { + let data_type = field.data_type(); + // Check if we need view type transformation + if matches!(data_type, &DataType::Utf8View | &DataType::BinaryView) { + needs_view_transform = true; + } + // Check if we need string type transformation + if matches!( + data_type, + &DataType::Utf8 | &DataType::LargeUtf8 | &DataType::Utf8View + ) { + needs_string_transform = true; + } + // Nested fields can need transformations even when their parent does not. + if matches!( + data_type, + DataType::Struct(_) + | DataType::List(_) + | DataType::LargeList(_) + | DataType::ListView(_) + | DataType::LargeListView(_) + | DataType::FixedSizeList(_, _) + | DataType::Map(_, _) + ) { + needs_nested_transform = true; + } + if enable_rle_to_dictionary + && matches!(data_type, &DataType::Dictionary(_, _)) + { + needs_dict_transform = true; + } + + (field.name(), data_type) + }) + .collect(); + + // Early return if no transformation needed + if !needs_view_transform + && !needs_string_transform + && !needs_nested_transform + && !needs_dict_transform + { + return None; + } + + let fields: Vec> = file_schema + .fields() + .iter() + .map(|field| { + let field_name = field.name(); + let field_type = field.data_type(); + + // Look up the corresponding field type in the table schema + if let Some(table_type) = table_fields.get(field_name) { + match (table_type, field_type) { + // table schema uses string type, coerce the file schema to use string type + ( + &DataType::Utf8, + DataType::Binary | DataType::LargeBinary | DataType::BinaryView, + ) => { + return field_with_new_type(field, DataType::Utf8); + } + // table schema uses large string type, coerce the file schema to use large string type + ( + &DataType::LargeUtf8, + DataType::Binary | DataType::LargeBinary | DataType::BinaryView, + ) => { + return field_with_new_type(field, DataType::LargeUtf8); + } + // table schema uses string view type, coerce the file schema to use view type + ( + &DataType::Utf8View, + DataType::Binary | DataType::LargeBinary | DataType::BinaryView, + ) => { + return field_with_new_type(field, DataType::Utf8View); + } + // Handle view type conversions + (&DataType::Utf8View, DataType::Utf8 | DataType::LargeUtf8) => { + return field_with_new_type(field, DataType::Utf8View); + } + (&DataType::BinaryView, DataType::Binary | DataType::LargeBinary) => { + return field_with_new_type(field, DataType::BinaryView); + } + // Apply the same coercions to matching fields inside structs. + (DataType::Struct(table_fields), DataType::Struct(file_fields)) => { + if let Some(schema) = apply_file_schema_type_coercions( + &Schema::new(table_fields.clone()), + &Schema::new(file_fields.clone()), + ) { + return field_with_new_type( + field, + DataType::Struct(schema.fields), + ); + } + } + // Container children match by position, regardless of their names. + (DataType::List(table_child), DataType::List(file_child)) + | ( + DataType::LargeList(table_child), + DataType::LargeList(file_child), + ) + | (DataType::ListView(table_child), DataType::ListView(file_child)) + | ( + DataType::LargeListView(table_child), + DataType::LargeListView(file_child), + ) + | ( + DataType::FixedSizeList(table_child, _), + DataType::FixedSizeList(file_child, _), + ) + | (DataType::Map(table_child, _), DataType::Map(file_child, _)) => { + if let Some(schema) = apply_file_schema_type_coercions( + &Schema::new(vec![field_with_new_type( + file_child, + table_child.data_type().clone(), + )]), + &Schema::new(vec![Arc::clone(file_child)]), + ) { + let child = Arc::clone(&schema.fields()[0]); + let new_type = match field_type { + DataType::List(_) => DataType::List(child), + DataType::LargeList(_) => DataType::LargeList(child), + DataType::ListView(_) => DataType::ListView(child), + DataType::LargeListView(_) => { + DataType::LargeListView(child) + } + DataType::FixedSizeList(_, size) => { + DataType::FixedSizeList(child, *size) + } + DataType::Map(_, sorted) => DataType::Map(child, *sorted), + _ => return Arc::clone(field), + }; + return field_with_new_type(field, new_type); + } + } + (DataType::Dictionary(_, _), _) + if enable_rle_to_dictionary + && can_promote_to_dictionary_type(field_type, table_type) => + { + return field_with_new_type(field, (*table_type).clone()); + } + _ => {} + } + } + // If no transformation is needed, keep the original field + Arc::clone(field) + }) + .collect(); + + if fields.iter().eq(file_schema.fields().iter()) { + return None; + } + + Some(Schema::new_with_metadata( + fields, + file_schema.metadata.clone(), + )) +} + /// Coerce `file_fields` towards `table_fields`, matching fields by name. /// /// File fields with no counterpart in `table_fields` are kept unchanged and @@ -242,6 +419,114 @@ fn coerce_map_entries( Some(field_with_new_type(file_entries, DataType::Struct(fields))) } +fn dictionary_value_type(data_type: &DataType) -> Option<&DataType> { + match data_type { + DataType::Dictionary(_, value_type) => Some(value_type.as_ref()), + _ => None, + } +} + +// Find the value type that can represent both sides without narrowing offsets +// or crossing string/binary families. +fn common_dictionary_value_type( + field_type: &DataType, + dictionary_value_type: &DataType, +) -> Option { + let field_type = match field_type { + DataType::Dictionary(_, field_value_type) => field_value_type.as_ref(), + _ => field_type, + }; + + match (field_type, dictionary_value_type) { + (DataType::Utf8, DataType::Utf8) => Some(DataType::Utf8), + (DataType::Utf8 | DataType::LargeUtf8, DataType::Utf8 | DataType::LargeUtf8) => { + Some(DataType::LargeUtf8) + } + (DataType::Binary, DataType::Binary) => Some(DataType::Binary), + ( + DataType::Binary | DataType::LargeBinary, + DataType::Binary | DataType::LargeBinary, + ) => Some(DataType::LargeBinary), + _ => None, + } +} + +/// Allows safe widening into the table dictionary type. +/// - `Utf8` to `Dictionary(Int32, LargeUtf8)`: allowed +/// - `LargeBinary` to `Dictionary(Int32, Binary)`: rejected +fn can_promote_to_dictionary_type( + file_field_type: &DataType, + table_dictionary_type: &DataType, +) -> bool { + dictionary_value_type(table_dictionary_type).is_some_and(|dictionary_value_type| { + common_dictionary_value_type(file_field_type, dictionary_value_type) + .is_some_and(|common_type| &common_type == dictionary_value_type) + }) +} + +/// Normalize per-file schemas so that a column promoted to `Dictionary` in +/// *any* file is promoted to the same `Dictionary` type in *all* files. +/// +/// This lets [`Schema::try_merge`] accept directories that mix dictionary and +/// plain encodings for the same column. +pub(crate) fn uniform_dict_schemas(schemas: Vec) -> Vec { + // First pass: record the dictionary type for every column that is Dictionary in + // at least one schema. + let mut dict_types: HashMap = HashMap::new(); + for schema in &schemas { + for field in schema.fields() { + if matches!(field.data_type(), DataType::Dictionary(_, _)) { + dict_types + .entry(field.name().clone()) + .or_insert_with(|| field.data_type().clone()); + } + } + } + if dict_types.is_empty() { + return schemas; + } + + // Widen the recorded dictionary value type before promoting plain fields. + for schema in &schemas { + for field in schema.fields() { + let Some(dict_type) = dict_types.get_mut(field.name()) else { + continue; + }; + let DataType::Dictionary(key_type, value_type) = dict_type else { + continue; + }; + let key_type = key_type.clone(); + let value_type = value_type.as_ref().clone(); + if let Some(common_type) = + common_dictionary_value_type(field.data_type(), &value_type) + { + *dict_type = DataType::Dictionary(key_type, Box::new(common_type)); + } + } + } + + // Only promote fields whose type family matches the dictionary value type, + // so schema normalization does not reinterpret binary bytes as UTF-8. + schemas + .into_iter() + .map(|schema| { + let fields: Vec> = schema + .fields() + .iter() + .map(|field| { + if let Some(dict_type) = dict_types.get(field.name()) + && can_promote_to_dictionary_type(field.data_type(), dict_type) + { + return field_with_new_type(field, dict_type.clone()); + } + Arc::clone(field) + }) + .collect(); + Schema::new_with_metadata(fields, schema.metadata().clone()) + }) + .collect() +} + /// Coerces the file schema's Timestamps to the provided TimeUnit if the /// Parquet schema contains INT96. /// @@ -548,6 +833,25 @@ pub fn transform_schema_to_view(schema: &Schema) -> Schema { DataType::Binary | DataType::LargeBinary => { field_with_new_type(field, DataType::BinaryView) } + // Also rewrite the value type inside dictionary columns so that + // dict<_, Utf8> / dict<_, LargeUtf8> become dict<_, Utf8View> and + // dict<_, Binary> / dict<_, LargeBinary> become dict<_, BinaryView>. + DataType::Dictionary(key_type, value_type) => { + let new_value = match value_type.as_ref() { + DataType::Utf8 | DataType::LargeUtf8 => Some(DataType::Utf8View), + DataType::Binary | DataType::LargeBinary => { + Some(DataType::BinaryView) + } + _ => None, + }; + match new_value { + Some(vt) => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(vt)), + ), + None => Arc::clone(field), + } + } _ => Arc::clone(field), }) .collect(); @@ -563,6 +867,22 @@ pub fn transform_binary_to_string(schema: &Schema) -> Schema { DataType::Binary => field_with_new_type(field, DataType::Utf8), DataType::LargeBinary => field_with_new_type(field, DataType::LargeUtf8), DataType::BinaryView => field_with_new_type(field, DataType::Utf8View), + // Keep binary_as_string consistent for dictionary value types. + DataType::Dictionary(key_type, value_type) => match value_type.as_ref() { + DataType::Binary => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(DataType::Utf8)), + ), + DataType::LargeBinary => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(DataType::LargeUtf8)), + ), + DataType::BinaryView => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(DataType::Utf8View)), + ), + _ => Arc::clone(field), + }, _ => Arc::clone(field), }) .collect(); @@ -1599,4 +1919,256 @@ mod tests { ), } } + + fn dict(value_type: DataType) -> DataType { + DataType::Dictionary(Box::new(DataType::Int32), Box::new(value_type)) + } + + fn one_field_schema(data_type: DataType) -> Schema { + Schema::new(vec![Field::new("col", data_type, true)]) + } + + #[test] + fn uniform_dict_schemas_respects_value_type_families() { + // String/binary dictionary promotion must preserve the concrete value type + // family so schema merging does not silently reinterpret bytes as UTF-8. + let cases = vec![ + ( + "utf8", + dict(DataType::Utf8), + DataType::Utf8, + dict(DataType::Utf8), + dict(DataType::Utf8), + ), + ( + "large_utf8", + dict(DataType::LargeUtf8), + DataType::LargeUtf8, + dict(DataType::LargeUtf8), + dict(DataType::LargeUtf8), + ), + ( + "utf8_large_utf8", + dict(DataType::Utf8), + DataType::LargeUtf8, + dict(DataType::LargeUtf8), + dict(DataType::LargeUtf8), + ), + ( + "binary", + dict(DataType::Binary), + DataType::Binary, + dict(DataType::Binary), + dict(DataType::Binary), + ), + ( + "large_binary", + dict(DataType::LargeBinary), + DataType::LargeBinary, + dict(DataType::LargeBinary), + dict(DataType::LargeBinary), + ), + ( + "binary_large_binary", + dict(DataType::Binary), + DataType::LargeBinary, + dict(DataType::LargeBinary), + dict(DataType::LargeBinary), + ), + ( + "binary_not_utf8", + dict(DataType::Utf8), + DataType::Binary, + dict(DataType::Utf8), + DataType::Binary, + ), + ( + "utf8_not_binary", + dict(DataType::Binary), + DataType::Utf8, + dict(DataType::Binary), + DataType::Utf8, + ), + ]; + + for (name, dict_type, plain_type, expected_dict_type, expected_plain_type) in + cases + { + let result = uniform_dict_schemas(vec![ + one_field_schema(dict_type.clone()), + one_field_schema(plain_type), + ]); + + assert_eq!( + result[0].field(0).data_type(), + &expected_dict_type, + "{name}" + ); + assert_eq!( + result[1].field(0).data_type(), + &expected_plain_type, + "{name}" + ); + } + } + + #[test] + fn rle_schema_coercion_respects_dictionary_value_type() { + // Scan-time schema coercion can only use dictionary types already chosen by + // schema inference, and must leave incompatible file fields unchanged. + let cases = vec![ + ( + "utf8", + dict(DataType::Utf8), + DataType::Utf8, + dict(DataType::Utf8), + ), + ( + "large_utf8", + dict(DataType::LargeUtf8), + DataType::LargeUtf8, + dict(DataType::LargeUtf8), + ), + ( + "large_utf8_to_utf8_dict_not_safe", + dict(DataType::Utf8), + DataType::LargeUtf8, + DataType::LargeUtf8, + ), + ( + "utf8_to_large_utf8_dict", + dict(DataType::LargeUtf8), + DataType::Utf8, + dict(DataType::LargeUtf8), + ), + ( + "binary", + dict(DataType::Binary), + DataType::Binary, + dict(DataType::Binary), + ), + ( + "large_binary", + dict(DataType::LargeBinary), + DataType::LargeBinary, + dict(DataType::LargeBinary), + ), + ( + "large_binary_to_binary_dict_not_safe", + dict(DataType::Binary), + DataType::LargeBinary, + DataType::LargeBinary, + ), + ( + "binary_to_large_binary_dict", + dict(DataType::LargeBinary), + DataType::Binary, + dict(DataType::LargeBinary), + ), + ( + "binary_not_utf8", + dict(DataType::Utf8), + DataType::Binary, + DataType::Binary, + ), + ( + "utf8_not_binary", + dict(DataType::Binary), + DataType::Utf8, + DataType::Utf8, + ), + ( + "binary_not_int64_dict", + dict(DataType::Int64), + DataType::Binary, + DataType::Binary, + ), + ]; + + for (name, table_type, file_type, expected_type) in cases { + let table_schema = one_field_schema(table_type); + let file_schema = one_field_schema(file_type); + + let output_type = apply_file_schema_type_coercions_with_rle( + &table_schema, + &file_schema, + true, + ) + .as_ref() + .map(|s| s.field(0).data_type().clone()) + .unwrap_or_else(|| file_schema.field(0).data_type().clone()); + + assert_eq!(output_type, expected_type, "{name}"); + } + + let table_schema = Schema::new(vec![ + Field::new("dict_col", dict(DataType::Utf8), true), + Field::new("string_col", DataType::Utf8, true), + ]); + let file_schema = Schema::new(vec![ + Field::new("dict_col", DataType::Utf8, true), + Field::new("string_col", DataType::Binary, true), + ]); + let result = + apply_file_schema_type_coercions_with_rle(&table_schema, &file_schema, false) + .unwrap(); + + assert_eq!(result.field(0).data_type(), &DataType::Utf8); + assert_eq!(result.field(1).data_type(), &DataType::Utf8); + } + + #[test] + fn transform_binary_to_string_rewrites_dict_binary_value_type() { + // binary_as_string must also rewrite dictionary value types. + let schema = Schema::new(vec![ + Field::new( + "rle_binary", + DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::Binary), + ), + true, + ), + Field::new( + "rle_large_binary", + DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::LargeBinary), + ), + true, + ), + Field::new("plain_binary", DataType::Binary, true), + Field::new( + "rle_utf8", + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + true, + ), + ]); + + let result = transform_binary_to_string(&schema); + + assert_eq!( + result.field(0).data_type(), + &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + "Dict(Int32, Binary) must become Dict(Int32, Utf8)" + ); + assert_eq!( + result.field(1).data_type(), + &DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::LargeUtf8) + ), + "Dict(Int32, LargeBinary) must become Dict(Int32, LargeUtf8)" + ); + assert_eq!( + result.field(2).data_type(), + &DataType::Utf8, + "plain Binary must become Utf8" + ); + assert_eq!( + result.field(3).data_type(), + &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + "Dict(Int32, Utf8) must be unchanged" + ); + } } diff --git a/datafusion/datasource-parquet/src/source.rs b/datafusion/datasource-parquet/src/source.rs index 42e7116dff73c..2bc16e5a191fa 100644 --- a/datafusion/datasource-parquet/src/source.rs +++ b/datafusion/datasource-parquet/src/source.rs @@ -677,6 +677,10 @@ impl FileSource for ParquetSource { encryption_factory: self.get_encryption_factory_with_config(), max_predicate_cache_size: self.max_predicate_cache_size(), max_in_list_size: self.max_in_list_size(), + enable_rle_to_dictionary: self + .table_parquet_options + .global + .enable_rle_to_dictionary, reverse_row_groups: self.reverse_row_groups, sort_order_for_reorder: self.sort_order_for_reorder.clone(), virtual_state, diff --git a/datafusion/proto-common/proto/datafusion_common.proto b/datafusion/proto-common/proto/datafusion_common.proto index bdabb0777ef81..6e55b011e49b4 100644 --- a/datafusion/proto-common/proto/datafusion_common.proto +++ b/datafusion/proto-common/proto/datafusion_common.proto @@ -626,6 +626,8 @@ message ParquetOptions { uint64 max_in_list_size = 38; + bool enable_rle_to_dictionary = 39; + string created_by = 16; oneof coerce_int96_opt { diff --git a/datafusion/proto-common/src/from_proto/mod.rs b/datafusion/proto-common/src/from_proto/mod.rs index 0bde696c7c46b..b58775d86448e 100644 --- a/datafusion/proto-common/src/from_proto/mod.rs +++ b/datafusion/proto-common/src/from_proto/mod.rs @@ -1289,6 +1289,7 @@ impl TryFrom<&protobuf::ParquetOptions> for ParquetOptions { } }).transpose()?, content_defined_chunking: value.content_defined_chunking.map(ParquetCdcOptions::try_from).transpose()?.unwrap_or_default(), + enable_rle_to_dictionary: value.enable_rle_to_dictionary, }) } } diff --git a/datafusion/proto-common/src/generated/pbjson.rs b/datafusion/proto-common/src/generated/pbjson.rs index a1d86c5fe3187..50db5816f53b7 100644 --- a/datafusion/proto-common/src/generated/pbjson.rs +++ b/datafusion/proto-common/src/generated/pbjson.rs @@ -4058,7 +4058,7 @@ impl serde::Serialize for ExplainAnalyzeCategoriesNode { if !self.only.is_empty() { let v = self.only.iter().copied().map(|v| { MetricCategory::try_from(v) - .map_err(|_| serde::ser::Error::custom(format!("Invalid variant {}", v))) + .map_err(|_| serde::ser::Error::custom(format!("Invalid variant {v}"))) }).collect::, _>>()?; struct_ser.serialize_field("only", &v)?; } @@ -6491,6 +6491,9 @@ impl serde::Serialize for ParquetOptions { if self.max_in_list_size != 0 { len += 1; } + if self.enable_rle_to_dictionary != false { + len += 1; + } if !self.created_by.is_empty() { len += 1; } @@ -6616,6 +6619,9 @@ impl serde::Serialize for ParquetOptions { #[allow(clippy::needless_borrows_for_generic_args)] struct_ser.serialize_field("maxInListSize", ToString::to_string(&self.max_in_list_size).as_str())?; } + if self.enable_rle_to_dictionary != false { + struct_ser.serialize_field("enableRleToDictionary", &self.enable_rle_to_dictionary)?; + } if !self.created_by.is_empty() { struct_ser.serialize_field("createdBy", &self.created_by)?; } @@ -6776,6 +6782,8 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { "maxRowGroupSize", "max_in_list_size", "maxInListSize", + "enable_rle_to_dictionary", + "enableRleToDictionary", "created_by", "createdBy", "content_defined_chunking", @@ -6829,6 +6837,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { DataPageRowCountLimit, MaxRowGroupSize, MaxInListSize, + EnableRleToDictionary, CreatedBy, ContentDefinedChunking, MetadataSizeHint, @@ -6886,6 +6895,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { "dataPageRowCountLimit" | "data_page_row_count_limit" => Ok(GeneratedField::DataPageRowCountLimit), "maxRowGroupSize" | "max_row_group_size" => Ok(GeneratedField::MaxRowGroupSize), "maxInListSize" | "max_in_list_size" => Ok(GeneratedField::MaxInListSize), + "enableRleToDictionary" | "enable_rle_to_dictionary" => Ok(GeneratedField::EnableRleToDictionary), "createdBy" | "created_by" => Ok(GeneratedField::CreatedBy), "contentDefinedChunking" | "content_defined_chunking" => Ok(GeneratedField::ContentDefinedChunking), "metadataSizeHint" | "metadata_size_hint" => Ok(GeneratedField::MetadataSizeHint), @@ -6941,6 +6951,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { let mut data_page_row_count_limit__ = None; let mut max_row_group_size__ = None; let mut max_in_list_size__ = None; + let mut enable_rle_to_dictionary__ = None; let mut created_by__ = None; let mut content_defined_chunking__ = None; let mut metadata_size_hint_opt__ = None; @@ -7100,6 +7111,12 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } + GeneratedField::EnableRleToDictionary => { + if enable_rle_to_dictionary__.is_some() { + return Err(serde::de::Error::duplicate_field("enableRleToDictionary")); + } + enable_rle_to_dictionary__ = Some(map_.next_value()?); + } GeneratedField::CreatedBy => { if created_by__.is_some() { return Err(serde::de::Error::duplicate_field("createdBy")); @@ -7214,6 +7231,7 @@ impl<'de> serde::Deserialize<'de> for ParquetOptions { data_page_row_count_limit: data_page_row_count_limit__.unwrap_or_default(), max_row_group_size: max_row_group_size__.unwrap_or_default(), max_in_list_size: max_in_list_size__.unwrap_or_default(), + enable_rle_to_dictionary: enable_rle_to_dictionary__.unwrap_or_default(), created_by: created_by__.unwrap_or_default(), content_defined_chunking: content_defined_chunking__, metadata_size_hint_opt: metadata_size_hint_opt__, diff --git a/datafusion/proto-common/src/generated/prost.rs b/datafusion/proto-common/src/generated/prost.rs index f856c1aa3959c..d88b9275baaa1 100644 --- a/datafusion/proto-common/src/generated/prost.rs +++ b/datafusion/proto-common/src/generated/prost.rs @@ -875,6 +875,8 @@ pub struct ParquetOptions { pub max_row_group_size: u64, #[prost(uint64, tag = "38")] pub max_in_list_size: u64, + #[prost(bool, tag = "39")] + pub enable_rle_to_dictionary: bool, #[prost(string, tag = "16")] pub created_by: ::prost::alloc::string::String, #[prost(message, optional, tag = "35")] diff --git a/datafusion/proto-common/src/to_proto/mod.rs b/datafusion/proto-common/src/to_proto/mod.rs index b8eb46f2c8da8..ab9ed414638af 100644 --- a/datafusion/proto-common/src/to_proto/mod.rs +++ b/datafusion/proto-common/src/to_proto/mod.rs @@ -968,6 +968,7 @@ impl TryFrom<&ParquetOptions> for protobuf::ParquetOptions { max_predicate_cache_size_opt: value.max_predicate_cache_size.map(|v| protobuf::parquet_options::MaxPredicateCacheSizeOpt::MaxPredicateCacheSize(v as u64)), max_row_group_bytes_opt: value.max_row_group_bytes.map(|v| protobuf::parquet_options::MaxRowGroupBytesOpt::MaxRowGroupBytes(v.get() as u64)), content_defined_chunking: Some((&value.content_defined_chunking).into()), + enable_rle_to_dictionary: value.enable_rle_to_dictionary, }) } } diff --git a/datafusion/proto-models/src/from_proto.rs b/datafusion/proto-models/src/from_proto.rs index fdf32396475b7..51969cb4d1a7e 100644 --- a/datafusion/proto-models/src/from_proto.rs +++ b/datafusion/proto-models/src/from_proto.rs @@ -389,6 +389,7 @@ impl TryFrom<&ParquetOptionsProto> for ParquetOptions { } }) .transpose()?, + enable_rle_to_dictionary: proto.enable_rle_to_dictionary, dictionary_enabled: proto.dictionary_enabled_opt.as_ref().map(|opt| { match opt { parquet_options::DictionaryEnabledOpt::DictionaryEnabled( diff --git a/datafusion/proto-models/src/generated/datafusion_proto_common.rs b/datafusion/proto-models/src/generated/datafusion_proto_common.rs index f856c1aa3959c..d88b9275baaa1 100644 --- a/datafusion/proto-models/src/generated/datafusion_proto_common.rs +++ b/datafusion/proto-models/src/generated/datafusion_proto_common.rs @@ -875,6 +875,8 @@ pub struct ParquetOptions { pub max_row_group_size: u64, #[prost(uint64, tag = "38")] pub max_in_list_size: u64, + #[prost(bool, tag = "39")] + pub enable_rle_to_dictionary: bool, #[prost(string, tag = "16")] pub created_by: ::prost::alloc::string::String, #[prost(message, optional, tag = "35")] diff --git a/datafusion/proto/src/logical_plan/file_formats.rs b/datafusion/proto/src/logical_plan/file_formats.rs index 7c7a8ef639457..b12db665fc107 100644 --- a/datafusion/proto/src/logical_plan/file_formats.rs +++ b/datafusion/proto/src/logical_plan/file_formats.rs @@ -334,6 +334,33 @@ mod parquet { ParquetOptions::default().writer_version ); } + + #[test] + fn enable_rle_to_dictionary_round_trips_through_codec() { + use datafusion_common::config::TableParquetOptions; + let mut options = TableParquetOptions::default(); + options.global.enable_rle_to_dictionary = true; + let original: Arc = Arc::new(ParquetFormatFactory { + options: Some(options), + }); + + let mut buf = Vec::new(); + ParquetLogicalExtensionCodec + .try_encode_file_format(&mut buf, Arc::clone(&original)) + .expect("encode parquet options"); + + let decoded = ParquetLogicalExtensionCodec + .try_decode_file_format(&buf, &TaskContext::default()) + .expect("decode parquet options"); + let decoded_options = decoded + .downcast_ref::() + .expect("parquet format factory") + .options + .as_ref() + .expect("parquet options"); + + assert!(decoded_options.global.enable_rle_to_dictionary); + } } } #[cfg(feature = "parquet")] diff --git a/datafusion/sqllogictest/test_files/information_schema.slt b/datafusion/sqllogictest/test_files/information_schema.slt index f6633a165df18..42a541c5c02f1 100644 --- a/datafusion/sqllogictest/test_files/information_schema.slt +++ b/datafusion/sqllogictest/test_files/information_schema.slt @@ -251,6 +251,7 @@ datafusion.execution.parquet.data_pagesize_limit 1048576 datafusion.execution.parquet.dictionary_enabled true datafusion.execution.parquet.dictionary_page_size_limit 1048576 datafusion.execution.parquet.enable_page_index true +datafusion.execution.parquet.enable_rle_to_dictionary false datafusion.execution.parquet.encoding NULL datafusion.execution.parquet.force_filter_selections false datafusion.execution.parquet.max_in_list_size 20 @@ -413,6 +414,7 @@ datafusion.execution.parquet.data_pagesize_limit 1048576 (writing) Sets best eff datafusion.execution.parquet.dictionary_enabled true (writing) Sets if dictionary encoding is enabled. If NULL, uses default parquet writer setting datafusion.execution.parquet.dictionary_page_size_limit 1048576 (writing) Sets best effort maximum dictionary page size, in bytes datafusion.execution.parquet.enable_page_index true (reading) If true, reads the Parquet data page level metadata (the Page Index), if present, to reduce the I/O and number of rows decoded. +datafusion.execution.parquet.enable_rle_to_dictionary false (reading) If true, top-level string and binary Parquet columns with dictionary pages are inferred and scanned as `Dictionary` / `Dictionary` instead of their plain value type. This applies only when DataFusion infers the table schema. Tables with a user-supplied schema are not promoted because the Parquet footer is not read at DDL time, so dictionary pages cannot be detected per column. See datafusion.execution.parquet.encoding NULL (writing) Sets default encoding for any column. Valid values are: plain, plain_dictionary, rle, bit_packed, delta_binary_packed, delta_length_byte_array, delta_byte_array, rle_dictionary, and byte_stream_split. These values are not case sensitive. If NULL, uses default parquet writer setting datafusion.execution.parquet.force_filter_selections false (reading) Force the use of RowSelections for filter results, when pushdown_filters is enabled. If false, the reader will automatically choose between a RowSelection and a Bitmap based on the number and pattern of selected rows. datafusion.execution.parquet.max_in_list_size 20 Maximum number of input values in an `IN (...)` list eligible for min/max pruning. Lists above this cap, or a cap of 0, skip this rewrite; other predicates and Bloom-filter pruning remain available. Within the cap, nonempty lists of at most 20 values use the existing per-value rewrite. Larger literal lists use a compact representation when the column type is string, variable-length binary, integer, decimal, date, time, timestamp, or duration. This applies to both `IN` and `NOT IN`, including lists with NULL members. `NOT IN` with NULL and all-NULL `IN` lists cannot match any rows. Compact lists containing NULL do not use the fully-matched-row-group optimization. Floating-point and other lists retain the existing per-value rewrite, so raising the cap can make those predicates expensive to build and evaluate. Defaults to 20. diff --git a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt new file mode 100644 index 0000000000000..685d904f8cd90 --- /dev/null +++ b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt @@ -0,0 +1,250 @@ +# 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. + +# Tests for datafusion.execution.parquet.enable_rle_to_dictionary. The flag +# applies only when DataFusion infers the table schema; explicit schemas are +# not promoted. + +# Write with dictionary_enabled=true to guarantee RLE_DICTIONARY encoding. +query I +COPY ( + SELECT column1 AS product_category, column2 AS status + FROM (VALUES + ('electronics', 'active'), + ('clothing', 'active'), + ('electronics', 'inactive'), + ('books', 'inactive'), + ('furniture', 'pending'), + ('books', 'active') + ) +) +TO 'test_files/scratch/parquet_rle_to_dictionary/products.parquet' +STORED AS PARQUET +OPTIONS ('format.dictionary_enabled' true); +---- +6 + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = false; + +statement ok +CREATE EXTERNAL TABLE products_utf8 +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/products.parquet'; + +query TT +SELECT DISTINCT arrow_typeof(product_category), arrow_typeof(status) FROM products_utf8; +---- +Utf8View Utf8View + +query TI rowsort +SELECT product_category, COUNT(*) FROM products_utf8 GROUP BY product_category; +---- +books 2 +clothing 1 +electronics 2 +furniture 1 + +statement ok +DROP TABLE products_utf8; + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = true; + +statement ok +CREATE EXTERNAL TABLE products_dict +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/products.parquet'; + +query TT +SELECT DISTINCT arrow_typeof(product_category), arrow_typeof(status) FROM products_dict; +---- +Dictionary(Int32, Utf8) Dictionary(Int32, Utf8) + +query TI rowsort +SELECT product_category, COUNT(*) FROM products_dict GROUP BY product_category; +---- +books 2 +clothing 1 +electronics 2 +furniture 1 + +# Predicates and aggregates must work against dictionary scan output. +query TT rowsort +SELECT product_category, status FROM products_dict WHERE status = 'active'; +---- +books active +clothing active +electronics active + +query TI rowsort +SELECT product_category, COUNT(*) FROM products_dict WHERE status != 'inactive' GROUP BY product_category; +---- +books 1 +clothing 1 +electronics 1 +furniture 1 + +query TT +SELECT DISTINCT arrow_typeof(product_category), arrow_typeof(product_count) +FROM ( + SELECT product_category, COUNT(*) AS product_count + FROM products_dict + WHERE status != 'inactive' + GROUP BY product_category +); +---- +Dictionary(Int32, Utf8) Int64 + +statement ok +DROP TABLE products_dict; + +# Per-column writer options: only columns with dictionary pages are promoted. + +statement ok +COPY ( + SELECT column1 AS env, column2 AS region, column3 AS build + FROM (VALUES + ('prod', 'us-east', 'debug'), + ('staging', 'us-west', 'release'), + ('prod', 'us-east', 'debug') + ) +) +TO 'test_files/scratch/parquet_rle_to_dictionary/selective.parquet' +STORED AS PARQUET +OPTIONS ( + 'format.dictionary_enabled' false, + 'format.dictionary_enabled::env' true, + 'format.dictionary_enabled::region' true +); + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = false; + +statement ok +CREATE EXTERNAL TABLE selective_utf8 +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/selective.parquet'; + +query TTT +SELECT DISTINCT arrow_typeof(env), arrow_typeof(region), arrow_typeof(build) +FROM selective_utf8; +---- +Utf8View Utf8View Utf8View + +statement ok +DROP TABLE selective_utf8; + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = true; + +statement ok +CREATE EXTERNAL TABLE selective_dict +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/selective.parquet'; + +query TTT +SELECT DISTINCT arrow_typeof(env), arrow_typeof(region), arrow_typeof(build) +FROM selective_dict; +---- +Dictionary(Int32, Utf8) Dictionary(Int32, Utf8) Utf8View + +statement ok +DROP TABLE selective_dict; + +# Cross-file mixed encoding: one file has dictionary pages and one is plain. +# Schema inference must merge both as Dictionary(Int32, Utf8). + +statement ok +COPY (SELECT column1 AS category FROM (VALUES ('electronics'), ('electronics'))) +TO 'test_files/scratch/parquet_rle_to_dictionary/mixed/rle.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' true); + +statement ok +COPY (SELECT column1 AS category FROM (VALUES ('books'), ('books'))) +TO 'test_files/scratch/parquet_rle_to_dictionary/mixed/plain.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' false); + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = true; + +statement ok +CREATE EXTERNAL TABLE mixed_encoding +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/mixed/'; + +query TI +SELECT arrow_typeof(category), COUNT(*) FROM mixed_encoding GROUP BY arrow_typeof(category); +---- +Dictionary(Int32, Utf8) 4 + +statement ok +DROP TABLE mixed_encoding; + +# Explicit-schema tables are not promoted. + +statement ok +CREATE EXTERNAL TABLE explicit_schema ( + product_category VARCHAR, + status VARCHAR +) +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/products.parquet'; + +query TT +SELECT DISTINCT arrow_typeof(product_category), arrow_typeof(status) FROM explicit_schema; +---- +Utf8View Utf8View + +# Predicates must still compile against the explicit plain schema. +query TT rowsort +SELECT product_category, status FROM explicit_schema WHERE status = 'active'; +---- +books active +clothing active +electronics active + +statement ok +DROP TABLE explicit_schema; + +# Family-compatibility safety: a plain Binary field must not be promoted to +# Dictionary(Int32, Utf8). Schema inference should reject this mixed dataset. + +statement ok +COPY (SELECT 'hello' AS payload FROM (VALUES (1), (2))) +TO 'test_files/scratch/parquet_rle_to_dictionary/compat/utf8_rle.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' true); + +statement ok +COPY (SELECT arrow_cast(X'6865', 'Binary') AS payload FROM (VALUES (1), (2))) +TO 'test_files/scratch/parquet_rle_to_dictionary/compat/binary_plain.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' false); + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = true; + +# Registering incompatible Dict Utf8 and plain Binary files must fail. +statement error Arrow error: Schema error: Fail to merge schema field 'payload' because the from data_type = Dictionary\(Int32, Utf8\) does not equal Binary +CREATE EXTERNAL TABLE mixed_str_binary +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/compat/'; + +statement ok +reset datafusion.execution.parquet.enable_rle_to_dictionary; + +statement ok +RESET datafusion.catalog.create_default_catalog_and_schema; diff --git a/docs/source/user-guide/configs.md b/docs/source/user-guide/configs.md index ba8ea52bb0850..be22d4c2c9bd7 100644 --- a/docs/source/user-guide/configs.md +++ b/docs/source/user-guide/configs.md @@ -94,6 +94,7 @@ The following configuration settings are available: | datafusion.execution.parquet.bloom_filter_on_read | true | (reading) Use any available bloom filters when reading parquet files | | datafusion.execution.parquet.max_predicate_cache_size | NULL | (reading) The maximum predicate cache size, in bytes. When `pushdown_filters` is enabled, sets the maximum memory used to cache the results of predicate evaluation between filter evaluation and output generation. Decreasing this value will reduce memory usage, but may increase IO and CPU usage. None means use the default parquet reader setting. 0 means no caching. | | datafusion.execution.parquet.max_in_list_size | 20 | Maximum number of input values in an `IN (...)` list eligible for min/max pruning. Lists above this cap, or a cap of 0, skip this rewrite; other predicates and Bloom-filter pruning remain available. Within the cap, nonempty lists of at most 20 values use the existing per-value rewrite. Larger literal lists use a compact representation when the column type is string, variable-length binary, integer, decimal, date, time, timestamp, or duration. This applies to both `IN` and `NOT IN`, including lists with NULL members. `NOT IN` with NULL and all-NULL `IN` lists cannot match any rows. Compact lists containing NULL do not use the fully-matched-row-group optimization. Floating-point and other lists retain the existing per-value rewrite, so raising the cap can make those predicates expensive to build and evaluate. Defaults to 20. | +| datafusion.execution.parquet.enable_rle_to_dictionary | false | (reading) If true, top-level string and binary Parquet columns with dictionary pages are inferred and scanned as `Dictionary` / `Dictionary` instead of their plain value type. This applies only when DataFusion infers the table schema. Tables with a user-supplied schema are not promoted because the Parquet footer is not read at DDL time, so dictionary pages cannot be detected per column. See | | datafusion.execution.parquet.data_pagesize_limit | 1048576 | (writing) Sets best effort maximum size of data page in bytes | | datafusion.execution.parquet.write_batch_size | 1024 | (writing) Sets write_batch_size in rows | | datafusion.execution.parquet.writer_version | 1.0 | (writing) Sets parquet writer version valid values are "1.0" and "2.0" | From 3054b529c1cd48128189c3c4f4f777ed9a7f5dc1 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 31 Aug 2026 14:49:57 -0400 Subject: [PATCH 2/9] remove rle_column_allowlist API, fix unconditional dict rewrite regression in schema coercion - Drop rle_column_allowlist from ParquetFormat and DFParquetMetadata: no callers existed and new public API is hard to remove post-release - Remove Dictionary arm from transform_schema_to_view and transform_binary_to_string: those arms fired unconditionally regardless of enable_rle_to_dictionary, causing pre-existing dict columns to get wrong value types when the flag was off - Remove the unit test that covered the now-deleted transform_binary_to_string dict arm - Drop Utf8View/BinaryView arms added to common_dictionary_value_type: the parquet reader never produces view types in the physical file schema so they were dead code - Clean up parquet_rle_to_dictionary.slt: remove redundant SET and spurious RESET --- .../datasource-parquet/src/file_format.rs | 32 +------ datafusion/datasource-parquet/src/metadata.rs | 28 +----- .../datasource-parquet/src/schema_coercion.rs | 90 ------------------- .../test_files/parquet_rle_to_dictionary.slt | 6 -- 4 files changed, 3 insertions(+), 153 deletions(-) diff --git a/datafusion/datasource-parquet/src/file_format.rs b/datafusion/datasource-parquet/src/file_format.rs index 8bb8a475a5f89..9fff6dad14b4a 100644 --- a/datafusion/datasource-parquet/src/file_format.rs +++ b/datafusion/datasource-parquet/src/file_format.rs @@ -145,9 +145,6 @@ impl Debug for ParquetFormatFactory { #[derive(Debug, Default)] pub struct ParquetFormat { options: TableParquetOptions, - /// When set, restricts RLE->Dictionary promotion to only these columns. - /// Overrides `enable_rle_to_dictionary`; see [`Self::with_rle_column_allowlist`]. - rle_column_allowlist: Option>, } impl ParquetFormat { @@ -212,22 +209,6 @@ impl ParquetFormat { &self.options } - /// Restrict RLE->Dictionary promotion to a named subset of columns. - /// - /// Only columns in `columns` that also have dictionary pages in the file - /// are promoted to `Dictionary(Int32, …)` in the inferred schema. - /// Overrides `enable_rle_to_dictionary` when set; an empty set disables - /// promotion entirely. See . - pub fn with_rle_column_allowlist(mut self, columns: HashSet) -> Self { - self.rle_column_allowlist = Some(columns); - self - } - - /// Returns the RLE->Dictionary column allowlist, if set. - pub fn rle_column_allowlist(&self) -> Option<&HashSet> { - self.rle_column_allowlist.as_ref() - } - /// Get [`schema_force_view_types`] /// /// [`schema_force_view_types`]: datafusion_common::config::ParquetOptions::schema_force_view_types @@ -375,7 +356,7 @@ impl FileFormat for ParquetFormat { &object.location, ) .await?; - let mut meta = DFParquetMetadata::new(store.as_ref(), object) + let meta = DFParquetMetadata::new(store.as_ref(), object) .with_metadata_size_hint(self.metadata_size_hint()) .with_decryption_properties(file_decryption_properties) .with_file_metadata_cache(Some(Arc::clone(&file_metadata_cache))) @@ -384,9 +365,6 @@ impl FileFormat for ParquetFormat { .with_enable_rle_to_dictionary( self.options.global.enable_rle_to_dictionary, ); - if let Some(allowlist) = &self.rle_column_allowlist { - meta = meta.with_rle_column_allowlist(allowlist.clone()); - } let result = meta.fetch_schema_with_location().await?; Ok::<_, DataFusionError>(result) }) @@ -537,13 +515,7 @@ impl FileFormat for ParquetFormat { .downcast_ref::() .cloned() .ok_or_else(|| internal_datafusion_err!("Expected ParquetSource"))?; - let mut source_options = self.options.clone(); - // An allowlist implies dict-typed fields in the table schema, so the reader - // must also have the flag set to coerce those columns at scan time. - if self.rle_column_allowlist.is_some() { - source_options.global.enable_rle_to_dictionary = true; - } - source = source.with_table_parquet_options(source_options); + source = source.with_table_parquet_options(self.options.clone()); // Use the CachedParquetFileReaderFactory let metadata_cache = state.runtime_env().cache_manager.get_file_metadata_cache(); diff --git a/datafusion/datasource-parquet/src/metadata.rs b/datafusion/datasource-parquet/src/metadata.rs index 18f2eaefb6b01..e4e12bae66642 100644 --- a/datafusion/datasource-parquet/src/metadata.rs +++ b/datafusion/datasource-parquet/src/metadata.rs @@ -154,9 +154,6 @@ pub struct DFParquetMetadata<'a> { pub coerce_int96_tz: Option>, /// If true, promote string/binary columns with dictionary pages to `Dictionary(Int32, ...)`. enable_rle_to_dictionary: bool, - /// When set, restricts RLE->Dictionary promotion to only these columns. - /// Takes priority over `enable_rle_to_dictionary`; `None` defers to the flag. - rle_column_allowlist: Option>, } impl<'a> DFParquetMetadata<'a> { @@ -175,7 +172,6 @@ impl<'a> DFParquetMetadata<'a> { coerce_int96: None, coerce_int96_tz: None, enable_rle_to_dictionary: false, - rle_column_allowlist: None, } } @@ -185,16 +181,6 @@ impl<'a> DFParquetMetadata<'a> { self } - /// Restrict RLE->Dictionary promotion to a named subset of columns. - /// Overrides `enable_rle_to_dictionary`; see [`crate::file_format::ParquetFormat::with_rle_column_allowlist`]. - pub(crate) fn with_rle_column_allowlist( - mut self, - columns: impl IntoIterator, - ) -> Self { - self.rle_column_allowlist = Some(columns.into_iter().collect()); - self - } - /// Set a hint for the number of trailing bytes to prefetch from the end /// of the file, equivalent to /// [`ParquetMetaDataReader::with_prefetch_hint`]. @@ -473,9 +459,7 @@ impl<'a> DFParquetMetadata<'a> { }) .unwrap_or(schema); - let should_promote = - self.enable_rle_to_dictionary || self.rle_column_allowlist.is_some(); - let schema = if should_promote { + let schema = if self.enable_rle_to_dictionary { let schema_descr = file_metadata.schema_descr(); // Top-level columns that have a dictionary page in at least one row group. let dict_cols: HashSet = metadata @@ -495,14 +479,6 @@ impl<'a> DFParquetMetadata<'a> { }) }) .collect(); - let dict_cols = if let Some(allowlist) = &self.rle_column_allowlist { - dict_cols - .into_iter() - .filter(|name| allowlist.contains(name)) - .collect::>() - } else { - dict_cols - }; if dict_cols.is_empty() { schema } else { @@ -516,10 +492,8 @@ impl<'a> DFParquetMetadata<'a> { let dict_value_type = match field.data_type() { DataType::Utf8 => Some(DataType::Utf8), DataType::LargeUtf8 => Some(DataType::LargeUtf8), - DataType::Utf8View => Some(DataType::Utf8View), DataType::Binary => Some(DataType::Binary), DataType::LargeBinary => Some(DataType::LargeBinary), - DataType::BinaryView => Some(DataType::BinaryView), _ => None, }; dict_value_type.map_or_else( diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index 807c04cc515ae..5c1716b5ef23f 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -833,25 +833,6 @@ pub fn transform_schema_to_view(schema: &Schema) -> Schema { DataType::Binary | DataType::LargeBinary => { field_with_new_type(field, DataType::BinaryView) } - // Also rewrite the value type inside dictionary columns so that - // dict<_, Utf8> / dict<_, LargeUtf8> become dict<_, Utf8View> and - // dict<_, Binary> / dict<_, LargeBinary> become dict<_, BinaryView>. - DataType::Dictionary(key_type, value_type) => { - let new_value = match value_type.as_ref() { - DataType::Utf8 | DataType::LargeUtf8 => Some(DataType::Utf8View), - DataType::Binary | DataType::LargeBinary => { - Some(DataType::BinaryView) - } - _ => None, - }; - match new_value { - Some(vt) => field_with_new_type( - field, - DataType::Dictionary(key_type.clone(), Box::new(vt)), - ), - None => Arc::clone(field), - } - } _ => Arc::clone(field), }) .collect(); @@ -867,22 +848,6 @@ pub fn transform_binary_to_string(schema: &Schema) -> Schema { DataType::Binary => field_with_new_type(field, DataType::Utf8), DataType::LargeBinary => field_with_new_type(field, DataType::LargeUtf8), DataType::BinaryView => field_with_new_type(field, DataType::Utf8View), - // Keep binary_as_string consistent for dictionary value types. - DataType::Dictionary(key_type, value_type) => match value_type.as_ref() { - DataType::Binary => field_with_new_type( - field, - DataType::Dictionary(key_type.clone(), Box::new(DataType::Utf8)), - ), - DataType::LargeBinary => field_with_new_type( - field, - DataType::Dictionary(key_type.clone(), Box::new(DataType::LargeUtf8)), - ), - DataType::BinaryView => field_with_new_type( - field, - DataType::Dictionary(key_type.clone(), Box::new(DataType::Utf8View)), - ), - _ => Arc::clone(field), - }, _ => Arc::clone(field), }) .collect(); @@ -2116,59 +2081,4 @@ mod tests { assert_eq!(result.field(0).data_type(), &DataType::Utf8); assert_eq!(result.field(1).data_type(), &DataType::Utf8); } - - #[test] - fn transform_binary_to_string_rewrites_dict_binary_value_type() { - // binary_as_string must also rewrite dictionary value types. - let schema = Schema::new(vec![ - Field::new( - "rle_binary", - DataType::Dictionary( - Box::new(DataType::Int32), - Box::new(DataType::Binary), - ), - true, - ), - Field::new( - "rle_large_binary", - DataType::Dictionary( - Box::new(DataType::Int32), - Box::new(DataType::LargeBinary), - ), - true, - ), - Field::new("plain_binary", DataType::Binary, true), - Field::new( - "rle_utf8", - DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), - true, - ), - ]); - - let result = transform_binary_to_string(&schema); - - assert_eq!( - result.field(0).data_type(), - &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), - "Dict(Int32, Binary) must become Dict(Int32, Utf8)" - ); - assert_eq!( - result.field(1).data_type(), - &DataType::Dictionary( - Box::new(DataType::Int32), - Box::new(DataType::LargeUtf8) - ), - "Dict(Int32, LargeBinary) must become Dict(Int32, LargeUtf8)" - ); - assert_eq!( - result.field(2).data_type(), - &DataType::Utf8, - "plain Binary must become Utf8" - ); - assert_eq!( - result.field(3).data_type(), - &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), - "Dict(Int32, Utf8) must be unchanged" - ); - } } diff --git a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt index 685d904f8cd90..8ff8361c58a55 100644 --- a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt +++ b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt @@ -234,9 +234,6 @@ COPY (SELECT arrow_cast(X'6865', 'Binary') AS payload FROM (VALUES (1), (2))) TO 'test_files/scratch/parquet_rle_to_dictionary/compat/binary_plain.parquet' STORED AS PARQUET OPTIONS ('format.dictionary_enabled' false); -statement ok -set datafusion.execution.parquet.enable_rle_to_dictionary = true; - # Registering incompatible Dict Utf8 and plain Binary files must fail. statement error Arrow error: Schema error: Fail to merge schema field 'payload' because the from data_type = Dictionary\(Int32, Utf8\) does not equal Binary CREATE EXTERNAL TABLE mixed_str_binary @@ -245,6 +242,3 @@ LOCATION 'test_files/scratch/parquet_rle_to_dictionary/compat/'; statement ok reset datafusion.execution.parquet.enable_rle_to_dictionary; - -statement ok -RESET datafusion.catalog.create_default_catalog_and_schema; From a8b74aa8e79eb1dfda90bffaf18da1bf87cf2173 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Wed, 2 Sep 2026 13:12:28 -0400 Subject: [PATCH 3/9] enable flag --- datafusion/common/src/config.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 179eb550fc231..0da478d141b14 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -1432,7 +1432,7 @@ config_namespace! { /// a user-supplied schema are not promoted because the Parquet footer is /// not read at DDL time, so dictionary pages cannot be detected per column. /// See - pub enable_rle_to_dictionary: bool, default = false + pub enable_rle_to_dictionary: bool, default = true // The following options affect writing to parquet files // and map to parquet::file::properties::WriterProperties From 2d26b55612bd39631e274f5c66eaaad4edffae7f Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Wed, 2 Sep 2026 13:43:05 -0400 Subject: [PATCH 4/9] flip back to false --- datafusion/common/src/config.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 0da478d141b14..179eb550fc231 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -1432,7 +1432,7 @@ config_namespace! { /// a user-supplied schema are not promoted because the Parquet footer is /// not read at DDL time, so dictionary pages cannot be detected per column. /// See - pub enable_rle_to_dictionary: bool, default = true + pub enable_rle_to_dictionary: bool, default = false // The following options affect writing to parquet files // and map to parquet::file::properties::WriterProperties From d571157be600c4c20d6d3f3b7dd5d09ba43c969c Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Wed, 2 Sep 2026 23:19:54 -0400 Subject: [PATCH 5/9] address key growth --- .../datasource-parquet/src/schema_coercion.rs | 155 ++++++++++++++++-- 1 file changed, 140 insertions(+), 15 deletions(-) diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index 5c1716b5ef23f..cf8f1257c29a6 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -419,12 +419,6 @@ fn coerce_map_entries( Some(field_with_new_type(file_entries, DataType::Struct(fields))) } -fn dictionary_value_type(data_type: &DataType) -> Option<&DataType> { - match data_type { - DataType::Dictionary(_, value_type) => Some(value_type.as_ref()), - _ => None, - } -} // Find the value type that can represent both sides without narrowing offsets // or crossing string/binary families. @@ -451,17 +445,61 @@ fn common_dictionary_value_type( } } -/// Allows safe widening into the table dictionary type. -/// - `Utf8` to `Dictionary(Int32, LargeUtf8)`: allowed -/// - `LargeBinary` to `Dictionary(Int32, Binary)`: rejected +// Same family (both signed or both unsigned): return the wider member. +// Mixed: return the smallest signed type covering the larger capacity; Int64 is the ceiling +// because Parquet only supports signed integer dictionary key types. +// Returns None only for non-integer key types. +fn common_dictionary_key_type(a: &DataType, b: &DataType) -> Option { + fn key_capacity(dt: &DataType) -> Option { + match dt { + DataType::Int8 => Some(1u128 << 7), + DataType::Int16 => Some(1u128 << 15), + DataType::Int32 => Some(1u128 << 31), + DataType::Int64 => Some(1u128 << 63), + DataType::UInt8 => Some(1u128 << 8), + DataType::UInt16 => Some(1u128 << 16), + DataType::UInt32 => Some(1u128 << 32), + DataType::UInt64 => Some(1u128 << 64), + _ => None, + } + } + fn is_signed(dt: &DataType) -> bool { + matches!( + dt, + DataType::Int8 | DataType::Int16 | DataType::Int32 | DataType::Int64 + ) + } + let cap_a = key_capacity(a)?; + let cap_b = key_capacity(b)?; + if is_signed(a) == is_signed(b) { + return Some(if cap_a >= cap_b { a.clone() } else { b.clone() }); + } + let max_cap = cap_a.max(cap_b); + Some(match max_cap { + c if c <= (1u128 << 7) => DataType::Int8, + c if c <= (1u128 << 15) => DataType::Int16, + c if c <= (1u128 << 31) => DataType::Int32, + _ => DataType::Int64, + }) +} + fn can_promote_to_dictionary_type( file_field_type: &DataType, table_dictionary_type: &DataType, ) -> bool { - dictionary_value_type(table_dictionary_type).is_some_and(|dictionary_value_type| { - common_dictionary_value_type(file_field_type, dictionary_value_type) - .is_some_and(|common_type| &common_type == dictionary_value_type) - }) + let DataType::Dictionary(table_key, table_value) = table_dictionary_type else { + return false; + }; + if common_dictionary_value_type(file_field_type, table_value) + .is_none_or(|common| &common != table_value.as_ref()) + { + return false; + } + if let DataType::Dictionary(file_key, _) = file_field_type { + return common_dictionary_key_type(file_key, table_key) + .is_some_and(|common| &common == table_key.as_ref()); + } + true } /// Normalize per-file schemas so that a column promoted to `Dictionary` in @@ -486,7 +524,6 @@ pub(crate) fn uniform_dict_schemas(schemas: Vec) -> Vec { return schemas; } - // Widen the recorded dictionary value type before promoting plain fields. for schema in &schemas { for field in schema.fields() { let Some(dict_type) = dict_types.get_mut(field.name()) else { @@ -497,10 +534,17 @@ pub(crate) fn uniform_dict_schemas(schemas: Vec) -> Vec { }; let key_type = key_type.clone(); let value_type = value_type.as_ref().clone(); + let new_key = if let DataType::Dictionary(file_key, _) = field.data_type() { + common_dictionary_key_type(&key_type, file_key) + .map(Box::new) + .unwrap_or_else(|| key_type.clone()) + } else { + key_type.clone() + }; if let Some(common_type) = common_dictionary_value_type(field.data_type(), &value_type) { - *dict_type = DataType::Dictionary(key_type, Box::new(common_type)); + *dict_type = DataType::Dictionary(new_key, Box::new(common_type)); } } } @@ -2081,4 +2125,85 @@ mod tests { assert_eq!(result.field(0).data_type(), &DataType::Utf8); assert_eq!(result.field(1).data_type(), &DataType::Utf8); } + + fn dk(key: DataType) -> DataType { + DataType::Dictionary(Box::new(key), Box::new(DataType::Utf8)) + } + + #[test] + fn uniform_dict_schemas_key_type_widening() { + // (name, file_a, file_b, expected_a, expected_b) + let cases = [ + ( + "signed wider wins", + dk(DataType::Int8), + dk(DataType::Int32), + dk(DataType::Int32), + dk(DataType::Int32), + ), + ( + "signed wider wins reversed", + dk(DataType::Int32), + dk(DataType::Int8), + dk(DataType::Int32), + dk(DataType::Int32), + ), + ( + "same unsigned stays", + dk(DataType::UInt8), + dk(DataType::UInt8), + dk(DataType::UInt8), + dk(DataType::UInt8), + ), + ( + "uint8+int32→int32", + dk(DataType::UInt8), + dk(DataType::Int32), + dk(DataType::Int32), + dk(DataType::Int32), + ), + ( + "uint8+int8→int16", + dk(DataType::UInt8), + dk(DataType::Int8), + dk(DataType::Int16), + dk(DataType::Int16), + ), + ( + "uint64+int64→int64", + dk(DataType::UInt64), + dk(DataType::Int64), + dk(DataType::Int64), + dk(DataType::Int64), + ), + ]; + for (name, a, b, exp_a, exp_b) in cases { + let result = + uniform_dict_schemas(vec![one_field_schema(a), one_field_schema(b)]); + assert_eq!(result[0].field(0).data_type(), &exp_a, "{name}"); + assert_eq!(result[1].field(0).data_type(), &exp_b, "{name}"); + } + } + + #[test] + fn rle_schema_coercion_rejects_key_narrowing() { + let coerce = |table: DataType, file: DataType| { + let t = one_field_schema(table); + let f = one_field_schema(file.clone()); + apply_file_schema_type_coercions_with_rle(&t, &f, true) + .as_ref() + .map(|s| s.field(0).data_type().clone()) + .unwrap_or(file) + }; + // Dict(Int32) file must not be narrowed to Dict(Int8) at scan time. + assert_eq!( + coerce(dk(DataType::Int8), dk(DataType::Int32)), + dk(DataType::Int32) + ); + // Plain Utf8 file is still promoted to a narrow-key dict (no file key to narrow from). + assert_eq!( + coerce(dk(DataType::Int8), DataType::Utf8), + dk(DataType::Int8) + ); + } } From 2b89440fcbf1e737a905bddbc1068bdbaf388f16 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 7 Sep 2026 00:10:32 -0400 Subject: [PATCH 6/9] add unit test to catch widening dictionary case, Add docs detailing new parquet option --- .../datasource-parquet/src/schema_coercion.rs | 150 +++++++++++++++++- .../test_files/parquet_rle_to_dictionary.slt | 90 +++++++++++ .../library-user-guide/upgrading/56.0.0.md | 36 +++++ 3 files changed, 270 insertions(+), 6 deletions(-) diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index cf8f1257c29a6..9ca334aa2c38a 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -499,7 +499,12 @@ fn can_promote_to_dictionary_type( return common_dictionary_key_type(file_key, table_key) .is_some_and(|common| &common == table_key.as_ref()); } - true + // Plain file field: only allow promotion if the table key is at least Int32 wide. + // A plain column can have any number of distinct values, so narrow keys (Int8, Int16) + // are unsafe — uniform_dict_schemas widens to Int32 when a plain file is present, so + // this guard handles the case where can_promote_to_dictionary_type is called directly. + common_dictionary_key_type(&DataType::Int32, table_key) + .is_some_and(|common| &common == table_key.as_ref()) } /// Normalize per-file schemas so that a column promoted to `Dictionary` in @@ -539,7 +544,14 @@ pub(crate) fn uniform_dict_schemas(schemas: Vec) -> Vec { .map(Box::new) .unwrap_or_else(|| key_type.clone()) } else { - key_type.clone() + // Plain field: treat as Int32 key capacity since we don't know how + // many distinct values it has. This prevents a narrow key (e.g. Int8) + // from being selected as the common type when one file uses a narrow + // dictionary and another uses a plain encoding with potentially more + // values than the key can represent. + common_dictionary_key_type(&key_type, &DataType::Int32) + .map(Box::new) + .unwrap_or_else(|| Box::new(DataType::Int32)) }; if let Some(common_type) = common_dictionary_value_type(field.data_type(), &value_type) @@ -883,7 +895,11 @@ pub fn transform_schema_to_view(schema: &Schema) -> Schema { Schema::new_with_metadata(transformed_fields, schema.metadata.clone()) } -/// Transform a schema so that any binary types are strings +/// Transform a schema so that any binary types are strings. +/// +/// Also handles `Dictionary(key, binary)` produced by `enable_rle_to_dictionary`, +/// converting the dictionary value type so `binary_as_string` applies consistently +/// regardless of whether the physical encoding triggered dictionary promotion. pub fn transform_binary_to_string(schema: &Schema) -> Schema { let transformed_fields: Vec> = schema .fields @@ -892,6 +908,21 @@ pub fn transform_binary_to_string(schema: &Schema) -> Schema { DataType::Binary => field_with_new_type(field, DataType::Utf8), DataType::LargeBinary => field_with_new_type(field, DataType::LargeUtf8), DataType::BinaryView => field_with_new_type(field, DataType::Utf8View), + DataType::Dictionary(key_type, value_type) => match value_type.as_ref() { + DataType::Binary => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(DataType::Utf8)), + ), + DataType::LargeBinary => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(DataType::LargeUtf8)), + ), + DataType::BinaryView => field_with_new_type( + field, + DataType::Dictionary(key_type.clone(), Box::new(DataType::Utf8View)), + ), + _ => Arc::clone(field), + }, _ => Arc::clone(field), }) .collect(); @@ -2185,6 +2216,51 @@ mod tests { } } + #[test] + fn uniform_dict_schemas_plain_field_widens_key_to_int32() { + // A plain field mixed with a narrow-key dictionary must widen the key to at + // least Int32 so that the plain file (with unknown cardinality) can be + // represented safely. + let cases = [ + ( + "int8_dict + plain → int32", + DataType::Dictionary(Box::new(DataType::Int8), Box::new(DataType::Utf8)), + DataType::Utf8, + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + ), + ( + "int16_dict + plain → int32", + DataType::Dictionary(Box::new(DataType::Int16), Box::new(DataType::Utf8)), + DataType::Utf8, + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + ), + ( + "int32_dict + plain stays int32", + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + DataType::Utf8, + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + ), + ( + "int64_dict + plain stays int64", + DataType::Dictionary(Box::new(DataType::Int64), Box::new(DataType::Utf8)), + DataType::Utf8, + DataType::Dictionary(Box::new(DataType::Int64), Box::new(DataType::Utf8)), + DataType::Dictionary(Box::new(DataType::Int64), Box::new(DataType::Utf8)), + ), + ]; + for (name, dict_type, plain_type, exp_dict, exp_plain) in cases { + let result = uniform_dict_schemas(vec![ + one_field_schema(dict_type), + one_field_schema(plain_type), + ]); + assert_eq!(result[0].field(0).data_type(), &exp_dict, "{name}"); + assert_eq!(result[1].field(0).data_type(), &exp_plain, "{name}"); + } + } + #[test] fn rle_schema_coercion_rejects_key_narrowing() { let coerce = |table: DataType, file: DataType| { @@ -2200,10 +2276,72 @@ mod tests { coerce(dk(DataType::Int8), dk(DataType::Int32)), dk(DataType::Int32) ); - // Plain Utf8 file is still promoted to a narrow-key dict (no file key to narrow from). + // Plain Utf8 file must NOT be promoted to a narrow-key dict: we cannot know how + // many distinct values the file has, so Int8 capacity is unsafe. + assert_eq!(coerce(dk(DataType::Int8), DataType::Utf8), DataType::Utf8); + } + + #[test] + fn transform_binary_to_string_handles_dictionary_value_types() { + let schema = Schema::new(vec![ + Field::new( + "a", + DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::Binary), + ), + true, + ), + Field::new( + "b", + DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::LargeBinary), + ), + true, + ), + Field::new( + "c", + DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::BinaryView), + ), + true, + ), + Field::new( + "d", + DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + true, + ), + Field::new("e", DataType::Binary, true), + ]); + + let result = transform_binary_to_string(&schema); + + assert_eq!( + result.field(0).data_type(), + &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), + ); + assert_eq!( + result.field(1).data_type(), + &DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::LargeUtf8) + ), + ); + assert_eq!( + result.field(2).data_type(), + &DataType::Dictionary( + Box::new(DataType::Int32), + Box::new(DataType::Utf8View) + ), + ); + // Dictionary(_, Utf8) is left unchanged assert_eq!( - coerce(dk(DataType::Int8), DataType::Utf8), - dk(DataType::Int8) + result.field(3).data_type(), + &DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)), ); + // Plain Binary is converted as before + assert_eq!(result.field(4).data_type(), &DataType::Utf8); } } diff --git a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt index 8ff8361c58a55..188266c57e1ad 100644 --- a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt +++ b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt @@ -195,6 +195,61 @@ Dictionary(Int32, Utf8) 4 statement ok DROP TABLE mixed_encoding; +# Mixed encoding with high cardinality: the plain file contains more than 128 +# distinct values so a narrow (Int8) dictionary key would overflow. After the +# fix, uniform_dict_schemas widens the shared key to at least Int32. +# Both files must be scannable and return the correct row count. + +statement ok +COPY ( + SELECT v AS category + FROM ( + VALUES + ('v000'),('v001'),('v002'),('v003'),('v004'),('v005'),('v006'),('v007'), + ('v008'),('v009'),('v010'),('v011'),('v012'),('v013'),('v014'),('v015'), + ('v016'),('v017'),('v018'),('v019'),('v020'),('v021'),('v022'),('v023'), + ('v024'),('v025'),('v026'),('v027'),('v028'),('v029'),('v030'),('v031'), + ('v032'),('v033'),('v034'),('v035'),('v036'),('v037'),('v038'),('v039'), + ('v040'),('v041'),('v042'),('v043'),('v044'),('v045'),('v046'),('v047'), + ('v048'),('v049'),('v050'),('v051'),('v052'),('v053'),('v054'),('v055'), + ('v056'),('v057'),('v058'),('v059'),('v060'),('v061'),('v062'),('v063'), + ('v064'),('v065'),('v066'),('v067'),('v068'),('v069'),('v070'),('v071'), + ('v072'),('v073'),('v074'),('v075'),('v076'),('v077'),('v078'),('v079'), + ('v080'),('v081'),('v082'),('v083'),('v084'),('v085'),('v086'),('v087'), + ('v088'),('v089'),('v090'),('v091'),('v092'),('v093'),('v094'),('v095'), + ('v096'),('v097'),('v098'),('v099'),('v100'),('v101'),('v102'),('v103'), + ('v104'),('v105'),('v106'),('v107'),('v108'),('v109'),('v110'),('v111'), + ('v112'),('v113'),('v114'),('v115'),('v116'),('v117'),('v118'),('v119'), + ('v120'),('v121'),('v122'),('v123'),('v124'),('v125'),('v126'),('v127'), + ('v128'),('v129') + ) AS t(v) +) +TO 'test_files/scratch/parquet_rle_to_dictionary/high_card/plain.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' false); + +statement ok +COPY (SELECT column1 AS category FROM (VALUES ('rle_a'), ('rle_b'), ('rle_a'))) +TO 'test_files/scratch/parquet_rle_to_dictionary/high_card/rle.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' true); + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = true; + +statement ok +CREATE EXTERNAL TABLE high_card_mixed +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/high_card/'; + +# Both files are scanned as Dictionary(Int32, Utf8) — Int32 key avoids +# overflow for the 130 distinct values in the plain file. +query TI +SELECT arrow_typeof(category), COUNT(*) FROM high_card_mixed GROUP BY arrow_typeof(category); +---- +Dictionary(Int32, Utf8) 133 + +statement ok +DROP TABLE high_card_mixed; + # Explicit-schema tables are not promoted. statement ok @@ -240,5 +295,40 @@ CREATE EXTERNAL TABLE mixed_str_binary STORED AS PARQUET LOCATION 'test_files/scratch/parquet_rle_to_dictionary/compat/'; +# binary_as_string + enable_rle_to_dictionary: Dictionary(_, Binary) value type +# must be promoted to Utf8 so the binary_as_string contract holds regardless of +# physical encoding. + +query I +COPY (SELECT arrow_cast(X'6865', 'Binary') AS payload FROM (VALUES (1), (2))) +TO 'test_files/scratch/parquet_rle_to_dictionary/binary_rle.parquet' +STORED AS PARQUET OPTIONS ('format.dictionary_enabled' true); +---- +2 + +statement ok +set datafusion.execution.parquet.enable_rle_to_dictionary = true; + +statement ok +set datafusion.execution.parquet.binary_as_string = true; + +statement ok +CREATE EXTERNAL TABLE binary_rle_as_string +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/binary_rle.parquet'; + +# Dictionary(Int32, Binary) must be coerced to Dictionary(Int32, Utf8) when +# binary_as_string is true, not left as Dictionary(Int32, Binary). +query T +SELECT DISTINCT arrow_typeof(payload) FROM binary_rle_as_string; +---- +Dictionary(Int32, Utf8) + +statement ok +DROP TABLE binary_rle_as_string; + +statement ok +reset datafusion.execution.parquet.binary_as_string; + statement ok reset datafusion.execution.parquet.enable_rle_to_dictionary; diff --git a/docs/source/library-user-guide/upgrading/56.0.0.md b/docs/source/library-user-guide/upgrading/56.0.0.md index 312668c485629..e981af06707fd 100644 --- a/docs/source/library-user-guide/upgrading/56.0.0.md +++ b/docs/source/library-user-guide/upgrading/56.0.0.md @@ -298,6 +298,42 @@ The same setting can be changed with SQL: SET datafusion.execution.soft_max_bytes_per_output_file = 134217728; ``` +### New field `enable_rle_to_dictionary` added to `ParquetOptions` + +`ParquetOptions` gained a new field `enable_rle_to_dictionary: bool` (default +`false`). It controls whether top-level string and binary Parquet columns with +dictionary pages are read as `Dictionary` / `Dictionary` instead of their plain type. + +`ParquetOptions` is a public struct without `#[non_exhaustive]`, so downstream +crates that construct it with an exhaustive struct literal will fail to compile +with `E0063: missing field 'enable_rle_to_dictionary'`. + +**Migration guide:** + +Add `enable_rle_to_dictionary: false` to any exhaustive `ParquetOptions { .. }` +literal, or use `..Default::default()` to future-proof against further additions: + +```rust,ignore +// Before +ParquetOptions { + enable_page_index: true, + // ... every other field ... +} + +// After: set it explicitly +ParquetOptions { + enable_page_index: true, + enable_rle_to_dictionary: false, + // ... every other field ... +} + +// After: or let remaining fields come from Default +ParquetOptions { + enable_page_index: true, + ..Default::default() +} +``` + ### `datafusion.optimizer.use_statistics_registry` is deprecated and ignored The `datafusion.optimizer.use_statistics_registry` config flag is deprecated and From bb984f7cb66e56609883fd3bb74fcee446faeeff Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Tue, 8 Sep 2026 17:13:13 -0400 Subject: [PATCH 7/9] address review: fix UInt64+Int64 key narrowing, exercise Int8 overflow in SLT - common_dictionary_key_type: add Int64 arm before catch-all so UInt64+Int64 selects UInt64 (wider) rather than silently narrowing to Int64 - update test expectation from Int64 to UInt64 for that case - SLT high-cardinality test: write rle.parquet via arrow_cast to Dictionary(Int8, Utf8) so the file embeds an Int8-keyed Arrow schema; add a standalone check confirming the round-trip; this exercises the actual Dict(Int8)+plain-130-values overflow path that the previous plain-Utf8 write never triggered --- .../datasource-parquet/src/schema_coercion.rs | 13 ++++----- .../test_files/parquet_rle_to_dictionary.slt | 27 ++++++++++++++++--- 2 files changed, 31 insertions(+), 9 deletions(-) diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index 9ca334aa2c38a..481c49b7bc1b0 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -446,8 +446,8 @@ fn common_dictionary_value_type( } // Same family (both signed or both unsigned): return the wider member. -// Mixed: return the smallest signed type covering the larger capacity; Int64 is the ceiling -// because Parquet only supports signed integer dictionary key types. +// Mixed: return the type with the larger capacity. For the UInt64+Int64 case +// this is UInt64 since arrow-rs supports both signed and unsigned dictionary keys. // Returns None only for non-integer key types. fn common_dictionary_key_type(a: &DataType, b: &DataType) -> Option { fn key_capacity(dt: &DataType) -> Option { @@ -479,7 +479,8 @@ fn common_dictionary_key_type(a: &DataType, b: &DataType) -> Option { c if c <= (1u128 << 7) => DataType::Int8, c if c <= (1u128 << 15) => DataType::Int16, c if c <= (1u128 << 31) => DataType::Int32, - _ => DataType::Int64, + c if c <= (1u128 << 63) => DataType::Int64, + _ => DataType::UInt64, }) } @@ -2201,11 +2202,11 @@ mod tests { dk(DataType::Int16), ), ( - "uint64+int64→int64", + "uint64+int64→uint64", dk(DataType::UInt64), dk(DataType::Int64), - dk(DataType::Int64), - dk(DataType::Int64), + dk(DataType::UInt64), + dk(DataType::UInt64), ), ]; for (name, a, b, exp_a, exp_b) in cases { diff --git a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt index 188266c57e1ad..a303f5eaec5dc 100644 --- a/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt +++ b/datafusion/sqllogictest/test_files/parquet_rle_to_dictionary.slt @@ -227,21 +227,42 @@ COPY ( TO 'test_files/scratch/parquet_rle_to_dictionary/high_card/plain.parquet' STORED AS PARQUET OPTIONS ('format.dictionary_enabled' false); +# Write rle.parquet with an explicit Dictionary(Int8, Utf8) Arrow schema so +# that the file carries Int8 dictionary keys in its embedded Arrow schema +# metadata. This exercises the real regression path: without the key-widening +# fix, uniform_dict_schemas would select Int8 as the merged key type, which +# cannot represent the 130 distinct values in the plain file. statement ok -COPY (SELECT column1 AS category FROM (VALUES ('rle_a'), ('rle_b'), ('rle_a'))) +COPY (SELECT arrow_cast(column1, 'Dictionary(Int8, Utf8)') AS category FROM (VALUES ('rle_a'), ('rle_b'), ('rle_a'))) TO 'test_files/scratch/parquet_rle_to_dictionary/high_card/rle.parquet' STORED AS PARQUET OPTIONS ('format.dictionary_enabled' true); statement ok set datafusion.execution.parquet.enable_rle_to_dictionary = true; +# Verify that rle.parquet alone reads back as Dictionary(Int8, Utf8) — +# confirming the Arrow schema metadata round-trips and was not widened on write. +statement ok +CREATE EXTERNAL TABLE rle_int8_only +STORED AS PARQUET +LOCATION 'test_files/scratch/parquet_rle_to_dictionary/high_card/rle.parquet'; + +query T +SELECT DISTINCT arrow_typeof(category) FROM rle_int8_only; +---- +Dictionary(Int8, Utf8) + +statement ok +DROP TABLE rle_int8_only; + statement ok CREATE EXTERNAL TABLE high_card_mixed STORED AS PARQUET LOCATION 'test_files/scratch/parquet_rle_to_dictionary/high_card/'; -# Both files are scanned as Dictionary(Int32, Utf8) — Int32 key avoids -# overflow for the 130 distinct values in the plain file. +# Both files are scanned as Dictionary(Int32, Utf8) — uniform_dict_schemas +# widens the Int8 key from rle.parquet to Int32 to avoid overflow for the +# 130 distinct values in the plain file. query TI SELECT arrow_typeof(category), COUNT(*) FROM high_card_mixed GROUP BY arrow_typeof(category); ---- From 50691d8ed7ba89c5f572f2de7770871067f1299a Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Sun, 27 Sep 2026 23:45:05 +0200 Subject: [PATCH 8/9] fix: remove extra blank line in schema_coercion.rs --- datafusion/datasource-parquet/src/schema_coercion.rs | 1 - 1 file changed, 1 deletion(-) diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index 481c49b7bc1b0..277a87bfb5974 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -419,7 +419,6 @@ fn coerce_map_entries( Some(field_with_new_type(file_entries, DataType::Struct(fields))) } - // Find the value type that can represent both sides without narrowing offsets // or crossing string/binary families. fn common_dictionary_value_type( From d5144977f65b7b2024f0fa5e955cebd6ee264583 Mon Sep 17 00:00:00 2001 From: RIchard Baah Date: Mon, 28 Sep 2026 09:08:23 +0200 Subject: [PATCH 9/9] refactor: use existing coercion helper pipeline in apply_file_schema_type_coercions_with_rle Replace the duplicate full implementation with a call to the existing coerce_fields_by_name helper, extended with an enable_rle flag threaded through coerce_data_type and coerce_child. This restores correct positional map-entry coercion via coerce_map_entries and avoids allocating a Vec for every field when nothing changes. --- .../datasource-parquet/src/schema_coercion.rs | 201 +++--------------- 1 file changed, 29 insertions(+), 172 deletions(-) diff --git a/datafusion/datasource-parquet/src/schema_coercion.rs b/datafusion/datasource-parquet/src/schema_coercion.rs index 277a87bfb5974..fbfc37a33f31b 100644 --- a/datafusion/datasource-parquet/src/schema_coercion.rs +++ b/datafusion/datasource-parquet/src/schema_coercion.rs @@ -77,7 +77,7 @@ pub fn apply_file_schema_type_coercions( file_schema: &Schema, ) -> Option { let fields = - coerce_fields_by_name(table_schema.fields(), file_schema.fields(), true)?; + coerce_fields_by_name(table_schema.fields(), file_schema.fields(), true, false)?; Some(Schema::new_with_metadata( fields, file_schema.metadata.clone(), @@ -92,169 +92,12 @@ pub(crate) fn apply_file_schema_type_coercions_with_rle( file_schema: &Schema, enable_rle_to_dictionary: bool, ) -> Option { - let mut needs_view_transform = false; - let mut needs_string_transform = false; - let mut needs_nested_transform = false; - let mut needs_dict_transform = false; - - // Create a mapping of table field names to their data types for fast lookup - // and simultaneously check if we need any transformations - let table_fields: HashMap<_, _> = table_schema - .fields() - .iter() - .map(|field| { - let data_type = field.data_type(); - // Check if we need view type transformation - if matches!(data_type, &DataType::Utf8View | &DataType::BinaryView) { - needs_view_transform = true; - } - // Check if we need string type transformation - if matches!( - data_type, - &DataType::Utf8 | &DataType::LargeUtf8 | &DataType::Utf8View - ) { - needs_string_transform = true; - } - // Nested fields can need transformations even when their parent does not. - if matches!( - data_type, - DataType::Struct(_) - | DataType::List(_) - | DataType::LargeList(_) - | DataType::ListView(_) - | DataType::LargeListView(_) - | DataType::FixedSizeList(_, _) - | DataType::Map(_, _) - ) { - needs_nested_transform = true; - } - if enable_rle_to_dictionary - && matches!(data_type, &DataType::Dictionary(_, _)) - { - needs_dict_transform = true; - } - - (field.name(), data_type) - }) - .collect(); - - // Early return if no transformation needed - if !needs_view_transform - && !needs_string_transform - && !needs_nested_transform - && !needs_dict_transform - { - return None; - } - - let fields: Vec> = file_schema - .fields() - .iter() - .map(|field| { - let field_name = field.name(); - let field_type = field.data_type(); - - // Look up the corresponding field type in the table schema - if let Some(table_type) = table_fields.get(field_name) { - match (table_type, field_type) { - // table schema uses string type, coerce the file schema to use string type - ( - &DataType::Utf8, - DataType::Binary | DataType::LargeBinary | DataType::BinaryView, - ) => { - return field_with_new_type(field, DataType::Utf8); - } - // table schema uses large string type, coerce the file schema to use large string type - ( - &DataType::LargeUtf8, - DataType::Binary | DataType::LargeBinary | DataType::BinaryView, - ) => { - return field_with_new_type(field, DataType::LargeUtf8); - } - // table schema uses string view type, coerce the file schema to use view type - ( - &DataType::Utf8View, - DataType::Binary | DataType::LargeBinary | DataType::BinaryView, - ) => { - return field_with_new_type(field, DataType::Utf8View); - } - // Handle view type conversions - (&DataType::Utf8View, DataType::Utf8 | DataType::LargeUtf8) => { - return field_with_new_type(field, DataType::Utf8View); - } - (&DataType::BinaryView, DataType::Binary | DataType::LargeBinary) => { - return field_with_new_type(field, DataType::BinaryView); - } - // Apply the same coercions to matching fields inside structs. - (DataType::Struct(table_fields), DataType::Struct(file_fields)) => { - if let Some(schema) = apply_file_schema_type_coercions( - &Schema::new(table_fields.clone()), - &Schema::new(file_fields.clone()), - ) { - return field_with_new_type( - field, - DataType::Struct(schema.fields), - ); - } - } - // Container children match by position, regardless of their names. - (DataType::List(table_child), DataType::List(file_child)) - | ( - DataType::LargeList(table_child), - DataType::LargeList(file_child), - ) - | (DataType::ListView(table_child), DataType::ListView(file_child)) - | ( - DataType::LargeListView(table_child), - DataType::LargeListView(file_child), - ) - | ( - DataType::FixedSizeList(table_child, _), - DataType::FixedSizeList(file_child, _), - ) - | (DataType::Map(table_child, _), DataType::Map(file_child, _)) => { - if let Some(schema) = apply_file_schema_type_coercions( - &Schema::new(vec![field_with_new_type( - file_child, - table_child.data_type().clone(), - )]), - &Schema::new(vec![Arc::clone(file_child)]), - ) { - let child = Arc::clone(&schema.fields()[0]); - let new_type = match field_type { - DataType::List(_) => DataType::List(child), - DataType::LargeList(_) => DataType::LargeList(child), - DataType::ListView(_) => DataType::ListView(child), - DataType::LargeListView(_) => { - DataType::LargeListView(child) - } - DataType::FixedSizeList(_, size) => { - DataType::FixedSizeList(child, *size) - } - DataType::Map(_, sorted) => DataType::Map(child, *sorted), - _ => return Arc::clone(field), - }; - return field_with_new_type(field, new_type); - } - } - (DataType::Dictionary(_, _), _) - if enable_rle_to_dictionary - && can_promote_to_dictionary_type(field_type, table_type) => - { - return field_with_new_type(field, (*table_type).clone()); - } - _ => {} - } - } - // If no transformation is needed, keep the original field - Arc::clone(field) - }) - .collect(); - - if fields.iter().eq(file_schema.fields().iter()) { - return None; - } - + let fields = coerce_fields_by_name( + table_schema.fields(), + file_schema.fields(), + true, + enable_rle_to_dictionary, + )?; Some(Schema::new_with_metadata( fields, file_schema.metadata.clone(), @@ -270,6 +113,7 @@ fn coerce_fields_by_name( table_fields: &Fields, file_fields: &Fields, binary_to_string: bool, + enable_rle: bool, ) -> Option { // Create a mapping of table field names to their data types for fast lookup let table_types: HashMap<_, _> = table_fields @@ -279,7 +123,7 @@ fn coerce_fields_by_name( coerce_fields(file_fields, |_, field| { let table_type = table_types.get(field.name())?; - coerce_data_type(table_type, field.data_type(), binary_to_string) + coerce_data_type(table_type, field.data_type(), binary_to_string, enable_rle) .map(|new_type| field_with_new_type(field, new_type)) }) } @@ -330,6 +174,7 @@ fn coerce_data_type( table_type: &DataType, file_type: &DataType, binary_to_string: bool, + enable_rle: bool, ) -> Option { use DataType::*; match (table_type, file_type) { @@ -347,24 +192,28 @@ fn coerce_data_type( (BinaryView, Binary | LargeBinary) => Some(BinaryView), // Struct children match by name (Struct(table_fields), Struct(file_fields)) => { - coerce_fields_by_name(table_fields, file_fields, binary_to_string).map(Struct) + coerce_fields_by_name(table_fields, file_fields, binary_to_string, enable_rle) + .map(Struct) } // List-like children match by position, regardless of their names. // The container kind and FixedSizeList width always come from the file. (List(table_child), List(file_child)) => { - coerce_child(table_child, file_child, binary_to_string).map(List) + coerce_child(table_child, file_child, binary_to_string, enable_rle).map(List) } (LargeList(table_child), LargeList(file_child)) => { - coerce_child(table_child, file_child, binary_to_string).map(LargeList) + coerce_child(table_child, file_child, binary_to_string, enable_rle) + .map(LargeList) } (ListView(table_child), ListView(file_child)) => { - coerce_child(table_child, file_child, binary_to_string).map(ListView) + coerce_child(table_child, file_child, binary_to_string, enable_rle) + .map(ListView) } (LargeListView(table_child), LargeListView(file_child)) => { - coerce_child(table_child, file_child, binary_to_string).map(LargeListView) + coerce_child(table_child, file_child, binary_to_string, enable_rle) + .map(LargeListView) } (FixedSizeList(table_child, _), FixedSizeList(file_child, size)) => { - coerce_child(table_child, file_child, binary_to_string) + coerce_child(table_child, file_child, binary_to_string, enable_rle) .map(|child| FixedSizeList(child, *size)) } // Map keys and values match by position: Parquet always names them @@ -373,6 +222,12 @@ fn coerce_data_type( coerce_map_entries(table_entries, file_entries) .map(|entries| Map(entries, *sorted)) } + // Promote a plain file column to the dictionary type the table expects. + (Dictionary(_, _), _) + if enable_rle && can_promote_to_dictionary_type(file_type, table_type) => + { + Some(table_type.clone()) + } _ => None, } } @@ -383,11 +238,13 @@ fn coerce_child( table_child: &FieldRef, file_child: &FieldRef, binary_to_string: bool, + enable_rle: bool, ) -> Option { coerce_data_type( table_child.data_type(), file_child.data_type(), binary_to_string, + enable_rle, ) .map(|new_type| field_with_new_type(file_child, new_type)) } @@ -413,7 +270,7 @@ fn coerce_map_entries( } let fields = coerce_fields(file_fields, |idx, file_child| { - coerce_child(&table_fields[idx], file_child, false) + coerce_child(&table_fields[idx], file_child, false, false) })?; Some(field_with_new_type(file_entries, DataType::Struct(fields)))