Skip to content

Commit bc5011f

Browse files
committed
feat: introduce optional RLE-dictionary preservation for Parquet reads
Adds `datafusion.execution.parquet.enable_rle_to_dictionary` (default false). When enabled, string and binary columns that are physically RLE_DICTIONARY-encoded in the Parquet file footer are surfaced as Arrow `Dictionary(Int32, Utf8/Binary)` arrays instead of being decoded to plain `Utf8/Binary`. This preserves the compact encoding in memory and can improve aggregation performance on low-cardinality columns. Schema promotion happens at planning time in `fetch_schema` so that physical operators (`FilterExec`, `AggregateExec`, etc.) are compiled against the correct output types before any file I/O begins. At scan time the same promotion is applied via `ArrowReaderOptions::with_schema` so arrow-rs produces dictionary arrays directly. `apply_file_schema_type_coercions` keeps its original two-parameter signature (no breaking change); the dictionary-aware variant is exposed as the crate-internal `apply_file_schema_type_coercions_with_rle`.
1 parent 1164d60 commit bc5011f

19 files changed

Lines changed: 1195 additions & 13 deletions

File tree

‎datafusion/common/src/config.rs‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1293,6 +1293,19 @@ config_namespace! {
12931293
/// Defaults to 20.
12941294
pub max_in_list_size: usize, default = 20
12951295

1296+
/// (reading) If true, string and binary columns that are dictionary-encoded
1297+
/// (RLE_DICTIONARY) in the Parquet file are read as
1298+
/// `Dictionary<Int32, Utf8>` / `Dictionary<Int32, Binary>` instead of
1299+
/// their plain value type. This reduces memory usage for low-cardinality
1300+
/// columns and can improve aggregation performance.
1301+
///
1302+
/// This applies only to **inferred-schema** tables (i.e. `CREATE EXTERNAL
1303+
/// TABLE` without an explicit column list). Tables with a user-supplied
1304+
/// schema are not promoted because the Parquet footer is not read at DDL
1305+
/// time, so RLE encoding cannot be detected per-column.
1306+
/// See <https://github.com/apache/datafusion/issues/24112>
1307+
pub enable_rle_to_dictionary: bool, default = false
1308+
12961309
// The following options affect writing to parquet files
12971310
// and map to parquet::file::properties::WriterProperties
12981311

‎datafusion/common/src/file_options/parquet_writer.rs‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -249,6 +249,7 @@ impl ParquetOptions {
249249
skip_arrow_metadata: _,
250250
max_predicate_cache_size: _,
251251
max_in_list_size: _,
252+
enable_rle_to_dictionary: _,
252253
} = self;
253254

254255
let mut builder = WriterProperties::builder()
@@ -509,6 +510,7 @@ mod tests {
509510
coerce_int96_tz: None,
510511
max_predicate_cache_size: defaults.max_predicate_cache_size,
511512
content_defined_chunking: defaults.content_defined_chunking.clone(),
513+
enable_rle_to_dictionary: defaults.enable_rle_to_dictionary,
512514
}
513515
}
514516

@@ -631,6 +633,8 @@ mod tests {
631633
coerce_int96: None,
632634
coerce_int96_tz: None,
633635
content_defined_chunking: props.content_defined_chunking().into(),
636+
enable_rle_to_dictionary: global_options_defaults
637+
.enable_rle_to_dictionary,
634638
},
635639
column_specific_options,
636640
key_value_metadata,

‎datafusion/core/tests/parquet/schema_coercion.rs‎

