1use std::sync::Arc;
21
22use crate::{
23 _internal_datafusion_err, DataFusionError, Result,
24 config::{ParquetCdcOptions, ParquetOptions, TableParquetOptions},
25};
26
27use arrow::datatypes::Schema;
28use parquet::arrow::encode_arrow_schema;
29use parquet::{
30 arrow::ARROW_SCHEMA_META_KEY,
31 basic::{BrotliLevel, GzipLevel, ZstdLevel},
32 file::{
33 metadata::KeyValue,
34 properties::{
35 DEFAULT_STATISTICS_ENABLED, EnabledStatistics, WriterProperties,
36 WriterPropertiesBuilder,
37 },
38 },
39 schema::types::ColumnPath,
40};
41
42#[derive(Clone, Debug)]
44pub struct ParquetWriterOptions {
45 pub writer_options: WriterProperties,
47}
48
49impl ParquetWriterOptions {
50 pub fn new(writer_options: WriterProperties) -> Self {
51 Self { writer_options }
52 }
53}
54
55impl ParquetWriterOptions {
56 pub fn writer_options(&self) -> &WriterProperties {
57 &self.writer_options
58 }
59}
60
61impl TableParquetOptions {
62 pub fn arrow_schema(&mut self, schema: &Arc<Schema>) {
65 self.key_value_metadata.insert(
66 ARROW_SCHEMA_META_KEY.into(),
67 Some(encode_arrow_schema(schema)),
68 );
69 }
70}
71
72impl TryFrom<&TableParquetOptions> for ParquetWriterOptions {
73 type Error = DataFusionError;
74
75 fn try_from(parquet_table_options: &TableParquetOptions) -> Result<Self> {
76 Ok(ParquetWriterOptions {
78 writer_options: WriterPropertiesBuilder::try_from(parquet_table_options)?
79 .build(),
80 })
81 }
82}
83
84impl TryFrom<&TableParquetOptions> for WriterPropertiesBuilder {
85 type Error = DataFusionError;
86
87 fn try_from(table_parquet_options: &TableParquetOptions) -> Result<Self> {
93 let TableParquetOptions {
95 global,
96 column_specific_options,
97 key_value_metadata,
98 ..
99 } = table_parquet_options;
100
101 let mut builder = global.into_writer_properties_builder()?;
102
103 if !global.skip_arrow_metadata
105 && !key_value_metadata.contains_key(ARROW_SCHEMA_META_KEY)
106 {
107 return Err(_internal_datafusion_err!(
108 "arrow schema was not added to the kv_metadata, even though it is required by configuration settings"
109 ));
110 }
111
112 if !key_value_metadata.is_empty() {
114 builder = builder.set_key_value_metadata(Some(
115 key_value_metadata
116 .to_owned()
117 .drain()
118 .map(|(key, value)| KeyValue { key, value })
119 .collect(),
120 ));
121 }
122
123 for (column, options) in column_specific_options {
125 let path = ColumnPath::new(column.split('.').map(|s| s.to_owned()).collect());
126
127 if let Some(bloom_filter_enabled) = options.bloom_filter_enabled {
128 builder = builder
129 .set_column_bloom_filter_enabled(path.clone(), bloom_filter_enabled);
130 }
131
132 if let Some(encoding) = &options.encoding {
133 let parsed_encoding = parse_encoding_string(encoding)?;
134 builder = builder.set_column_encoding(path.clone(), parsed_encoding);
135 }
136
137 if let Some(dictionary_enabled) = options.dictionary_enabled {
138 builder = builder
139 .set_column_dictionary_enabled(path.clone(), dictionary_enabled);
140 }
141
142 if let Some(compression) = &options.compression {
143 let parsed_compression = parse_compression_string(compression)?;
144 builder =
145 builder.set_column_compression(path.clone(), parsed_compression);
146 }
147
148 if let Some(statistics_enabled) = &options.statistics_enabled {
149 let parsed_value = parse_statistics_string(statistics_enabled)?;
150 builder =
151 builder.set_column_statistics_enabled(path.clone(), parsed_value);
152 }
153
154 if let Some(bloom_filter_fpp) = options.bloom_filter_fpp {
155 builder =
156 builder.set_column_bloom_filter_fpp(path.clone(), bloom_filter_fpp);
157 }
158
159 if let Some(bloom_filter_ndv) = options.bloom_filter_ndv {
160 builder = builder
161 .set_column_bloom_filter_max_ndv(path.clone(), bloom_filter_ndv);
162 }
163 }
164
165 Ok(builder)
166 }
167}
168
169impl From<&ParquetCdcOptions> for Option<parquet::file::properties::CdcOptions> {
175 fn from(value: &ParquetCdcOptions) -> Self {
176 value
177 .enabled
178 .then_some(parquet::file::properties::CdcOptions {
179 min_chunk_size: value.min_chunk_size,
180 max_chunk_size: value.max_chunk_size,
181 norm_level: value.norm_level,
182 })
183 }
184}
185
186impl From<Option<&parquet::file::properties::CdcOptions>> for ParquetCdcOptions {
192 fn from(value: Option<&parquet::file::properties::CdcOptions>) -> Self {
193 match value {
194 Some(cdc) => ParquetCdcOptions {
195 enabled: true,
196 min_chunk_size: cdc.min_chunk_size,
197 max_chunk_size: cdc.max_chunk_size,
198 norm_level: cdc.norm_level,
199 },
200 None => ParquetCdcOptions::default(),
201 }
202 }
203}
204
205impl ParquetOptions {
206 pub fn into_writer_properties_builder(&self) -> Result<WriterPropertiesBuilder> {
213 let ParquetOptions {
214 data_pagesize_limit,
215 write_batch_size,
216 writer_version,
217 compression,
218 dictionary_enabled,
219 dictionary_page_size_limit,
220 statistics_enabled,
221 max_row_group_size,
222 max_row_group_bytes,
223 created_by,
224 column_index_truncate_length,
225 statistics_truncate_length,
226 data_page_row_count_limit,
227 encoding,
228 bloom_filter_on_write,
229 bloom_filter_fpp,
230 bloom_filter_ndv,
231 content_defined_chunking,
232
233 enable_page_index: _,
235 pruning: _,
236 skip_metadata: _,
237 metadata_size_hint: _,
238 pushdown_filters: _,
239 reorder_filters: _,
240 force_filter_selections: _, allow_single_file_parallelism: _,
242 maximum_parallel_row_group_writers: _,
243 maximum_buffered_record_batches_per_stream: _,
244 bloom_filter_on_read: _, schema_force_view_types: _,
246 binary_as_string: _, coerce_int96: _, coerce_int96_tz: _, skip_arrow_metadata: _,
250 max_predicate_cache_size: _,
251 max_in_list_size: _,
252 } = self;
253
254 let mut builder = WriterProperties::builder()
255 .set_data_page_size_limit(*data_pagesize_limit)
256 .set_write_batch_size(*write_batch_size)
257 .set_writer_version((*writer_version).into())
258 .set_dictionary_page_size_limit(*dictionary_page_size_limit)
259 .set_statistics_enabled(
260 statistics_enabled
261 .as_ref()
262 .and_then(|s| parse_statistics_string(s).ok())
263 .unwrap_or(DEFAULT_STATISTICS_ENABLED),
264 )
265 .set_max_row_group_row_count(Some(*max_row_group_size))
266 .set_max_row_group_bytes(max_row_group_bytes.as_ref().map(|v| v.get()))
267 .set_created_by(created_by.clone())
268 .set_column_index_truncate_length(*column_index_truncate_length)
269 .set_statistics_truncate_length(*statistics_truncate_length)
270 .set_data_page_row_count_limit(*data_page_row_count_limit)
271 .set_bloom_filter_enabled(*bloom_filter_on_write);
272
273 if let Some(bloom_filter_fpp) = bloom_filter_fpp {
274 builder = builder.set_bloom_filter_fpp(*bloom_filter_fpp);
275 };
276 if let Some(bloom_filter_ndv) = bloom_filter_ndv {
277 builder = builder.set_bloom_filter_max_ndv(*bloom_filter_ndv);
278 };
279 if let Some(dictionary_enabled) = dictionary_enabled {
280 builder = builder.set_dictionary_enabled(*dictionary_enabled);
281 };
282
283 if let Some(compression) = compression {
286 builder = builder.set_compression(parse_compression_string(compression)?);
287 }
288 if let Some(encoding) = encoding {
289 builder = builder.set_encoding(parse_encoding_string(encoding)?);
290 }
291 builder = builder.set_content_defined_chunking(content_defined_chunking.into());
292
293 Ok(builder)
294 }
295}
296
297pub(crate) fn parse_encoding_string(
299 str_setting: &str,
300) -> Result<parquet::basic::Encoding> {
301 let str_setting_lower: &str = &str_setting.to_lowercase();
302 match str_setting_lower {
303 "plain" => Ok(parquet::basic::Encoding::PLAIN),
304 "plain_dictionary" => Ok(parquet::basic::Encoding::PLAIN_DICTIONARY),
305 "rle" => Ok(parquet::basic::Encoding::RLE),
306 #[expect(deprecated)]
307 "bit_packed" => Ok(parquet::basic::Encoding::BIT_PACKED),
308 "delta_binary_packed" => Ok(parquet::basic::Encoding::DELTA_BINARY_PACKED),
309 "delta_length_byte_array" => {
310 Ok(parquet::basic::Encoding::DELTA_LENGTH_BYTE_ARRAY)
311 }
312 "delta_byte_array" => Ok(parquet::basic::Encoding::DELTA_BYTE_ARRAY),
313 "rle_dictionary" => Ok(parquet::basic::Encoding::RLE_DICTIONARY),
314 "byte_stream_split" => Ok(parquet::basic::Encoding::BYTE_STREAM_SPLIT),
315 _ => Err(DataFusionError::Configuration(format!(
316 "Unknown or unsupported parquet encoding: \
317 {str_setting}. Valid values are: plain, plain_dictionary, rle, \
318 bit_packed, delta_binary_packed, delta_length_byte_array, \
319 delta_byte_array, rle_dictionary, and byte_stream_split."
320 ))),
321 }
322}
323
324fn split_compression_string(str_setting: &str) -> Result<(String, Option<u32>)> {
327 let str_setting = str_setting.replace('\'', "");
329 let split_setting = str_setting.split_once('(');
330
331 match split_setting {
332 Some((codec, rh)) => {
333 let level = &rh[..rh.len() - 1].parse::<u32>().map_err(|_| {
334 DataFusionError::Configuration(format!(
335 "Could not parse compression string. \
336 Got codec: {codec} and unknown level from {str_setting}"
337 ))
338 })?;
339 Ok((codec.to_owned(), Some(*level)))
340 }
341 None => Ok((str_setting.to_owned(), None)),
342 }
343}
344
345fn check_level_is_none(codec: &str, level: &Option<u32>) -> Result<()> {
348 if level.is_some() {
349 return Err(DataFusionError::Configuration(format!(
350 "Compression {codec} does not support specifying a level"
351 )));
352 }
353 Ok(())
354}
355
356fn require_level(codec: &str, level: Option<u32>) -> Result<u32> {
359 level.ok_or(DataFusionError::Configuration(format!(
360 "{codec} compression requires specifying a level such as {codec}(4)"
361 )))
362}
363
364pub fn parse_compression_string(
366 str_setting: &str,
367) -> Result<parquet::basic::Compression> {
368 let str_setting_lower: &str = &str_setting.to_lowercase();
369 let (codec, level) = split_compression_string(str_setting_lower)?;
370 let codec = codec.as_str();
371 match codec {
372 "uncompressed" => {
373 check_level_is_none(codec, &level)?;
374 Ok(parquet::basic::Compression::UNCOMPRESSED)
375 }
376 "snappy" => {
377 check_level_is_none(codec, &level)?;
378 Ok(parquet::basic::Compression::SNAPPY)
379 }
380 "gzip" => {
381 let level = require_level(codec, level)?;
382 Ok(parquet::basic::Compression::GZIP(GzipLevel::try_new(
383 level,
384 )?))
385 }
386 "brotli" => {
387 let level = require_level(codec, level)?;
388 Ok(parquet::basic::Compression::BROTLI(BrotliLevel::try_new(
389 level,
390 )?))
391 }
392 "lz4" => {
393 check_level_is_none(codec, &level)?;
394 Ok(parquet::basic::Compression::LZ4)
395 }
396 "zstd" => {
397 let level = require_level(codec, level)?;
398 Ok(parquet::basic::Compression::ZSTD(ZstdLevel::try_new(
399 level as i32,
400 )?))
401 }
402 "lz4_raw" => {
403 check_level_is_none(codec, &level)?;
404 Ok(parquet::basic::Compression::LZ4_RAW)
405 }
406 _ => Err(DataFusionError::Configuration(format!(
407 "Unknown or unsupported parquet compression: \
408 {str_setting}. Valid values are: uncompressed, snappy, gzip(level), \
409 brotli(level), lz4, zstd(level), and lz4_raw."
410 ))),
411 }
412}
413
414pub(crate) fn parse_statistics_string(str_setting: &str) -> Result<EnabledStatistics> {
415 let str_setting_lower: &str = &str_setting.to_lowercase();
416 match str_setting_lower {
417 "none" => Ok(EnabledStatistics::None),
418 "chunk" => Ok(EnabledStatistics::Chunk),
419 "page" => Ok(EnabledStatistics::Page),
420 _ => Err(DataFusionError::Configuration(format!(
421 "Unknown or unsupported parquet statistics setting {str_setting} \
422 valid options are none, page, and chunk"
423 ))),
424 }
425}
426
427#[cfg(feature = "parquet")]
428#[cfg(test)]
429mod tests {
430 use super::*;
431 #[cfg(feature = "parquet_encryption")]
432 use crate::config::ConfigFileEncryptionProperties;
433 use crate::config::{
434 MaxRowGroupBytes, ParquetCdcOptions, ParquetColumnOptions,
435 ParquetEncryptionOptions, ParquetOptions,
436 };
437 use crate::parquet_config::DFParquetWriterVersion;
438 use parquet::basic::Compression;
439 use parquet::file::properties::{
440 BloomFilterProperties, DEFAULT_BLOOM_FILTER_FPP, DEFAULT_BLOOM_FILTER_NDV,
441 DEFAULT_MAX_ROW_GROUP_ROW_COUNT, EnabledStatistics,
442 };
443 use std::collections::HashMap;
444
445 const COL_NAME: &str = "configured";
446
447 fn column_options_with_non_defaults(
449 src_col_defaults: &ParquetOptions,
450 ) -> ParquetColumnOptions {
451 ParquetColumnOptions {
452 compression: Some("zstd(22)".into()),
453 dictionary_enabled: src_col_defaults.dictionary_enabled.map(|v| !v),
454 statistics_enabled: Some("none".into()),
455 encoding: Some("RLE".into()),
456 bloom_filter_enabled: Some(true),
457 bloom_filter_fpp: Some(0.72),
458 bloom_filter_ndv: Some(72),
459 }
460 }
461
462 fn parquet_options_with_non_defaults() -> ParquetOptions {
463 let defaults = ParquetOptions::default();
464 let writer_version = if defaults.writer_version.eq(&DFParquetWriterVersion::V1_0)
465 {
466 DFParquetWriterVersion::V2_0
467 } else {
468 DFParquetWriterVersion::V1_0
469 };
470
471 ParquetOptions {
472 data_pagesize_limit: 42,
473 write_batch_size: 42,
474 writer_version,
475 compression: Some("zstd(22)".into()),
476 dictionary_enabled: Some(!defaults.dictionary_enabled.unwrap_or(false)),
477 dictionary_page_size_limit: 43,
478 statistics_enabled: Some("chunk".into()),
479 max_row_group_size: 42,
480 max_row_group_bytes: Some(MaxRowGroupBytes::try_new(42).unwrap()),
481 created_by: "wordy".into(),
482 column_index_truncate_length: Some(42),
483 statistics_truncate_length: Some(42),
484 data_page_row_count_limit: 42,
485 encoding: Some("BYTE_STREAM_SPLIT".into()),
486 bloom_filter_on_write: !defaults.bloom_filter_on_write,
487 bloom_filter_fpp: Some(0.42),
488 bloom_filter_ndv: Some(42),
489
490 enable_page_index: defaults.enable_page_index,
492 pruning: defaults.pruning,
493 max_in_list_size: defaults.max_in_list_size,
494 skip_metadata: defaults.skip_metadata,
495 metadata_size_hint: defaults.metadata_size_hint,
496 pushdown_filters: defaults.pushdown_filters,
497 reorder_filters: defaults.reorder_filters,
498 force_filter_selections: defaults.force_filter_selections,
499 allow_single_file_parallelism: defaults.allow_single_file_parallelism,
500 maximum_parallel_row_group_writers: defaults
501 .maximum_parallel_row_group_writers,
502 maximum_buffered_record_batches_per_stream: defaults
503 .maximum_buffered_record_batches_per_stream,
504 bloom_filter_on_read: defaults.bloom_filter_on_read,
505 schema_force_view_types: defaults.schema_force_view_types,
506 binary_as_string: defaults.binary_as_string,
507 skip_arrow_metadata: defaults.skip_arrow_metadata,
508 coerce_int96: None,
509 coerce_int96_tz: None,
510 max_predicate_cache_size: defaults.max_predicate_cache_size,
511 content_defined_chunking: defaults.content_defined_chunking.clone(),
512 }
513 }
514
515 fn extract_column_options(
516 props: &WriterProperties,
517 col: ColumnPath,
518 ) -> ParquetColumnOptions {
519 let bloom_filter_default_props = props.bloom_filter_properties(&col);
520
521 ParquetColumnOptions {
522 bloom_filter_enabled: Some(bloom_filter_default_props.is_some()),
523 encoding: props.encoding(&col).map(|s| s.to_string()),
524 dictionary_enabled: Some(props.dictionary_enabled(&col)),
525 compression: match props.compression(&col) {
526 Compression::ZSTD(lvl) => {
527 Some(format!("zstd({})", lvl.compression_level()))
528 }
529 _ => None,
530 },
531 statistics_enabled: Some(
532 match props.statistics_enabled(&col) {
533 EnabledStatistics::None => "none",
534 EnabledStatistics::Chunk => "chunk",
535 EnabledStatistics::Page => "page",
536 }
537 .into(),
538 ),
539 bloom_filter_fpp: bloom_filter_default_props.map(|p| p.fpp()),
540 bloom_filter_ndv: bloom_filter_default_props.map(|p| p.ndv()),
541 }
542 }
543
544 fn session_config_from_writer_props(props: &WriterProperties) -> TableParquetOptions {
547 let default_col = ColumnPath::from("col doesn't have specific config");
548 let default_col_props = extract_column_options(props, default_col);
549
550 let configured_col = ColumnPath::from(COL_NAME);
551 let configured_col_props = extract_column_options(props, configured_col);
552
553 let key_value_metadata = props
554 .key_value_metadata()
555 .map(|pairs| {
556 HashMap::from_iter(
557 pairs
558 .iter()
559 .cloned()
560 .map(|KeyValue { key, value }| (key, value)),
561 )
562 })
563 .unwrap_or_default();
564
565 let global_options_defaults = ParquetOptions::default();
566
567 let column_specific_options = if configured_col_props.eq(&default_col_props) {
568 HashMap::default()
569 } else {
570 HashMap::from([(COL_NAME.into(), configured_col_props)])
571 };
572
573 #[cfg(feature = "parquet_encryption")]
574 let fep = props
575 .file_encryption_properties()
576 .map(ConfigFileEncryptionProperties::from);
577
578 #[cfg(not(feature = "parquet_encryption"))]
579 let fep = None;
580
581 TableParquetOptions {
582 global: ParquetOptions {
583 data_pagesize_limit: props.data_page_size_limit(),
585 write_batch_size: props.write_batch_size(),
586 writer_version: props.writer_version().into(),
587 dictionary_page_size_limit: props.dictionary_page_size_limit(),
588 max_row_group_size: props
589 .max_row_group_row_count()
590 .unwrap_or(DEFAULT_MAX_ROW_GROUP_ROW_COUNT),
591 max_row_group_bytes: props
592 .max_row_group_bytes()
593 .and_then(|v| MaxRowGroupBytes::try_new(v).ok()),
594 created_by: props.created_by().to_string(),
595 column_index_truncate_length: props.column_index_truncate_length(),
596 statistics_truncate_length: props.statistics_truncate_length(),
597 data_page_row_count_limit: props.data_page_row_count_limit(),
598
599 encoding: default_col_props.encoding,
601 compression: default_col_props.compression,
602 dictionary_enabled: default_col_props.dictionary_enabled,
603 statistics_enabled: default_col_props.statistics_enabled,
604 bloom_filter_on_write: default_col_props
605 .bloom_filter_enabled
606 .unwrap_or_default(),
607 bloom_filter_fpp: default_col_props.bloom_filter_fpp,
608 bloom_filter_ndv: default_col_props.bloom_filter_ndv,
609
610 enable_page_index: global_options_defaults.enable_page_index,
612 pruning: global_options_defaults.pruning,
613 max_in_list_size: global_options_defaults.max_in_list_size,
614 skip_metadata: global_options_defaults.skip_metadata,
615 metadata_size_hint: global_options_defaults.metadata_size_hint,
616 pushdown_filters: global_options_defaults.pushdown_filters,
617 reorder_filters: global_options_defaults.reorder_filters,
618 force_filter_selections: global_options_defaults.force_filter_selections,
619 allow_single_file_parallelism: global_options_defaults
620 .allow_single_file_parallelism,
621 maximum_parallel_row_group_writers: global_options_defaults
622 .maximum_parallel_row_group_writers,
623 maximum_buffered_record_batches_per_stream: global_options_defaults
624 .maximum_buffered_record_batches_per_stream,
625 bloom_filter_on_read: global_options_defaults.bloom_filter_on_read,
626 max_predicate_cache_size: global_options_defaults
627 .max_predicate_cache_size,
628 schema_force_view_types: global_options_defaults.schema_force_view_types,
629 binary_as_string: global_options_defaults.binary_as_string,
630 skip_arrow_metadata: global_options_defaults.skip_arrow_metadata,
631 coerce_int96: None,
632 coerce_int96_tz: None,
633 content_defined_chunking: props.content_defined_chunking().into(),
634 },
635 column_specific_options,
636 key_value_metadata,
637 crypto: ParquetEncryptionOptions {
638 file_encryption: fep,
639 file_decryption: None,
640 factory_id: None,
641 factory_options: Default::default(),
642 },
643 }
644 }
645
646 #[test]
647 fn table_parquet_opts_to_writer_props_skip_arrow_metadata() {
648 let mut table_parquet_opts = TableParquetOptions::default();
650 assert!(
651 !table_parquet_opts.global.skip_arrow_metadata,
652 "default false, to not skip the arrow schema requirement"
653 );
654
655 let should_error = WriterPropertiesBuilder::try_from(&table_parquet_opts);
657 assert!(
658 should_error.is_err(),
659 "should error without the required arrow schema in kv_metadata",
660 );
661
662 table_parquet_opts = table_parquet_opts.with_skip_arrow_metadata(true);
664 let should_succeed = WriterPropertiesBuilder::try_from(&table_parquet_opts);
665 assert!(
666 should_succeed.is_ok(),
667 "should work with the arrow schema skipped by config",
668 );
669
670 table_parquet_opts = table_parquet_opts.with_skip_arrow_metadata(false);
672 table_parquet_opts.arrow_schema(&Arc::new(Schema::empty()));
674 let should_succeed = WriterPropertiesBuilder::try_from(&table_parquet_opts);
675 assert!(
676 should_succeed.is_ok(),
677 "should work with the arrow schema included in TableParquetOptions",
678 );
679 }
680
681 #[test]
682 fn table_parquet_opts_to_writer_props() {
683 let parquet_options = parquet_options_with_non_defaults();
685
686 let key = ARROW_SCHEMA_META_KEY.to_string();
688 let value = Some("bar".into());
689 let table_parquet_opts = TableParquetOptions {
690 global: parquet_options.clone(),
691 column_specific_options: [(
692 COL_NAME.into(),
693 column_options_with_non_defaults(&parquet_options),
694 )]
695 .into(),
696 key_value_metadata: [(key, value)].into(),
697 crypto: Default::default(),
698 };
699
700 let writer_props = WriterPropertiesBuilder::try_from(&table_parquet_opts)
701 .unwrap()
702 .build();
703 assert_eq!(
704 table_parquet_opts,
705 session_config_from_writer_props(&writer_props),
706 "the writer_props should have the same configuration as the session's TableParquetOptions",
707 );
708 }
709
710 #[test]
713 fn test_defaults_match() {
714 let mut default_table_writer_opts = TableParquetOptions::default();
716 let default_parquet_opts = ParquetOptions::default();
717 assert_eq!(
718 default_table_writer_opts.global, default_parquet_opts,
719 "should have matching defaults for TableParquetOptions.global and ParquetOptions",
720 );
721
722 default_table_writer_opts =
724 default_table_writer_opts.with_skip_arrow_metadata(true);
725
726 let default_writer_props = WriterProperties::new();
728
729 let from_datafusion_defaults =
731 WriterPropertiesBuilder::try_from(&default_table_writer_opts)
732 .unwrap()
733 .build();
734
735 assert_ne!(
737 default_writer_props.created_by(),
738 from_datafusion_defaults.created_by(),
739 "should have different created_by sources",
740 );
741 assert!(
742 default_writer_props
743 .created_by()
744 .starts_with("parquet-rs version"),
745 "should indicate that writer_props defaults came from the extern parquet crate",
746 );
747 assert!(
748 default_table_writer_opts
749 .global
750 .created_by
751 .starts_with("datafusion version"),
752 "should indicate that table_parquet_opts defaults came from datafusion",
753 );
754
755 assert_eq!(
757 default_writer_props.compression(&"default".into()),
758 Compression::UNCOMPRESSED,
759 "extern parquet's default is None"
760 );
761 assert!(
762 matches!(
763 from_datafusion_defaults.compression(&"default".into()),
764 Compression::ZSTD(_)
765 ),
766 "datafusion's default is zstd"
767 );
768
769 let same_created_by = default_table_writer_opts.global.created_by.clone();
771 let mut from_extern_parquet =
772 session_config_from_writer_props(&default_writer_props);
773 from_extern_parquet.global.created_by = same_created_by;
774 from_extern_parquet.global.compression = Some("zstd(3)".into());
775 from_extern_parquet.global.skip_arrow_metadata = true;
776
777 assert_eq!(
778 default_table_writer_opts, from_extern_parquet,
779 "the default writer_props should have the same configuration as the session's default TableParquetOptions",
780 );
781 }
782
783 #[test]
784 fn test_bloom_filter_defaults() {
785 let mut default_table_writer_opts = TableParquetOptions::default();
787 default_table_writer_opts.global.bloom_filter_on_write = true;
788 default_table_writer_opts.arrow_schema(&Arc::new(Schema::empty())); let from_datafusion_defaults =
790 WriterPropertiesBuilder::try_from(&default_table_writer_opts)
791 .unwrap()
792 .build();
793
794 let default_writer_props = WriterProperties::builder()
796 .set_bloom_filter_enabled(true)
797 .build();
798
799 assert_eq!(
800 default_writer_props.bloom_filter_properties(&"default".into()),
801 from_datafusion_defaults.bloom_filter_properties(&"default".into()),
802 "parquet and datafusion props, should have the same bloom filter props",
803 );
804 assert_eq!(
805 default_writer_props.bloom_filter_properties(&"default".into()),
806 Some(&BloomFilterProperties::default()),
807 "should use the default bloom filter props"
808 );
809 }
810
811 #[test]
812 fn test_bloom_filter_set_fpp_only() {
813 let mut default_table_writer_opts = TableParquetOptions::default();
815 default_table_writer_opts.global.bloom_filter_on_write = true;
816 default_table_writer_opts.global.bloom_filter_fpp = Some(0.42);
817 default_table_writer_opts.arrow_schema(&Arc::new(Schema::empty())); let from_datafusion_defaults =
819 WriterPropertiesBuilder::try_from(&default_table_writer_opts)
820 .unwrap()
821 .build();
822
823 let default_writer_props = WriterProperties::builder()
825 .set_bloom_filter_enabled(true)
826 .set_bloom_filter_fpp(0.42)
827 .build();
828
829 assert_eq!(
830 default_writer_props.bloom_filter_properties(&"default".into()),
831 from_datafusion_defaults.bloom_filter_properties(&"default".into()),
832 "parquet and datafusion props, should have the same bloom filter props",
833 );
834 assert_eq!(
835 default_writer_props.bloom_filter_properties(&"default".into()),
836 Some(
837 &BloomFilterProperties::builder()
838 .with_fpp(0.42)
839 .with_max_ndv(DEFAULT_BLOOM_FILTER_NDV)
840 .build()
841 ),
842 "should have only the fpp set, and the ndv at default",
843 );
844 }
845
846 #[test]
847 fn test_cdc_enabled_with_custom_options() {
848 let mut opts = TableParquetOptions::default();
849 opts.global.content_defined_chunking = ParquetCdcOptions {
850 enabled: true,
851 min_chunk_size: 128 * 1024,
852 max_chunk_size: 512 * 1024,
853 norm_level: 2,
854 };
855 opts.arrow_schema(&Arc::new(Schema::empty()));
856
857 let props = WriterPropertiesBuilder::try_from(&opts).unwrap().build();
858 let cdc = props.content_defined_chunking().expect("CDC should be set");
859 assert_eq!(cdc.min_chunk_size, 128 * 1024);
860 assert_eq!(cdc.max_chunk_size, 512 * 1024);
861 assert_eq!(cdc.norm_level, 2);
862 }
863
864 #[test]
865 fn test_cdc_disabled_by_default() {
866 let mut opts = TableParquetOptions::default();
867 opts.arrow_schema(&Arc::new(Schema::empty()));
868
869 let props = WriterPropertiesBuilder::try_from(&opts).unwrap().build();
870 assert!(props.content_defined_chunking().is_none());
871 }
872
873 #[test]
874 fn test_cdc_params_ignored_when_disabled() {
875 let mut opts = TableParquetOptions::default();
877 opts.global.content_defined_chunking = ParquetCdcOptions {
878 enabled: false,
879 min_chunk_size: 128 * 1024,
880 max_chunk_size: 512 * 1024,
881 norm_level: 2,
882 };
883 opts.arrow_schema(&Arc::new(Schema::empty()));
884
885 let props = WriterPropertiesBuilder::try_from(&opts).unwrap().build();
886 assert!(props.content_defined_chunking().is_none());
887 }
888
889 #[test]
890 fn test_cdc_round_trip_through_writer_props() {
891 let mut opts = TableParquetOptions::default();
892 opts.global.content_defined_chunking = ParquetCdcOptions {
893 enabled: true,
894 min_chunk_size: 64 * 1024,
895 max_chunk_size: 2 * 1024 * 1024,
896 norm_level: -1,
897 };
898 opts.arrow_schema(&Arc::new(Schema::empty()));
899
900 let props = WriterPropertiesBuilder::try_from(&opts).unwrap().build();
901 let recovered = session_config_from_writer_props(&props);
902
903 let cdc = recovered.global.content_defined_chunking;
904 assert!(cdc.enabled);
905 assert_eq!(cdc.min_chunk_size, 64 * 1024);
906 assert_eq!(cdc.max_chunk_size, 2 * 1024 * 1024);
907 assert_eq!(cdc.norm_level, -1);
908 }
909
910 #[test]
911 fn test_max_row_group_bytes_disabled_by_default() {
912 let mut opts = TableParquetOptions::default();
913 opts.arrow_schema(&Arc::new(Schema::empty()));
914
915 let props = WriterPropertiesBuilder::try_from(&opts).unwrap().build();
916 assert_eq!(props.max_row_group_bytes(), None);
917 }
918
919 #[test]
920 fn test_max_row_group_bytes_propagated_to_writer_props() {
921 let mut opts = TableParquetOptions::default();
922 opts.global.max_row_group_bytes =
923 Some(MaxRowGroupBytes::try_new(64 * 1024 * 1024).unwrap());
924 opts.arrow_schema(&Arc::new(Schema::empty()));
925
926 let props = WriterPropertiesBuilder::try_from(&opts).unwrap().build();
927 assert_eq!(props.max_row_group_bytes(), Some(64 * 1024 * 1024));
928 }
929
930 #[test]
931 fn test_bloom_filter_set_ndv_only() {
932 let mut default_table_writer_opts = TableParquetOptions::default();
934 default_table_writer_opts.global.bloom_filter_on_write = true;
935 default_table_writer_opts.global.bloom_filter_ndv = Some(42);
936 default_table_writer_opts.arrow_schema(&Arc::new(Schema::empty())); let from_datafusion_defaults =
938 WriterPropertiesBuilder::try_from(&default_table_writer_opts)
939 .unwrap()
940 .build();
941
942 let default_writer_props = WriterProperties::builder()
944 .set_bloom_filter_enabled(true)
945 .set_bloom_filter_max_ndv(42)
946 .build();
947
948 assert_eq!(
949 default_writer_props.bloom_filter_properties(&"default".into()),
950 from_datafusion_defaults.bloom_filter_properties(&"default".into()),
951 "parquet and datafusion props, should have the same bloom filter props",
952 );
953 assert_eq!(
954 default_writer_props.bloom_filter_properties(&"default".into()),
955 Some(
956 &BloomFilterProperties::builder()
957 .with_fpp(DEFAULT_BLOOM_FILTER_FPP)
958 .with_max_ndv(42)
959 .build()
960 ),
961 "should have only the ndv set, and the fpp at default",
962 );
963 }
964}