Skip to main content

datafusion_common/file_options/
parquet_writer.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Options related to how parquet files should be written
19
20use 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/// Options for writing parquet files
43#[derive(Clone, Debug)]
44pub struct ParquetWriterOptions {
45    /// parquet-rs writer properties
46    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    /// Add the arrow schema to the parquet kv_metadata.
63    /// If already exists, then overwrites.
64    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        // ParquetWriterOptions will have defaults for the remaining fields (e.g. sorting_columns)
77        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    /// Convert the session's [`TableParquetOptions`] into a single write action's [`WriterPropertiesBuilder`].
88    ///
89    /// The returned [`WriterPropertiesBuilder`] includes customizations applicable per column.
90    /// Note that any encryption options are ignored as building the `FileEncryptionProperties`
91    /// might require other inputs besides the [`TableParquetOptions`].
92    fn try_from(table_parquet_options: &TableParquetOptions) -> Result<Self> {
93        // Table options include kv_metadata and col-specific options
94        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        // check that the arrow schema is present in the kv_metadata, if configured to do so
104        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        // add kv_meta, if any
113        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        // Apply column-specific options:
124        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
169/// Convert DataFusion's [`ParquetCdcOptions`] into parquet-rs's `Option<CdcOptions>`.
170///
171/// parquet-rs has no `enabled` flag; CDC is on when the option is `Some`. So a
172/// disabled [`ParquetCdcOptions`] maps to `None`, and an enabled one to `Some`
173/// with the chunking parameters.
174impl 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
186/// Convert parquet-rs's `Option<&CdcOptions>` back into DataFusion's
187/// [`ParquetCdcOptions`].
188///
189/// The presence of parquet-rs options means CDC was enabled, so `Some` maps to
190/// `enabled: true`; `None` yields the disabled default.
191impl 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    /// Convert the global session options, [`ParquetOptions`], into a single write action's [`WriterPropertiesBuilder`].
207    ///
208    /// The returned [`WriterPropertiesBuilder`] can then be further modified with additional options
209    /// applied per column; a customization which is not applicable for [`ParquetOptions`].
210    ///
211    /// Note that this method does not include the key_value_metadata from [`TableParquetOptions`].
212    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            // not in WriterProperties
234            enable_page_index: _,
235            pruning: _,
236            skip_metadata: _,
237            metadata_size_hint: _,
238            pushdown_filters: _,
239            reorder_filters: _,
240            force_filter_selections: _, // not used for writer props
241            allow_single_file_parallelism: _,
242            maximum_parallel_row_group_writers: _,
243            maximum_buffered_record_batches_per_stream: _,
244            bloom_filter_on_read: _, // reads not used for writer props
245            schema_force_view_types: _,
246            binary_as_string: _, // not used for writer props
247            coerce_int96: _,     // not used for writer props
248            coerce_int96_tz: _,  // not used for writer props
249            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        // We do not have access to default ColumnProperties set in Arrow.
284        // Therefore, only overwrite if these settings exist.
285        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
297/// Parses datafusion.execution.parquet.encoding String to a parquet::basic::Encoding
298pub(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
324/// Splits compression string into compression codec and optional compression_level
325/// I.e. gzip(2) -> gzip, 2
326fn split_compression_string(str_setting: &str) -> Result<(String, Option<u32>)> {
327    // ignore string literal chars passed from sqlparser i.e. remove single quotes
328    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
345/// Helper to ensure compression codecs which don't support levels
346/// don't have one set. E.g. snappy(2) is invalid.
347fn 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
356/// Helper to ensure compression codecs which require a level
357/// do have one set. E.g. zstd is invalid, zstd(3) is valid
358fn 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
364/// Parses datafusion.execution.parquet.compression String to a parquet::basic::Compression
365pub 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    /// Take the column defaults provided in [`ParquetOptions`], and generate a non-default col config.
448    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            // not in WriterProperties, but itemizing here to not skip newly added props
491            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    /// For testing only, take a single write's props and convert back into the session config.
545    /// (use identity to confirm correct.)
546    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                // global options
584                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                // global options which set the default column props
600                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                // not in WriterProperties
611                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        // TableParquetOptions, all props set to default
649        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        // see errors without the schema added, using default settings
656        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        // succeeds if we permit skipping the arrow schema
663        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        // Set the arrow schema back to required
671        table_parquet_opts = table_parquet_opts.with_skip_arrow_metadata(false);
672        // add the arrow schema to the kv_meta
673        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        // ParquetOptions, all props set to non-default
684        let parquet_options = parquet_options_with_non_defaults();
685
686        // TableParquetOptions, using ParquetOptions for global settings
687        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    /// Ensure that the configuration defaults for writing parquet files are
711    /// consistent with the options in arrow-rs
712    #[test]
713    fn test_defaults_match() {
714        // ensure the global settings are the same
715        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        // selectively skip the arrow_schema metadata, since the WriterProperties default has an empty kv_meta (no arrow schema)
723        default_table_writer_opts =
724            default_table_writer_opts.with_skip_arrow_metadata(true);
725
726        // WriterProperties::default, a.k.a. using extern parquet's defaults
727        let default_writer_props = WriterProperties::new();
728
729        // WriterProperties::try_from(TableParquetOptions::default), a.k.a. using datafusion's defaults
730        let from_datafusion_defaults =
731            WriterPropertiesBuilder::try_from(&default_table_writer_opts)
732                .unwrap()
733                .build();
734
735        // Expected: how the defaults should not match
736        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        // Expected: the datafusion default compression is different from arrow-rs's parquet
756        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        // Expected: the remaining should match
770        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        // the TableParquetOptions::default, with only the bloom filter turned on
786        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())); // add the required arrow schema
789        let from_datafusion_defaults =
790            WriterPropertiesBuilder::try_from(&default_table_writer_opts)
791                .unwrap()
792                .build();
793
794        // the WriterProperties::default, with only the bloom filter turned on
795        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        // the TableParquetOptions::default, with only fpp set
814        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())); // add the required arrow schema
818        let from_datafusion_defaults =
819            WriterPropertiesBuilder::try_from(&default_table_writer_opts)
820                .unwrap()
821                .build();
822
823        // the WriterProperties::default, with only fpp set
824        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        // Parameters are customized but `enabled` is false, so CDC stays off.
876        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        // the TableParquetOptions::default, with only ndv set
933        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())); // add the required arrow schema
937        let from_datafusion_defaults =
938            WriterPropertiesBuilder::try_from(&default_table_writer_opts)
939                .unwrap()
940                .build();
941
942        // the WriterProperties::default, with only ndv set
943        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}