Lines changed: 151 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -24,9 +24,10 @@ use arrow::array::{
2424
use arrow::datatypes::{DataType, Field, Schema};
2525
use datafusion::datasource::physical_plan::ParquetSource;
2626
use datafusion::physical_plan::collect;
27-
use datafusion::prelude::SessionContext;
27+
use datafusion::prelude::{ParquetReadOptions, SessionConfig, SessionContext};
2828
use datafusion::test::object_store::local_unpartitioned_file;
2929
use datafusion_common::Result;
30+
use datafusion_common::ScalarValue;
3031
use datafusion_common::test_util::batches_to_sort_string;
3132
use datafusion_execution::object_store::ObjectStoreUrl;
3233

@@ -145,6 +146,154 @@ async fn multi_parquet_coercion_projection() {
145146
");
146147
}
147148

149+
/// Writes `batch` to a temp parquet file where `rle_cols` use RLE_DICTIONARY encoding and all other columns use plain.
150+
fn store_mixed_encoding_parquet(
151+
batch: &RecordBatch,
152+
rle_cols: &[&str],
153+
) -> (ObjectMeta, NamedTempFile) {
154+
use parquet::schema::types::ColumnPath;
155+
let mut output = tempfile::Builder::new()
156+
.suffix(".parquet")
157+
.tempfile()
158+
.expect("creating temp file");
159+
let mut builder = WriterProperties::builder().set_dictionary_enabled(false);
160+
for col in rle_cols {
161+
builder = builder
162+
.set_column_dictionary_enabled(ColumnPath::new(vec![col.to_string()]), true);
163+
}
164+
let mut writer =
165+
ArrowWriter::try_new(&mut output, batch.schema(), Some(builder.build()))
166+
.expect("creating writer");
167+
writer.write(batch).expect("Writing batch");
168+
writer.close().unwrap();
169+
let meta = local_unpartitioned_file(&output);
170+
(meta, output)
171+
}
172+
173+
/// Promotion happens at logical planning time: after `register_parquet` with the
174+
/// flag on, `ctx.table()` already shows RLE columns as `Dictionary(Int32, Utf8)`
175+
/// before any query executes.
176+
#[tokio::test]
177+
async fn rle_dictionary_logical_schema_promotion() {
178+
let schema = Arc::new(Schema::new(vec![
179+
Field::new("env", DataType::Utf8, true),
180+
Field::new("version", DataType::Utf8, true),
181+
]));
182+
let batch = RecordBatch::try_new(
183+
Arc::clone(&schema),
184+
vec![
185+
Arc::new(StringArray::from(vec!["prod", "staging"])) as ArrayRef,
186+
Arc::new(StringArray::from(vec!["1.0", "1.1"])) as ArrayRef,
187+
],
188+
)
189+
.unwrap();
190+
let (_meta, file) = store_mixed_encoding_parquet(&batch, &["env"]);
191+
192+
let config = SessionConfig::new().set(
193+
"datafusion.execution.parquet.enable_rle_to_dictionary",
194+
&ScalarValue::Boolean(Some(true)),
195+
);
196+
let ctx = SessionContext::new_with_config(config);
197+
ctx.register_parquet(
198+
"t",
199+
file.path().to_str().unwrap(),
200+
ParquetReadOptions::default(),
201+
)
202+
.await
203+
.unwrap();
204+
205+
let logical_schema = ctx.table("t").await.unwrap().schema().clone();
206+
let dict_type =
207+
DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8));
208+
assert_eq!(
209+
logical_schema
210+
.field_with_unqualified_name("env")
211+
.unwrap()
212+
.data_type(),
213+
&dict_type,
214+
);
215+
assert_eq!(
216+
logical_schema
217+
.field_with_unqualified_name("version")
218+
.unwrap()
219+
.data_type(),
220+
&DataType::Utf8View,
221+
);
222+
}
223+
224+
/// Five string columns, three RLE_DICTIONARY encoded and two plain.
225+
/// Flag on: only the three RLE columns become `Dictionary(Int32, Utf8)`; plain
226+
/// columns stay `Utf8View`. Flag off: all columns stay `Utf8View`.
227+
#[tokio::test]
228+
async fn rle_dictionary_selective_promotion() {
229+
let schema = Arc::new(Schema::new(vec![
230+
Field::new("env", DataType::Utf8, true),
231+
Field::new("region", DataType::Utf8, true),
232+
Field::new("tier", DataType::Utf8, true),
233+
Field::new("version", DataType::Utf8, true),
234+
Field::new("build", DataType::Utf8, true),
235+
]));
236+
let make_col = |vals: Vec<&str>| -> ArrayRef { Arc::new(StringArray::from(vals)) };
237+
let batch = RecordBatch::try_new(
238+
Arc::clone(&schema),
239+
vec![
240+
make_col(vec!["prod", "staging", "prod", "canary", "staging", "prod"]),
241+
make_col(vec![
242+
"us-east", "us-west", "us-east", "eu", "us-west", "us-east",
243+
]),
244+
make_col(vec!["free", "pro", "free", "enterprise", "pro", "free"]),
245+
make_col(vec!["1.0", "1.1", "1.2", "1.0", "1.1", "1.2"]),
246+
make_col(vec![
247+
"debug", "release", "debug", "release", "debug", "release",
248+
]),
249+
],
250+
)
251+
.unwrap();
252+
253+
let (_meta, file) = store_mixed_encoding_parquet(&batch, &["env", "region", "tier"]);
254+
let path = file.path().to_str().unwrap();
255+
let sql = "SELECT env, region, tier, version, build FROM t LIMIT 1";
256+
let dict_type =
257+
DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8));
258+
259+
// Flag off: all string columns stay Utf8View.
260+
let ctx = SessionContext::new();
261+
ctx.register_parquet("t", path, ParquetReadOptions::default())
262+
.await
263+
.unwrap();
264+
let result = ctx.sql(sql).await.unwrap().collect().await.unwrap();
265+
assert!(!result.is_empty());
266+
let s = result[0].schema();
267+
for col in ["env", "region", "tier", "version", "build"] {
268+
assert_eq!(
269+
s.field_with_name(col).unwrap().data_type(),
270+
&DataType::Utf8View
271+
);
272+
}
273+
274+
// Flag on: only RLE-encoded columns become Dict; plain columns stay Utf8View.
275+
let config = SessionConfig::new().set(
276+
"datafusion.execution.parquet.enable_rle_to_dictionary",
277+
&ScalarValue::Boolean(Some(true)),
278+
);
279+
let ctx = SessionContext::new_with_config(config);
280+
ctx.register_parquet("t", path, ParquetReadOptions::default())
281+
.await
282+
.unwrap();
283+
let result = ctx.sql(sql).await.unwrap().collect().await.unwrap();
284+
assert!(!result.is_empty());
285+
let s = result[0].schema();
286+
for col in ["env", "region", "tier"] {
287+
assert_eq!(s.field_with_name(col).unwrap().data_type(), &dict_type);
288+
}
289+
for col in ["version", "build"] {
290+
assert_eq!(
291+
s.field_with_name(col).unwrap().data_type(),
292+
&DataType::Utf8View
293+
);
294+
}
295+
}
296+
148297
/// Writes `batches` to a temporary parquet file
149298
pub fn store_parquet(
150299
batches: Vec<RecordBatch>,
@@ -154,14 +303,10 @@ pub fn store_parquet(
154303
.into_iter()
155304
.map(|batch| {
156305
let mut output = NamedTempFile::new().expect("creating temp file");
157-
158-
let builder = WriterProperties::builder();
159-
let props = builder.build();
160-
306+
let props = WriterProperties::builder().build();
161307
let mut writer =
162308
ArrowWriter::try_new(&mut output, batch.schema(), Some(props))
163309
.expect("creating writer");
164-
165310
writer.write(&batch).expect("Writing batch");
166311
writer.close().unwrap();
167312
output

‎datafusion/datasource-parquet/src/file_format.rs‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -363,6 +363,9 @@ impl FileFormat for ParquetFormat {
363363
.with_file_metadata_cache(Some(Arc::clone(&file_metadata_cache)))
364364
.with_coerce_int96(coerce_int96)
365365
.with_coerce_int96_tz(coerce_int96_tz.clone())
366+
.with_enable_rle_to_dictionary(
367+
self.options.global.enable_rle_to_dictionary,
368+
)
366369
.fetch_schema_with_location()
367370
.await?;
368371
Ok::<_, DataFusionError>(result)
@@ -388,7 +391,14 @@ impl FileFormat for ParquetFormat {
388391
schemas
389392
.sort_unstable_by(|(location1, _), (location2, _)| location1.cmp(location2));
390393

391-
let schemas = schemas.into_iter().map(|(_, schema)| schema);
394+
// Strip paths; when rle-to-dictionary is on, a column may be
395+
// RLE-encoded in one file but plain in another. Normalise to the
396+
// Dictionary type before merging so Schema::try_merge does not see a
397+
// type mismatch.
398+
let mut schemas: Vec<Schema> = schemas.into_iter().map(|(_, s)| s).collect();
399+
if self.options.global.enable_rle_to_dictionary {
400+
schemas = crate::schema_coercion::uniform_dict_schemas(schemas);
401+
}
392402

393403
let schema = if self.skip_metadata() {
394404
Schema::try_merge(clear_metadata(schemas))
@@ -510,7 +520,7 @@ impl FileFormat for ParquetFormat {
510520
source = source.with_parquet_file_reader_factory(cached_parquet_read_factory);
511521

512522
if let Some(metadata_size_hint) = metadata_size_hint {
513-
source = source.with_metadata_size_hint(metadata_size_hint)
523+
source = source.with_metadata_size_hint(metadata_size_hint);
514524
}
515525

516526
source = self.set_source_encryption_factory(source, state)?;
@@ -724,6 +734,7 @@ impl From<&ParquetFormatFactory> for protobuf::TableParquetOptions {
724734
compression_opt: global_options.global.compression.map(|compression| {
725735
parquet_options::CompressionOpt::Compression(compression)
726736
}),
737+
enable_rle_to_dictionary: global_options.global.enable_rle_to_dictionary,
727738
dictionary_enabled_opt: global_options.global.dictionary_enabled.map(|enabled| {
728739
parquet_options::DictionaryEnabledOpt::DictionaryEnabled(enabled)
729740
}),

‎datafusion/datasource-parquet/src/metadata.rs‎

Lines changed: 80 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ use arrow::datatypes::{DataType, Schema, SchemaRef, TimeUnit};
2727
use datafusion_common::encryption::FileDecryptionProperties;
2828
use datafusion_common::stats::Precision;
2929
use datafusion_common::{
30-
ColumnStatistics, DataFusionError, HashMap, Result, ScalarValue, Statistics,
30+
ColumnStatistics, DataFusionError, HashMap, HashSet, Result, ScalarValue, Statistics,
3131
internal_datafusion_err,
3232
};
3333
use datafusion_execution::cache::cache_manager::{
@@ -91,6 +91,9 @@ pub struct DFParquetMetadata<'a> {
9191
pub coerce_int96: Option<TimeUnit>,
9292
/// Optional timezone applied to INT96-coerced timestamps.
9393
pub coerce_int96_tz: Option<Arc<str>>,
94+
/// When true, columns that are physically RLE_DICTIONARY encoded in the parquet file
95+
/// are promoted to `Dictionary(Int32, …)` in the inferred schema.
96+
enable_rle_to_dictionary: bool,
9497
}
9598

9699
impl<'a> DFParquetMetadata<'a> {
@@ -108,9 +111,17 @@ impl<'a> DFParquetMetadata<'a> {
108111
page_index_policy: None,
109112
coerce_int96: None,
110113
coerce_int96_tz: None,
114+
enable_rle_to_dictionary: false,
111115
}
112116
}
113117

118+
/// When set, columns that are physically RLE_DICTIONARY encoded are promoted to
119+
/// `Dictionary(Int32, …)` in the schema returned by [`Self::fetch_schema`].
120+
pub fn with_enable_rle_to_dictionary(mut self, enable: bool) -> Self {
121+
self.enable_rle_to_dictionary = enable;
122+
self
123+
}
124+
114125
/// Set a hint for the number of trailing bytes to prefetch from the end
115126
/// of the file, equivalent to
116127
/// [`ParquetMetaDataReader::with_prefetch_hint`].
@@ -384,6 +395,74 @@ impl<'a> DFParquetMetadata<'a> {
384395
.coerce()
385396
})
386397
.unwrap_or(schema);
398+
399+
// Promote columns that are physically RLE_DICTIONARY encoded to Dictionary(Int32, …).
400+
// This happens before force_view_types so Dict columns survive the Utf8→Utf8View pass.
401+
let schema = if self.enable_rle_to_dictionary {
402+
let rle_cols: HashSet<String> = {
403+
let schema_descr = file_metadata.schema_descr();
404+
metadata
405+
.row_groups()
406+
.iter()
407+
.flat_map(|rg| {
408+
rg.columns()
409+
.iter()
410+
.enumerate()
411+
.filter_map(|(col_idx, col)| {
412+
col.dictionary_page_offset()?;
413+
let desc = schema_descr.column(col_idx);
414+
// Only promote top-level columns. For nested columns
415+
// the Parquet leaf name (e.g. "element" inside a List)
416+
// does not match the Arrow top-level field name (e.g.
417+
// "tags"), so using the leaf name here would miss the
418+
// target field or promote an unrelated one.
419+
let parts = desc.path().parts();
420+
(parts.len() == 1).then(|| parts[0].clone())
421+
})
422+
})
423+
.collect()
424+
};
425+
if rle_cols.is_empty() {
426+
schema
427+
} else {
428+
let promoted: Vec<_> = schema
429+
.fields()
430+
.iter()
431+
.map(|field| {
432+
if !rle_cols.contains(field.name()) {
433+
return Arc::clone(field);
434+
}
435+
let dict_value_type = match field.data_type() {
436+
DataType::Utf8 | DataType::LargeUtf8 => Some(DataType::Utf8),
437+
DataType::Binary | DataType::LargeBinary => {
438+
Some(DataType::Binary)
439+
}
440+
_ => None,
441+
};
442+
dict_value_type.map_or_else(
443+
|| Arc::clone(field),
444+
|value_type| {
445+
Arc::new(
446+
arrow::datatypes::Field::new(
447+
field.name(),
448+
DataType::Dictionary(
449+
Box::new(DataType::Int32),
450+
Box::new(value_type),
451+
),
452+
field.is_nullable(),
453+
)
454+
.with_metadata(field.metadata().clone()),
455+
)
456+
},
457+
)
458+
})
459+
.collect();
460+
Schema::new_with_metadata(promoted, schema.metadata().clone())
461+
}
462+
} else {
463+
schema
464+
};
465+
387466
Ok(schema)
388467
}
389468

0 commit comments

Comments
 (0)