Skip to main content

iceberg/writer/file_writer/
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//! The module contains the file writer for parquet file format.
19
20use std::collections::HashMap;
21use std::sync::Arc;
22
23use arrow_schema::SchemaRef as ArrowSchemaRef;
24use bytes::Bytes;
25use futures::future::BoxFuture;
26use itertools::Itertools;
27use parquet::arrow::AsyncArrowWriter;
28use parquet::arrow::async_reader::AsyncFileReader;
29use parquet::arrow::async_writer::AsyncFileWriter as ArrowAsyncFileWriter;
30use parquet::basic::{BrotliLevel, Compression, GzipLevel, ZstdLevel};
31use parquet::encryption::encrypt::FileEncryptionProperties;
32use parquet::file::metadata::ParquetMetaData;
33use parquet::file::properties::{CdcOptions, WriterProperties};
34use parquet::file::statistics::Statistics;
35
36use super::{FileWriter, FileWriterBuilder};
37use crate::arrow::{
38    ArrowFileReader, DEFAULT_MAP_FIELD_NAME, FieldMatchMode, NanValueCountVisitor,
39    get_parquet_stat_max_as_datum, get_parquet_stat_min_as_datum,
40};
41use crate::compression::CompressionCodec;
42use crate::encryption::{EncryptionManager, StandardKeyMetadata};
43use crate::error::invalid_data;
44use crate::io::{FileIO, FileWrite, OutputFile};
45use crate::spec::{
46    DataContentType, DataFileBuilder, DataFileFormat, Datum, ListType, Literal, MapType,
47    NestedFieldRef, PartitionSpec, PrimitiveType, Schema, SchemaRef, SchemaVisitor, Struct,
48    StructType, TableMetadata, TableProperties, Type, VariantType, visit_schema,
49};
50use crate::transform::create_transform_function;
51use crate::writer::{CurrentFileStatus, DataFile};
52use crate::{Error, ErrorKind, Result};
53
54/// ParquetWriterBuilder is used to builder a [`ParquetWriter`]
55#[derive(Clone, Debug)]
56pub struct ParquetWriterBuilder {
57    props: WriterProperties,
58    schema: SchemaRef,
59    match_mode: FieldMatchMode,
60    encryption_manager: Option<Arc<EncryptionManager>>,
61}
62
63impl ParquetWriterBuilder {
64    /// Create a new `ParquetWriterBuilder`
65    /// To construct the write result, the schema should contain the `PARQUET_FIELD_ID_META_KEY` metadata for each field.
66    ///
67    /// When writing into an existing Iceberg table, prefer
68    /// [`Self::from_table_properties`], which derives `WriterProperties` from
69    /// the table's `write.parquet.*` properties.
70    pub fn new(props: WriterProperties, schema: SchemaRef) -> Self {
71        Self::new_with_match_mode(props, schema, FieldMatchMode::Id)
72    }
73
74    /// Create a new `ParquetWriterBuilder` with custom match mode
75    pub fn new_with_match_mode(
76        props: WriterProperties,
77        schema: SchemaRef,
78        match_mode: FieldMatchMode,
79    ) -> Self {
80        Self {
81            props,
82            schema,
83            match_mode,
84            encryption_manager: None,
85        }
86    }
87
88    /// Build a `ParquetWriterBuilder` from Iceberg table properties and a
89    /// schema, translating `write.parquet.*` settings into `WriterProperties`
90    /// instead of using parquet-rs defaults.
91    pub fn from_table_properties(
92        table_props: &TableProperties<'_>,
93        schema: SchemaRef,
94    ) -> Result<Self> {
95        let cdc = if table_props.cdc_enabled()? {
96            Some(CdcOptions {
97                min_chunk_size: table_props.cdc_min_chunk_size()?,
98                max_chunk_size: table_props.cdc_max_chunk_size()?,
99                norm_level: table_props.cdc_norm_level()?,
100            })
101        } else {
102            None
103        };
104        let compression = parquet_compression(table_props.parquet_compression_codec()?)?;
105        let props = WriterProperties::builder()
106            .set_content_defined_chunking(cdc)
107            .set_compression(compression)
108            .set_max_row_group_bytes(Some(table_props.parquet_row_group_size_bytes()?))
109            .set_data_page_size_limit(table_props.parquet_page_size_bytes()?)
110            .set_data_page_row_count_limit(table_props.parquet_page_row_limit()?)
111            .set_dictionary_page_size_limit(table_props.parquet_dict_size_bytes()?)
112            .build();
113        Ok(Self::new_with_match_mode(props, schema, FieldMatchMode::Id))
114    }
115
116    /// Set the field match mode used to map Arrow fields to Iceberg fields.
117    ///
118    /// Defaults to [`FieldMatchMode::Id`]. Use [`FieldMatchMode::Name`] when the
119    /// incoming Arrow schema does not carry Iceberg field-id metadata.
120    pub fn with_match_mode(mut self, match_mode: FieldMatchMode) -> Self {
121        self.match_mode = match_mode;
122        self
123    }
124
125    /// Set the [`EncryptionManager`] used to encrypt written files.
126    pub fn with_encryption_manager(mut self, encryption_manager: Arc<EncryptionManager>) -> Self {
127        self.encryption_manager = Some(encryption_manager);
128        self
129    }
130}
131
132fn parquet_compression(codec: CompressionCodec) -> Result<Compression> {
133    let compression = match codec {
134        CompressionCodec::None => Compression::UNCOMPRESSED,
135        CompressionCodec::Snappy => Compression::SNAPPY,
136        CompressionCodec::Lzo => Compression::LZO,
137        CompressionCodec::Lz4 => Compression::LZ4,
138        CompressionCodec::Lz4Raw => Compression::LZ4_RAW,
139        CompressionCodec::Zstd(level) => {
140            let level =
141                ZstdLevel::try_new(level as i32).map_err(|e| invalid_level_error("zstd", e))?;
142            Compression::ZSTD(level)
143        }
144        CompressionCodec::Gzip(level) => {
145            let level =
146                GzipLevel::try_new(level as u32).map_err(|e| invalid_level_error("gzip", e))?;
147            Compression::GZIP(level)
148        }
149        CompressionCodec::Brotli(level) => {
150            let level =
151                BrotliLevel::try_new(level as u32).map_err(|e| invalid_level_error("brotli", e))?;
152            Compression::BROTLI(level)
153        }
154    };
155    Ok(compression)
156}
157
158fn invalid_level_error(codec: &str, source: impl Into<anyhow::Error>) -> Error {
159    invalid_data!("Invalid {codec} compression level").with_source(source)
160}
161
162impl FileWriterBuilder for ParquetWriterBuilder {
163    type R = ParquetWriter;
164
165    async fn build(&self, output_file: OutputFile) -> Result<Self::R> {
166        let key_metadata = self
167            .encryption_manager
168            .as_ref()
169            .map(|em| em.generate_key_metadata());
170        let writer_properties =
171            resolve_writer_properties(self.props.clone(), key_metadata.as_ref())?;
172        Ok(ParquetWriter {
173            schema: self.schema.clone(),
174            inner_writer: None,
175            writer_properties,
176            current_row_num: 0,
177            output_file,
178            nan_value_count_visitor: NanValueCountVisitor::new_with_match_mode(self.match_mode),
179            key_metadata,
180        })
181    }
182}
183
184/// A mapping from Parquet column path names to internal field id
185struct IndexByParquetPathName {
186    name_to_id: HashMap<String, i32>,
187
188    field_names: Vec<String>,
189
190    field_id: i32,
191}
192
193impl IndexByParquetPathName {
194    /// Creates a new, empty `IndexByParquetPathName`
195    pub fn new() -> Self {
196        Self {
197            name_to_id: HashMap::new(),
198            field_names: Vec::new(),
199            field_id: 0,
200        }
201    }
202
203    /// Retrieves the internal field ID
204    pub fn get(&self, name: &str) -> Option<&i32> {
205        self.name_to_id.get(name)
206    }
207}
208
209impl Default for IndexByParquetPathName {
210    fn default() -> Self {
211        Self::new()
212    }
213}
214
215impl SchemaVisitor for IndexByParquetPathName {
216    type T = ();
217
218    fn before_struct_field(&mut self, field: &NestedFieldRef) -> Result<()> {
219        self.field_names.push(field.name.to_string());
220        self.field_id = field.id;
221        Ok(())
222    }
223
224    fn after_struct_field(&mut self, _field: &NestedFieldRef) -> Result<()> {
225        self.field_names.pop();
226        Ok(())
227    }
228
229    fn before_list_element(&mut self, field: &NestedFieldRef) -> Result<()> {
230        self.field_names.push(format!("list.{}", field.name));
231        self.field_id = field.id;
232        Ok(())
233    }
234
235    fn after_list_element(&mut self, _field: &NestedFieldRef) -> Result<()> {
236        self.field_names.pop();
237        Ok(())
238    }
239
240    fn before_map_key(&mut self, field: &NestedFieldRef) -> Result<()> {
241        self.field_names
242            .push(format!("{DEFAULT_MAP_FIELD_NAME}.key"));
243        self.field_id = field.id;
244        Ok(())
245    }
246
247    fn after_map_key(&mut self, _field: &NestedFieldRef) -> Result<()> {
248        self.field_names.pop();
249        Ok(())
250    }
251
252    fn before_map_value(&mut self, field: &NestedFieldRef) -> Result<()> {
253        self.field_names
254            .push(format!("{DEFAULT_MAP_FIELD_NAME}.value"));
255        self.field_id = field.id;
256        Ok(())
257    }
258
259    fn after_map_value(&mut self, _field: &NestedFieldRef) -> Result<()> {
260        self.field_names.pop();
261        Ok(())
262    }
263
264    fn schema(&mut self, _schema: &Schema, _value: Self::T) -> Result<Self::T> {
265        Ok(())
266    }
267
268    fn field(&mut self, _field: &NestedFieldRef, _value: Self::T) -> Result<Self::T> {
269        Ok(())
270    }
271
272    fn r#struct(&mut self, _struct: &StructType, _results: Vec<Self::T>) -> Result<Self::T> {
273        Ok(())
274    }
275
276    fn list(&mut self, _list: &ListType, _value: Self::T) -> Result<Self::T> {
277        Ok(())
278    }
279
280    fn map(&mut self, _map: &MapType, _key_value: Self::T, _value: Self::T) -> Result<Self::T> {
281        Ok(())
282    }
283
284    fn primitive(&mut self, _p: &PrimitiveType) -> Result<Self::T> {
285        let full_name = self.field_names.iter().map(String::as_str).join(".");
286        let field_id = self.field_id;
287        if let Some(existing_field_id) = self.name_to_id.get(full_name.as_str()) {
288            return Err(invalid_data!(
289                "Invalid schema: multiple fields for name {full_name}: {field_id} and {existing_field_id}"
290            ));
291        } else {
292            self.name_to_id.insert(full_name, field_id);
293        }
294
295        Ok(())
296    }
297
298    fn variant(&mut self, _v: &VariantType) -> Result<Self::T> {
299        Err(Error::new(
300            ErrorKind::FeatureUnsupported,
301            "Writing variant columns to Parquet is not supported yet",
302        ))
303    }
304}
305
306/// `ParquetWriter`` is used to write arrow data into parquet file on storage.
307pub struct ParquetWriter {
308    schema: SchemaRef,
309    output_file: OutputFile,
310    inner_writer: Option<AsyncArrowWriter<AsyncFileWriter>>,
311    writer_properties: WriterProperties,
312    current_row_num: usize,
313    nan_value_count_visitor: NanValueCountVisitor,
314    key_metadata: Option<StandardKeyMetadata>,
315}
316
317/// Used to aggregate min and max value of each column.
318struct MinMaxColAggregator {
319    lower_bounds: HashMap<i32, Datum>,
320    upper_bounds: HashMap<i32, Datum>,
321    schema: SchemaRef,
322}
323
324impl MinMaxColAggregator {
325    /// Creates new and empty `MinMaxColAggregator`
326    fn new(schema: SchemaRef) -> Self {
327        Self {
328            lower_bounds: HashMap::new(),
329            upper_bounds: HashMap::new(),
330            schema,
331        }
332    }
333
334    fn update_state_min(&mut self, field_id: i32, datum: Datum) {
335        self.lower_bounds
336            .entry(field_id)
337            .and_modify(|e| {
338                if *e > datum {
339                    *e = datum.clone()
340                }
341            })
342            .or_insert(datum);
343    }
344
345    fn update_state_max(&mut self, field_id: i32, datum: Datum) {
346        self.upper_bounds
347            .entry(field_id)
348            .and_modify(|e| {
349                if *e < datum {
350                    *e = datum.clone()
351                }
352            })
353            .or_insert(datum);
354    }
355
356    /// Update statistics
357    fn update(&mut self, field_id: i32, value: &Statistics) -> Result<()> {
358        let Some(ty) = self
359            .schema
360            .field_by_id(field_id)
361            .map(|f| f.field_type.as_ref())
362        else {
363            // Following java implementation: https://github.com/apache/iceberg/blob/29a2c456353a6120b8c882ed2ab544975b168d7b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetUtil.java#L163
364            // Ignore the field if it is not in schema.
365            return Ok(());
366        };
367        let Type::Primitive(ty) = ty.clone() else {
368            return Err(Error::new(
369                ErrorKind::Unexpected,
370                format!("Composed type {ty} is not supported for min max aggregation."),
371            ));
372        };
373
374        if value.min_is_exact() {
375            let Some(min_datum) = get_parquet_stat_min_as_datum(&ty, value)? else {
376                return Err(Error::new(
377                    ErrorKind::Unexpected,
378                    format!("Statistics {value} is not match with field type {ty}."),
379                ));
380            };
381
382            self.update_state_min(field_id, min_datum);
383        }
384
385        if value.max_is_exact() {
386            let Some(max_datum) = get_parquet_stat_max_as_datum(&ty, value)? else {
387                return Err(Error::new(
388                    ErrorKind::Unexpected,
389                    format!("Statistics {value} is not match with field type {ty}."),
390                ));
391            };
392
393            self.update_state_max(field_id, max_datum);
394        }
395
396        Ok(())
397    }
398
399    /// Returns lower and upper bounds
400    fn produce(self) -> (HashMap<i32, Datum>, HashMap<i32, Datum>) {
401        (self.lower_bounds, self.upper_bounds)
402    }
403}
404
405impl ParquetWriter {
406    /// Converts parquet files to data files
407    #[allow(dead_code)]
408    pub(crate) async fn parquet_files_to_data_files(
409        file_io: &FileIO,
410        file_paths: Vec<String>,
411        table_metadata: &TableMetadata,
412    ) -> Result<Vec<DataFile>> {
413        // TODO: support adding to partitioned table
414        let mut data_files: Vec<DataFile> = Vec::new();
415
416        for file_path in file_paths {
417            let input_file = file_io.new_input(&file_path)?;
418            let file_metadata = input_file.metadata().await?;
419            let file_size_in_bytes = file_metadata.size as usize;
420            let reader = input_file.reader().await?;
421
422            let mut parquet_reader = ArrowFileReader::new(file_metadata, reader);
423            let parquet_metadata = parquet_reader
424                .get_metadata(None)
425                .await
426                .map_err(|err| invalid_data!("Error reading Parquet metadata: {err}"))?;
427            let mut builder = ParquetWriter::parquet_to_data_file_builder(
428                table_metadata.current_schema().clone(),
429                parquet_metadata,
430                file_size_in_bytes,
431                file_path,
432                // TODO: Implement nan_value_counts here
433                HashMap::new(),
434                None,
435            )?;
436            builder.partition_spec_id(table_metadata.default_partition_spec_id());
437            let data_file = builder.build().unwrap();
438            data_files.push(data_file);
439        }
440
441        Ok(data_files)
442    }
443
444    /// `ParquetMetadata` to data file builder
445    pub(crate) fn parquet_to_data_file_builder(
446        schema: SchemaRef,
447        metadata: Arc<ParquetMetaData>,
448        written_size: usize,
449        file_path: String,
450        nan_value_counts: HashMap<i32, u64>,
451        key_metadata: Option<StandardKeyMetadata>,
452    ) -> Result<DataFileBuilder> {
453        let index_by_parquet_path = {
454            let mut visitor = IndexByParquetPathName::new();
455            visit_schema(&schema, &mut visitor)?;
456            visitor
457        };
458
459        let (column_sizes, value_counts, null_value_counts, (lower_bounds, upper_bounds)) = {
460            let mut per_col_size: HashMap<i32, u64> = HashMap::new();
461            let mut per_col_val_num: HashMap<i32, u64> = HashMap::new();
462            let mut per_col_null_val_num: HashMap<i32, u64> = HashMap::new();
463            let mut min_max_agg = MinMaxColAggregator::new(schema);
464
465            for row_group in metadata.row_groups() {
466                for column_chunk_metadata in row_group.columns() {
467                    let parquet_path = column_chunk_metadata.column_descr().path().string();
468
469                    let Some(&field_id) = index_by_parquet_path.get(&parquet_path) else {
470                        continue;
471                    };
472
473                    *per_col_size.entry(field_id).or_insert(0) +=
474                        column_chunk_metadata.compressed_size() as u64;
475                    *per_col_val_num.entry(field_id).or_insert(0) +=
476                        column_chunk_metadata.num_values() as u64;
477
478                    if let Some(statistics) = column_chunk_metadata.statistics() {
479                        if let Some(null_count) = statistics.null_count_opt() {
480                            *per_col_null_val_num.entry(field_id).or_insert(0) += null_count;
481                        }
482
483                        min_max_agg.update(field_id, statistics)?;
484                    }
485                }
486            }
487            (
488                per_col_size,
489                per_col_val_num,
490                per_col_null_val_num,
491                min_max_agg.produce(),
492            )
493        };
494
495        let key_metadata = match key_metadata {
496            Some(m) => Some(m.encode()?.into_vec()),
497            None => None,
498        };
499
500        let mut builder = DataFileBuilder::default();
501        builder
502            .content(DataContentType::Data)
503            .file_path(file_path)
504            .file_format(DataFileFormat::Parquet)
505            .partition(Struct::empty())
506            .record_count(metadata.file_metadata().num_rows() as u64)
507            .file_size_in_bytes(written_size as u64)
508            .column_sizes(column_sizes)
509            .value_counts(value_counts)
510            .null_value_counts(null_value_counts)
511            .nan_value_counts(nan_value_counts)
512            // # NOTE:
513            // - We can ignore implementing distinct_counts due to this: https://lists.apache.org/thread/j52tsojv0x4bopxyzsp7m7bqt23n5fnd
514            .lower_bounds(lower_bounds)
515            .upper_bounds(upper_bounds)
516            .key_metadata(key_metadata)
517            .split_offsets(Some(
518                metadata
519                    .row_groups()
520                    .iter()
521                    .filter_map(|group| group.file_offset())
522                    .collect(),
523            ));
524
525        Ok(builder)
526    }
527
528    #[allow(dead_code)]
529    fn partition_value_from_bounds(
530        table_spec: Arc<PartitionSpec>,
531        lower_bounds: &HashMap<i32, Datum>,
532        upper_bounds: &HashMap<i32, Datum>,
533    ) -> Result<Struct> {
534        let mut partition_literals: Vec<Option<Literal>> = Vec::new();
535
536        for field in table_spec.fields() {
537            if let (Some(lower), Some(upper)) = (
538                lower_bounds.get(&field.source_id),
539                upper_bounds.get(&field.source_id),
540            ) {
541                if !field.transform.preserves_order() {
542                    return Err(invalid_data!(
543                        "cannot infer partition value for non linear partition field (needs to preserve order): {} with transform {}",
544                        field.name,
545                        field.transform
546                    ));
547                }
548
549                if lower != upper {
550                    return Err(invalid_data!(
551                        "multiple partition values for field {}: lower: {:?}, upper: {:?}",
552                        field.name,
553                        lower,
554                        upper
555                    ));
556                }
557
558                let transform_fn = create_transform_function(&field.transform)?;
559                let transform_literal =
560                    Literal::from(transform_fn.transform_literal_result(lower)?);
561
562                partition_literals.push(Some(transform_literal));
563            } else {
564                partition_literals.push(None);
565            }
566        }
567
568        let partition_struct = Struct::from_iter(partition_literals);
569
570        Ok(partition_struct)
571    }
572}
573
574fn resolve_writer_properties(
575    writer_properties: WriterProperties,
576    key_metadata: Option<&StandardKeyMetadata>,
577) -> Result<WriterProperties> {
578    let Some(key_metadata) = key_metadata else {
579        return Ok(writer_properties);
580    };
581
582    if writer_properties.file_encryption_properties().is_some() {
583        return Err(Error::new(
584            ErrorKind::Unexpected,
585            "Parquet writer properties already have file encryption properties set",
586        ));
587    }
588
589    let mut builder =
590        FileEncryptionProperties::builder(key_metadata.encryption_key().as_bytes().to_vec());
591    if let Some(aad) = key_metadata.aad_prefix() {
592        builder = builder.with_aad_prefix(aad.to_vec());
593    }
594    let file_encryption_properties = builder.build().map_err(|e| {
595        Error::new(
596            ErrorKind::Unexpected,
597            "Failed to build parquet file encryption properties",
598        )
599        .with_source(e)
600    })?;
601
602    Ok(writer_properties
603        .into_builder()
604        .with_file_encryption_properties(file_encryption_properties)
605        .build())
606}
607
608impl FileWriter for ParquetWriter {
609    async fn write(&mut self, batch: &arrow_array::RecordBatch) -> Result<()> {
610        // Skip empty batch
611        if batch.num_rows() == 0 {
612            return Ok(());
613        }
614
615        self.current_row_num += batch.num_rows();
616
617        let batch_c = batch.clone();
618        self.nan_value_count_visitor
619            .compute(self.schema.clone(), batch_c)?;
620
621        // Lazy initialize the writer
622        let writer = if let Some(writer) = &mut self.inner_writer {
623            writer
624        } else {
625            let arrow_schema: ArrowSchemaRef = Arc::new(self.schema.as_ref().try_into()?);
626            let inner_writer = self.output_file.writer().await?;
627            let async_writer = AsyncFileWriter::new(inner_writer);
628            let writer = AsyncArrowWriter::try_new(
629                async_writer,
630                arrow_schema.clone(),
631                Some(self.writer_properties.clone()),
632            )
633            .map_err(|err| {
634                Error::new(ErrorKind::Unexpected, "Failed to build parquet writer.")
635                    .with_source(err)
636            })?;
637            self.inner_writer = Some(writer);
638            self.inner_writer.as_mut().unwrap()
639        };
640
641        writer.write(batch).await.map_err(|err| {
642            Error::new(
643                ErrorKind::Unexpected,
644                "Failed to write using parquet writer.",
645            )
646            .with_source(err)
647        })?;
648
649        Ok(())
650    }
651
652    async fn close(mut self) -> Result<Vec<DataFileBuilder>> {
653        let mut writer = match self.inner_writer.take() {
654            Some(writer) => writer,
655            None => return Ok(vec![]),
656        };
657
658        let metadata = writer.finish().await.map_err(|err| {
659            Error::new(ErrorKind::Unexpected, "Failed to finish parquet writer.").with_source(err)
660        })?;
661
662        let written_size = writer.bytes_written();
663
664        if self.current_row_num == 0 {
665            self.output_file.delete().await.map_err(|err| {
666                Error::new(
667                    ErrorKind::Unexpected,
668                    "Failed to delete empty parquet file.",
669                )
670                .with_source(err)
671            })?;
672            Ok(vec![])
673        } else {
674            let parquet_metadata = Arc::new(metadata);
675
676            Ok(vec![Self::parquet_to_data_file_builder(
677                self.schema,
678                parquet_metadata,
679                written_size,
680                self.output_file.location().to_string(),
681                self.nan_value_count_visitor.nan_value_counts,
682                self.key_metadata,
683            )?])
684        }
685    }
686}
687
688impl CurrentFileStatus for ParquetWriter {
689    fn current_file_path(&self) -> String {
690        self.output_file.location().to_string()
691    }
692
693    fn current_row_num(&self) -> usize {
694        self.current_row_num
695    }
696
697    fn current_written_size(&self) -> usize {
698        if let Some(inner) = self.inner_writer.as_ref() {
699            // inner/AsyncArrowWriter contains sync and async writers
700            // written size = bytes flushed to inner's async writer + bytes buffered in the inner's sync writer
701            inner.bytes_written() + inner.in_progress_size()
702        } else {
703            // inner writer is not initialized yet
704            0
705        }
706    }
707}
708
709/// AsyncFileWriter is a wrapper of FileWrite to make it compatible with tokio::io::AsyncWrite.
710///
711/// # NOTES
712///
713/// We keep this wrapper been used inside only.
714struct AsyncFileWriter(Box<dyn FileWrite>);
715
716impl AsyncFileWriter {
717    /// Create a new `AsyncFileWriter` with the given writer.
718    pub fn new(writer: Box<dyn FileWrite>) -> Self {
719        Self(writer)
720    }
721}
722
723impl ArrowAsyncFileWriter for AsyncFileWriter {
724    fn write(&mut self, bs: Bytes) -> BoxFuture<'_, parquet::errors::Result<()>> {
725        Box::pin(async {
726            self.0
727                .write(bs)
728                .await
729                .map_err(|err| parquet::errors::ParquetError::External(Box::new(err)))
730        })
731    }
732
733    fn complete(&mut self) -> BoxFuture<'_, parquet::errors::Result<()>> {
734        Box::pin(async {
735            // TODO(encryption): retain the stored file size in data-file key metadata.
736            self.0
737                .close()
738                .await
739                .map(|_| ())
740                .map_err(|err| parquet::errors::ParquetError::External(Box::new(err)))
741        })
742    }
743}
744
745#[cfg(test)]
746mod tests {
747    use std::collections::HashMap;
748    use std::sync::Arc;
749
750    use anyhow::Result;
751    use arrow_array::builder::{Float32Builder, Int32Builder, MapBuilder};
752    use arrow_array::types::{Float32Type, Int64Type};
753    use arrow_array::{
754        Array, ArrayRef, BooleanArray, Decimal128Array, Float32Array, Float64Array, Int32Array,
755        Int64Array, ListArray, MapArray, RecordBatch, StructArray,
756    };
757    use arrow_schema::{DataType, Field, Fields, SchemaRef as ArrowSchemaRef};
758    use arrow_select::concat::concat_batches;
759    use futures::TryStreamExt;
760    use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
761    use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
762    use parquet::basic::{BrotliLevel, Compression, GzipLevel, ZstdLevel};
763    use parquet::file::statistics::ValueStatistics;
764    use parquet::schema::types::ColumnPath;
765    use tempfile::TempDir;
766    use uuid::Uuid;
767
768    use super::*;
769    use crate::Runtime;
770    use crate::arrow::{ArrowReaderBuilder, schema_to_arrow_schema};
771    use crate::io::FileIO;
772    use crate::scan::{FileScanTask, FileScanTaskStream};
773    use crate::spec::decimal_utils::{decimal_mantissa, decimal_new, decimal_scale};
774    use crate::spec::{PrimitiveLiteral, Struct, *};
775    use crate::test_utils::make_encryption_manager;
776    use crate::writer::file_writer::location_generator::{
777        DefaultFileNameGenerator, DefaultLocationGenerator, FileNameGenerator, LocationGenerator,
778    };
779    use crate::writer::tests::check_parquet_data_file;
780
781    fn schema_for_all_type() -> Schema {
782        Schema::builder()
783            .with_schema_id(1)
784            .with_fields(vec![
785                NestedField::optional(0, "boolean", Type::Primitive(PrimitiveType::Boolean)).into(),
786                NestedField::optional(1, "int", Type::Primitive(PrimitiveType::Int)).into(),
787                NestedField::optional(2, "long", Type::Primitive(PrimitiveType::Long)).into(),
788                NestedField::optional(3, "float", Type::Primitive(PrimitiveType::Float)).into(),
789                NestedField::optional(4, "double", Type::Primitive(PrimitiveType::Double)).into(),
790                NestedField::optional(5, "string", Type::Primitive(PrimitiveType::String)).into(),
791                NestedField::optional(6, "binary", Type::Primitive(PrimitiveType::Binary)).into(),
792                NestedField::optional(7, "date", Type::Primitive(PrimitiveType::Date)).into(),
793                NestedField::optional(8, "time", Type::Primitive(PrimitiveType::Time)).into(),
794                NestedField::optional(9, "timestamp", Type::Primitive(PrimitiveType::Timestamp))
795                    .into(),
796                NestedField::optional(
797                    10,
798                    "timestamptz",
799                    Type::Primitive(PrimitiveType::Timestamptz),
800                )
801                .into(),
802                NestedField::optional(
803                    11,
804                    "timestamp_ns",
805                    Type::Primitive(PrimitiveType::TimestampNs),
806                )
807                .into(),
808                NestedField::optional(
809                    12,
810                    "timestamptz_ns",
811                    Type::Primitive(PrimitiveType::TimestamptzNs),
812                )
813                .into(),
814                NestedField::optional(
815                    13,
816                    "decimal",
817                    Type::Primitive(PrimitiveType::Decimal {
818                        precision: 10,
819                        scale: 5,
820                    }),
821                )
822                .into(),
823                NestedField::optional(14, "uuid", Type::Primitive(PrimitiveType::Uuid)).into(),
824                NestedField::optional(15, "fixed", Type::Primitive(PrimitiveType::Fixed(10)))
825                    .into(),
826                // Parquet Statistics will use different representation for Decimal with precision 38 and scale 5,
827                // so we need to add a new field for it.
828                NestedField::optional(
829                    16,
830                    "decimal_38",
831                    Type::Primitive(PrimitiveType::Decimal {
832                        precision: 38,
833                        scale: 5,
834                    }),
835                )
836                .into(),
837            ])
838            .build()
839            .unwrap()
840    }
841
842    fn nested_schema_for_test() -> Schema {
843        // Int, Struct(Int,Int), String, List(Int), Struct(Struct(Int)), Map(String, List(Int))
844        Schema::builder()
845            .with_schema_id(1)
846            .with_fields(vec![
847                NestedField::required(0, "col0", Type::Primitive(PrimitiveType::Long)).into(),
848                NestedField::required(
849                    1,
850                    "col1",
851                    Type::Struct(StructType::new(vec![
852                        NestedField::required(5, "col_1_5", Type::Primitive(PrimitiveType::Long))
853                            .into(),
854                        NestedField::required(6, "col_1_6", Type::Primitive(PrimitiveType::Long))
855                            .into(),
856                    ])),
857                )
858                .into(),
859                NestedField::required(2, "col2", Type::Primitive(PrimitiveType::String)).into(),
860                NestedField::required(
861                    3,
862                    "col3",
863                    Type::List(ListType::new(
864                        NestedField::required(7, "element", Type::Primitive(PrimitiveType::Long))
865                            .into(),
866                    )),
867                )
868                .into(),
869                NestedField::required(
870                    4,
871                    "col4",
872                    Type::Struct(StructType::new(vec![
873                        NestedField::required(
874                            8,
875                            "col_4_8",
876                            Type::Struct(StructType::new(vec![
877                                NestedField::required(
878                                    9,
879                                    "col_4_8_9",
880                                    Type::Primitive(PrimitiveType::Long),
881                                )
882                                .into(),
883                            ])),
884                        )
885                        .into(),
886                    ])),
887                )
888                .into(),
889                NestedField::required(
890                    10,
891                    "col5",
892                    Type::Map(MapType::new(
893                        NestedField::required(11, "key", Type::Primitive(PrimitiveType::String))
894                            .into(),
895                        NestedField::required(
896                            12,
897                            "value",
898                            Type::List(ListType::new(
899                                NestedField::required(
900                                    13,
901                                    "item",
902                                    Type::Primitive(PrimitiveType::Long),
903                                )
904                                .into(),
905                            )),
906                        )
907                        .into(),
908                    )),
909                )
910                .into(),
911            ])
912            .build()
913            .unwrap()
914    }
915
916    #[tokio::test]
917    async fn test_index_by_parquet_path() {
918        let expect = HashMap::from([
919            ("col0".to_string(), 0),
920            ("col1.col_1_5".to_string(), 5),
921            ("col1.col_1_6".to_string(), 6),
922            ("col2".to_string(), 2),
923            ("col3.list.element".to_string(), 7),
924            ("col4.col_4_8.col_4_8_9".to_string(), 9),
925            ("col5.key_value.key".to_string(), 11),
926            ("col5.key_value.value.list.item".to_string(), 13),
927        ]);
928        let mut visitor = IndexByParquetPathName::new();
929        visit_schema(&nested_schema_for_test(), &mut visitor).unwrap();
930        assert_eq!(visitor.name_to_id, expect);
931    }
932
933    #[test]
934    fn test_index_by_parquet_path_variant_is_unsupported() {
935        // Writing variant columns to Parquet is not supported yet; indexing a schema that
936        // contains one must error rather than silently miss-map columns.
937        let schema = Schema::builder()
938            .with_fields(vec![
939                NestedField::optional(1, "v", Type::Variant(VariantType)).into(),
940            ])
941            .build()
942            .unwrap();
943        let mut visitor = IndexByParquetPathName::new();
944        let err = visit_schema(&schema, &mut visitor).unwrap_err();
945        assert_eq!(err.kind(), ErrorKind::FeatureUnsupported, "{err}");
946    }
947
948    #[tokio::test]
949    async fn test_parquet_writer() -> Result<()> {
950        let temp_dir = TempDir::new().unwrap();
951        let file_io = FileIO::new_with_fs();
952        let location_gen = DefaultLocationGenerator::with_data_location(
953            temp_dir.path().to_str().unwrap().to_string(),
954        );
955        let file_name_gen =
956            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
957
958        // prepare data
959        let schema = {
960            let fields =
961                vec![
962                    Field::new("col", DataType::Int64, true).with_metadata(HashMap::from([(
963                        PARQUET_FIELD_ID_META_KEY.to_string(),
964                        "0".to_string(),
965                    )])),
966                ];
967            Arc::new(arrow_schema::Schema::new(fields))
968        };
969        let col = Arc::new(Int64Array::from_iter_values(0..1024)) as ArrayRef;
970        let null_col = Arc::new(Int64Array::new_null(1024)) as ArrayRef;
971        let to_write = RecordBatch::try_new(schema.clone(), vec![col]).unwrap();
972        let to_write_null = RecordBatch::try_new(schema.clone(), vec![null_col]).unwrap();
973
974        let output_file = file_io.new_output(
975            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
976        )?;
977
978        // write data
979        let mut pw = ParquetWriterBuilder::new(
980            WriterProperties::builder()
981                .set_max_row_group_row_count(Some(128))
982                .build(),
983            Arc::new(to_write.schema().as_ref().try_into().unwrap()),
984        )
985        .build(output_file)
986        .await?;
987        pw.write(&to_write).await?;
988        pw.write(&to_write_null).await?;
989        let res = pw.close().await?;
990        assert_eq!(res.len(), 1);
991        let data_file = res
992            .into_iter()
993            .next()
994            .unwrap()
995            // Put dummy field for build successfully.
996            .content(DataContentType::Data)
997            .partition(Struct::empty())
998            .partition_spec_id(0)
999            .build()
1000            .unwrap();
1001
1002        // check data file
1003        assert_eq!(data_file.record_count(), 2048);
1004        assert_eq!(*data_file.value_counts(), HashMap::from([(0, 2048)]));
1005        assert_eq!(
1006            *data_file.lower_bounds(),
1007            HashMap::from([(0, Datum::long(0))])
1008        );
1009        assert_eq!(
1010            *data_file.upper_bounds(),
1011            HashMap::from([(0, Datum::long(1023))])
1012        );
1013        assert_eq!(*data_file.null_value_counts(), HashMap::from([(0, 1024)]));
1014
1015        // check the written file
1016        let expect_batch = concat_batches(&schema, vec![&to_write, &to_write_null]).unwrap();
1017        check_parquet_data_file(&file_io, &data_file, &expect_batch).await;
1018
1019        Ok(())
1020    }
1021
1022    #[tokio::test]
1023    async fn test_parquet_writer_encrypted_roundtrip() -> Result<()> {
1024        let temp_dir = TempDir::new().unwrap();
1025        let file_io = FileIO::new_with_fs();
1026        let location_gen = DefaultLocationGenerator::with_data_location(
1027            temp_dir.path().to_str().unwrap().to_string(),
1028        );
1029        let file_name_gen =
1030            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1031
1032        let arrow_schema = Arc::new(arrow_schema::Schema::new(vec![
1033            Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1034                PARQUET_FIELD_ID_META_KEY.to_string(),
1035                "1".to_string(),
1036            )])),
1037        ]));
1038        let iceberg_schema: SchemaRef = Arc::new(arrow_schema.as_ref().try_into().unwrap());
1039        let batch = RecordBatch::try_new(arrow_schema, vec![Arc::new(Int32Array::from(vec![
1040            10, 20, 30,
1041        ])) as ArrayRef])
1042        .unwrap();
1043
1044        let output_file = file_io.new_output(
1045            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1046        )?;
1047        let raw_properties = HashMap::from([(
1048            TableProperties::PROPERTY_ENCRYPTION_KEY_ID.to_string(),
1049            "test-key".to_string(),
1050        )]);
1051        let table_properties = TableProperties::new(&raw_properties);
1052        let mut parquet_writer =
1053            ParquetWriterBuilder::from_table_properties(&table_properties, iceberg_schema.clone())?
1054                .with_encryption_manager(make_encryption_manager("test-key"))
1055                .build(output_file)
1056                .await?;
1057        parquet_writer.write(&batch).await?;
1058        let data_file = parquet_writer
1059            .close()
1060            .await?
1061            .into_iter()
1062            .next()
1063            .unwrap()
1064            .content(DataContentType::Data)
1065            .partition(Struct::empty())
1066            .partition_spec_id(0)
1067            .build()
1068            .unwrap();
1069
1070        let raw = file_io
1071            .new_input(data_file.file_path.clone())?
1072            .read()
1073            .await?;
1074        assert!(
1075            ParquetRecordBatchReaderBuilder::try_new(raw).is_err(),
1076            "an encrypted parquet file must not be readable without decryption"
1077        );
1078
1079        let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
1080        let task = FileScanTask::builder()
1081            .with_file_size_in_bytes(data_file.file_size_in_bytes())
1082            .with_start(0)
1083            .with_length(0)
1084            .with_data_file_path(data_file.file_path.clone())
1085            .with_data_file_format(DataFileFormat::Parquet)
1086            .with_schema(iceberg_schema)
1087            .with_project_field_ids(vec![1])
1088            .with_case_sensitive(false)
1089            .with_key_metadata(data_file.key_metadata().map(Box::from))
1090            .build()
1091            .unwrap();
1092        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1093        let batches: Vec<RecordBatch> = reader
1094            .read(tasks)
1095            .unwrap()
1096            .stream()
1097            .try_collect()
1098            .await
1099            .unwrap();
1100
1101        assert_eq!(batches.len(), 1);
1102        let ids = batches[0]
1103            .column(0)
1104            .as_any()
1105            .downcast_ref::<Int32Array>()
1106            .unwrap();
1107        assert_eq!(ids.values(), &[10, 20, 30]);
1108
1109        Ok(())
1110    }
1111
1112    #[tokio::test]
1113    async fn test_parquet_writer_with_complex_schema() -> Result<()> {
1114        let temp_dir = TempDir::new().unwrap();
1115        let file_io = FileIO::new_with_fs();
1116        let location_gen = DefaultLocationGenerator::with_data_location(
1117            temp_dir.path().to_str().unwrap().to_string(),
1118        );
1119        let file_name_gen =
1120            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1121
1122        // prepare data
1123        let schema = nested_schema_for_test();
1124        let arrow_schema: ArrowSchemaRef = Arc::new((&schema).try_into().unwrap());
1125        let col0 = Arc::new(Int64Array::from_iter_values(0..1024)) as ArrayRef;
1126        let col1 = Arc::new(StructArray::new(
1127            {
1128                if let DataType::Struct(fields) = arrow_schema.field(1).data_type() {
1129                    fields.clone()
1130                } else {
1131                    unreachable!()
1132                }
1133            },
1134            vec![
1135                Arc::new(Int64Array::from_iter_values(0..1024)),
1136                Arc::new(Int64Array::from_iter_values(0..1024)),
1137            ],
1138            None,
1139        ));
1140        let col2 = Arc::new(arrow_array::StringArray::from_iter_values(
1141            (0..1024).map(|n| n.to_string()),
1142        )) as ArrayRef;
1143        let col3 = Arc::new({
1144            let list_parts = ListArray::from_iter_primitive::<Int64Type, _, _>(
1145                (0..1024).map(|n| Some(vec![Some(n)])),
1146            )
1147            .into_parts();
1148            ListArray::new(
1149                {
1150                    if let DataType::List(field) = arrow_schema.field(3).data_type() {
1151                        field.clone()
1152                    } else {
1153                        unreachable!()
1154                    }
1155                },
1156                list_parts.1,
1157                list_parts.2,
1158                list_parts.3,
1159            )
1160        }) as ArrayRef;
1161        let col4 = Arc::new(StructArray::new(
1162            {
1163                if let DataType::Struct(fields) = arrow_schema.field(4).data_type() {
1164                    fields.clone()
1165                } else {
1166                    unreachable!()
1167                }
1168            },
1169            vec![Arc::new(StructArray::new(
1170                {
1171                    if let DataType::Struct(fields) = arrow_schema.field(4).data_type() {
1172                        if let DataType::Struct(fields) = fields[0].data_type() {
1173                            fields.clone()
1174                        } else {
1175                            unreachable!()
1176                        }
1177                    } else {
1178                        unreachable!()
1179                    }
1180                },
1181                vec![Arc::new(Int64Array::from_iter_values(0..1024))],
1182                None,
1183            ))],
1184            None,
1185        ));
1186        let col5 = Arc::new({
1187            let mut map_array_builder = MapBuilder::new(
1188                None,
1189                arrow_array::builder::StringBuilder::new(),
1190                arrow_array::builder::ListBuilder::new(arrow_array::builder::PrimitiveBuilder::<
1191                    Int64Type,
1192                >::new()),
1193            );
1194            for i in 0..1024 {
1195                map_array_builder.keys().append_value(i.to_string());
1196                map_array_builder
1197                    .values()
1198                    .append_value(vec![Some(i as i64); i + 1]);
1199                map_array_builder.append(true)?;
1200            }
1201            let (_, offset_buffer, struct_array, null_buffer, ordered) =
1202                map_array_builder.finish().into_parts();
1203            let struct_array = {
1204                let (_, mut arrays, nulls) = struct_array.into_parts();
1205                let list_array = {
1206                    let list_array = arrays[1]
1207                        .as_any()
1208                        .downcast_ref::<ListArray>()
1209                        .unwrap()
1210                        .clone();
1211                    let (_, offsets, array, nulls) = list_array.into_parts();
1212                    let list_field = {
1213                        if let DataType::Map(map_field, _) = arrow_schema.field(5).data_type() {
1214                            if let DataType::Struct(fields) = map_field.data_type() {
1215                                if let DataType::List(list_field) = fields[1].data_type() {
1216                                    list_field.clone()
1217                                } else {
1218                                    unreachable!()
1219                                }
1220                            } else {
1221                                unreachable!()
1222                            }
1223                        } else {
1224                            unreachable!()
1225                        }
1226                    };
1227                    ListArray::new(list_field, offsets, array, nulls)
1228                };
1229                arrays[1] = Arc::new(list_array) as ArrayRef;
1230                StructArray::new(
1231                    {
1232                        if let DataType::Map(map_field, _) = arrow_schema.field(5).data_type() {
1233                            if let DataType::Struct(fields) = map_field.data_type() {
1234                                fields.clone()
1235                            } else {
1236                                unreachable!()
1237                            }
1238                        } else {
1239                            unreachable!()
1240                        }
1241                    },
1242                    arrays,
1243                    nulls,
1244                )
1245            };
1246            MapArray::new(
1247                {
1248                    if let DataType::Map(map_field, _) = arrow_schema.field(5).data_type() {
1249                        map_field.clone()
1250                    } else {
1251                        unreachable!()
1252                    }
1253                },
1254                offset_buffer,
1255                struct_array,
1256                null_buffer,
1257                ordered,
1258            )
1259        }) as ArrayRef;
1260        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![
1261            col0, col1, col2, col3, col4, col5,
1262        ])
1263        .unwrap();
1264        let output_file = file_io.new_output(
1265            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1266        )?;
1267
1268        // write data
1269        let mut pw =
1270            ParquetWriterBuilder::new(WriterProperties::builder().build(), Arc::new(schema))
1271                .build(output_file)
1272                .await?;
1273        pw.write(&to_write).await?;
1274        let res = pw.close().await?;
1275        assert_eq!(res.len(), 1);
1276        let data_file = res
1277            .into_iter()
1278            .next()
1279            .unwrap()
1280            // Put dummy field for build successfully.
1281            .content(DataContentType::Data)
1282            .partition(Struct::empty())
1283            .partition_spec_id(0)
1284            .build()
1285            .unwrap();
1286
1287        // check data file
1288        assert_eq!(data_file.record_count(), 1024);
1289        assert_eq!(
1290            *data_file.value_counts(),
1291            HashMap::from([
1292                (0, 1024),
1293                (5, 1024),
1294                (6, 1024),
1295                (2, 1024),
1296                (7, 1024),
1297                (9, 1024),
1298                (11, 1024),
1299                (13, (1..1025).sum()),
1300            ])
1301        );
1302        assert_eq!(
1303            *data_file.lower_bounds(),
1304            HashMap::from([
1305                (0, Datum::long(0)),
1306                (5, Datum::long(0)),
1307                (6, Datum::long(0)),
1308                (2, Datum::string("0")),
1309                (7, Datum::long(0)),
1310                (9, Datum::long(0)),
1311                (11, Datum::string("0")),
1312                (13, Datum::long(0))
1313            ])
1314        );
1315        assert_eq!(
1316            *data_file.upper_bounds(),
1317            HashMap::from([
1318                (0, Datum::long(1023)),
1319                (5, Datum::long(1023)),
1320                (6, Datum::long(1023)),
1321                (2, Datum::string("999")),
1322                (7, Datum::long(1023)),
1323                (9, Datum::long(1023)),
1324                (11, Datum::string("999")),
1325                (13, Datum::long(1023))
1326            ])
1327        );
1328
1329        // check the written file
1330        check_parquet_data_file(&file_io, &data_file, &to_write).await;
1331
1332        Ok(())
1333    }
1334
1335    #[tokio::test]
1336    async fn test_all_type_for_write() -> Result<()> {
1337        let temp_dir = TempDir::new().unwrap();
1338        let file_io = FileIO::new_with_fs();
1339        let location_gen = DefaultLocationGenerator::with_data_location(
1340            temp_dir.path().to_str().unwrap().to_string(),
1341        );
1342        let file_name_gen =
1343            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1344
1345        // prepare data
1346        // generate iceberg schema for all type
1347        let schema = schema_for_all_type();
1348        let arrow_schema: ArrowSchemaRef = Arc::new((&schema).try_into().unwrap());
1349        let col0 = Arc::new(BooleanArray::from(vec![
1350            Some(true),
1351            Some(false),
1352            None,
1353            Some(true),
1354        ])) as ArrayRef;
1355        let col1 = Arc::new(Int32Array::from(vec![Some(1), Some(2), None, Some(4)])) as ArrayRef;
1356        let col2 = Arc::new(Int64Array::from(vec![Some(1), Some(2), None, Some(4)])) as ArrayRef;
1357        let col3 = Arc::new(Float32Array::from(vec![
1358            Some(0.5),
1359            Some(2.0),
1360            None,
1361            Some(3.5),
1362        ])) as ArrayRef;
1363        let col4 = Arc::new(Float64Array::from(vec![
1364            Some(0.5),
1365            Some(2.0),
1366            None,
1367            Some(3.5),
1368        ])) as ArrayRef;
1369        let col5 = Arc::new(arrow_array::StringArray::from(vec![
1370            Some("a"),
1371            Some("b"),
1372            None,
1373            Some("d"),
1374        ])) as ArrayRef;
1375        let col6 = Arc::new(arrow_array::LargeBinaryArray::from_opt_vec(vec![
1376            Some(b"one"),
1377            None,
1378            Some(b""),
1379            Some(b"zzzz"),
1380        ])) as ArrayRef;
1381        let col7 = Arc::new(arrow_array::Date32Array::from(vec![
1382            Some(0),
1383            Some(1),
1384            None,
1385            Some(3),
1386        ])) as ArrayRef;
1387        let col8 = Arc::new(arrow_array::Time64MicrosecondArray::from(vec![
1388            Some(0),
1389            Some(1),
1390            None,
1391            Some(3),
1392        ])) as ArrayRef;
1393        let col9 = Arc::new(arrow_array::TimestampMicrosecondArray::from(vec![
1394            Some(0),
1395            Some(1),
1396            None,
1397            Some(3),
1398        ])) as ArrayRef;
1399        let col10 = Arc::new(
1400            arrow_array::TimestampMicrosecondArray::from(vec![Some(0), Some(1), None, Some(3)])
1401                .with_timezone_utc(),
1402        ) as ArrayRef;
1403        let col11 = Arc::new(arrow_array::TimestampNanosecondArray::from(vec![
1404            Some(0),
1405            Some(1),
1406            None,
1407            Some(3),
1408        ])) as ArrayRef;
1409        let col12 = Arc::new(
1410            arrow_array::TimestampNanosecondArray::from(vec![Some(0), Some(1), None, Some(3)])
1411                .with_timezone_utc(),
1412        ) as ArrayRef;
1413        let col13 = Arc::new(
1414            Decimal128Array::from(vec![Some(1), Some(2), None, Some(100)])
1415                .with_precision_and_scale(10, 5)
1416                .unwrap(),
1417        ) as ArrayRef;
1418        let col14 = Arc::new(
1419            arrow_array::FixedSizeBinaryArray::try_from_sparse_iter_with_size(
1420                vec![
1421                    Some(Uuid::from_u128(0).as_bytes().to_vec()),
1422                    Some(Uuid::from_u128(1).as_bytes().to_vec()),
1423                    None,
1424                    Some(Uuid::from_u128(3).as_bytes().to_vec()),
1425                ]
1426                .into_iter(),
1427                16,
1428            )
1429            .unwrap(),
1430        ) as ArrayRef;
1431        let col15 = Arc::new(
1432            arrow_array::FixedSizeBinaryArray::try_from_sparse_iter_with_size(
1433                vec![
1434                    Some(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10]),
1435                    Some(vec![11, 12, 13, 14, 15, 16, 17, 18, 19, 20]),
1436                    None,
1437                    Some(vec![21, 22, 23, 24, 25, 26, 27, 28, 29, 30]),
1438                ]
1439                .into_iter(),
1440                10,
1441            )
1442            .unwrap(),
1443        ) as ArrayRef;
1444        let col16 = Arc::new(
1445            Decimal128Array::from(vec![Some(1), Some(2), None, Some(100)])
1446                .with_precision_and_scale(38, 5)
1447                .unwrap(),
1448        ) as ArrayRef;
1449        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![
1450            col0, col1, col2, col3, col4, col5, col6, col7, col8, col9, col10, col11, col12, col13,
1451            col14, col15, col16,
1452        ])
1453        .unwrap();
1454        let output_file = file_io.new_output(
1455            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1456        )?;
1457
1458        // write data
1459        let mut pw =
1460            ParquetWriterBuilder::new(WriterProperties::builder().build(), Arc::new(schema))
1461                .build(output_file)
1462                .await?;
1463        pw.write(&to_write).await?;
1464        let res = pw.close().await?;
1465        assert_eq!(res.len(), 1);
1466        let data_file = res
1467            .into_iter()
1468            .next()
1469            .unwrap()
1470            // Put dummy field for build successfully.
1471            .content(DataContentType::Data)
1472            .partition(Struct::empty())
1473            .partition_spec_id(0)
1474            .build()
1475            .unwrap();
1476
1477        // check data file
1478        assert_eq!(data_file.record_count(), 4);
1479        assert!(data_file.value_counts().iter().all(|(_, &v)| { v == 4 }));
1480        assert!(
1481            data_file
1482                .null_value_counts()
1483                .iter()
1484                .all(|(_, &v)| { v == 1 })
1485        );
1486        assert_eq!(
1487            *data_file.lower_bounds(),
1488            HashMap::from([
1489                (0, Datum::bool(false)),
1490                (1, Datum::int(1)),
1491                (2, Datum::long(1)),
1492                (3, Datum::float(0.5_f32)),
1493                (4, Datum::double(0.5)),
1494                (5, Datum::string("a")),
1495                (6, Datum::binary(vec![])),
1496                (7, Datum::date(0)),
1497                (8, Datum::time_micros(0).unwrap()),
1498                (9, Datum::timestamp_micros(0)),
1499                (10, Datum::timestamptz_micros(0)),
1500                (11, Datum::timestamp_nanos(0)),
1501                (12, Datum::timestamptz_nanos(0)),
1502                (
1503                    13,
1504                    Datum::new(
1505                        PrimitiveType::Decimal {
1506                            precision: 10,
1507                            scale: 5
1508                        },
1509                        PrimitiveLiteral::Int128(1)
1510                    )
1511                ),
1512                (14, Datum::uuid(Uuid::from_u128(0))),
1513                (15, Datum::fixed(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10])),
1514                (
1515                    16,
1516                    Datum::new(
1517                        PrimitiveType::Decimal {
1518                            precision: 38,
1519                            scale: 5
1520                        },
1521                        PrimitiveLiteral::Int128(1)
1522                    )
1523                ),
1524            ])
1525        );
1526        assert_eq!(
1527            *data_file.upper_bounds(),
1528            HashMap::from([
1529                (0, Datum::bool(true)),
1530                (1, Datum::int(4)),
1531                (2, Datum::long(4)),
1532                (3, Datum::float(3.5_f32)),
1533                (4, Datum::double(3.5)),
1534                (5, Datum::string("d")),
1535                (6, Datum::binary(vec![122, 122, 122, 122])),
1536                (7, Datum::date(3)),
1537                (8, Datum::time_micros(3).unwrap()),
1538                (9, Datum::timestamp_micros(3)),
1539                (10, Datum::timestamptz_micros(3)),
1540                (11, Datum::timestamp_nanos(3)),
1541                (12, Datum::timestamptz_nanos(3)),
1542                (
1543                    13,
1544                    Datum::new(
1545                        PrimitiveType::Decimal {
1546                            precision: 10,
1547                            scale: 5
1548                        },
1549                        PrimitiveLiteral::Int128(100)
1550                    )
1551                ),
1552                (14, Datum::uuid(Uuid::from_u128(3))),
1553                (
1554                    15,
1555                    Datum::fixed(vec![21, 22, 23, 24, 25, 26, 27, 28, 29, 30])
1556                ),
1557                (
1558                    16,
1559                    Datum::new(
1560                        PrimitiveType::Decimal {
1561                            precision: 38,
1562                            scale: 5
1563                        },
1564                        PrimitiveLiteral::Int128(100)
1565                    )
1566                ),
1567            ])
1568        );
1569
1570        // check the written file
1571        check_parquet_data_file(&file_io, &data_file, &to_write).await;
1572
1573        Ok(())
1574    }
1575
1576    #[tokio::test]
1577    async fn test_decimal_bound() -> Result<()> {
1578        let temp_dir = TempDir::new().unwrap();
1579        let file_io = FileIO::new_with_fs();
1580        let location_gen = DefaultLocationGenerator::with_data_location(
1581            temp_dir.path().to_str().unwrap().to_string(),
1582        );
1583        let file_name_gen =
1584            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1585
1586        // test 1.1 and 2.2
1587        let schema = Arc::new(
1588            Schema::builder()
1589                .with_fields(vec![
1590                    NestedField::optional(
1591                        0,
1592                        "decimal",
1593                        Type::Primitive(PrimitiveType::Decimal {
1594                            precision: 28,
1595                            scale: 10,
1596                        }),
1597                    )
1598                    .into(),
1599                ])
1600                .build()
1601                .unwrap(),
1602        );
1603        let arrow_schema: ArrowSchemaRef = Arc::new(schema_to_arrow_schema(&schema).unwrap());
1604        let output_file = file_io.new_output(
1605            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1606        )?;
1607        let mut pw = ParquetWriterBuilder::new(WriterProperties::builder().build(), schema.clone())
1608            .build(output_file)
1609            .await?;
1610        let col0 = Arc::new(
1611            Decimal128Array::from(vec![Some(22000000000), Some(11000000000)])
1612                .with_data_type(DataType::Decimal128(28, 10)),
1613        ) as ArrayRef;
1614        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![col0]).unwrap();
1615        pw.write(&to_write).await?;
1616        let res = pw.close().await?;
1617        assert_eq!(res.len(), 1);
1618        let data_file = res
1619            .into_iter()
1620            .next()
1621            .unwrap()
1622            .content(DataContentType::Data)
1623            .partition(Struct::empty())
1624            .partition_spec_id(0)
1625            .build()
1626            .unwrap();
1627        assert_eq!(
1628            data_file.upper_bounds().get(&0),
1629            Some(Datum::decimal_with_precision(decimal_new(22000000000_i64, 10), 28).unwrap())
1630                .as_ref()
1631        );
1632        assert_eq!(
1633            data_file.lower_bounds().get(&0),
1634            Some(Datum::decimal_with_precision(decimal_new(11000000000_i64, 10), 28).unwrap())
1635                .as_ref()
1636        );
1637
1638        // test -1.1 and -2.2
1639        let schema = Arc::new(
1640            Schema::builder()
1641                .with_fields(vec![
1642                    NestedField::optional(
1643                        0,
1644                        "decimal",
1645                        Type::Primitive(PrimitiveType::Decimal {
1646                            precision: 28,
1647                            scale: 10,
1648                        }),
1649                    )
1650                    .into(),
1651                ])
1652                .build()
1653                .unwrap(),
1654        );
1655        let arrow_schema: ArrowSchemaRef = Arc::new(schema_to_arrow_schema(&schema).unwrap());
1656        let output_file = file_io.new_output(
1657            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1658        )?;
1659        let mut pw = ParquetWriterBuilder::new(WriterProperties::builder().build(), schema.clone())
1660            .build(output_file)
1661            .await?;
1662        let col0 = Arc::new(
1663            Decimal128Array::from(vec![Some(-22000000000), Some(-11000000000)])
1664                .with_data_type(DataType::Decimal128(28, 10)),
1665        ) as ArrayRef;
1666        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![col0]).unwrap();
1667        pw.write(&to_write).await?;
1668        let res = pw.close().await?;
1669        assert_eq!(res.len(), 1);
1670        let data_file = res
1671            .into_iter()
1672            .next()
1673            .unwrap()
1674            .content(DataContentType::Data)
1675            .partition(Struct::empty())
1676            .partition_spec_id(0)
1677            .build()
1678            .unwrap();
1679        assert_eq!(
1680            data_file.upper_bounds().get(&0),
1681            Some(Datum::decimal_with_precision(decimal_new(-11000000000_i64, 10), 28).unwrap())
1682                .as_ref()
1683        );
1684        assert_eq!(
1685            data_file.lower_bounds().get(&0),
1686            Some(Datum::decimal_with_precision(decimal_new(-22000000000_i64, 10), 28).unwrap())
1687                .as_ref()
1688        );
1689
1690        // test 38-digit precision decimal values (Iceberg spec max)
1691        // Note: fastnum D128::MAX/MIN have impractical exponents, so we use meaningful values
1692        use crate::spec::decimal_utils::decimal_from_str_exact;
1693        let decimal_max = decimal_from_str_exact("99999999999999999999999999999999999999").unwrap();
1694        let decimal_min =
1695            decimal_from_str_exact("-99999999999999999999999999999999999999").unwrap();
1696        assert_eq!(decimal_scale(&decimal_max), decimal_scale(&decimal_min));
1697        let schema = Arc::new(
1698            Schema::builder()
1699                .with_fields(vec![
1700                    NestedField::optional(
1701                        0,
1702                        "decimal",
1703                        Type::Primitive(PrimitiveType::Decimal {
1704                            precision: 38,
1705                            scale: decimal_scale(&decimal_max),
1706                        }),
1707                    )
1708                    .into(),
1709                ])
1710                .build()
1711                .unwrap(),
1712        );
1713        let arrow_schema: ArrowSchemaRef = Arc::new(schema_to_arrow_schema(&schema).unwrap());
1714        let output_file = file_io.new_output(
1715            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1716        )?;
1717        let mut pw = ParquetWriterBuilder::new(WriterProperties::builder().build(), schema)
1718            .build(output_file)
1719            .await?;
1720        let col0 = Arc::new(
1721            Decimal128Array::from(vec![
1722                Some(decimal_mantissa(&decimal_max)),
1723                Some(decimal_mantissa(&decimal_min)),
1724            ])
1725            .with_data_type(DataType::Decimal128(38, 0)),
1726        ) as ArrayRef;
1727        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![col0]).unwrap();
1728        pw.write(&to_write).await?;
1729        let res = pw.close().await?;
1730        assert_eq!(res.len(), 1);
1731        let data_file = res
1732            .into_iter()
1733            .next()
1734            .unwrap()
1735            .content(DataContentType::Data)
1736            .partition(Struct::empty())
1737            .partition_spec_id(0)
1738            .build()
1739            .unwrap();
1740        assert_eq!(
1741            data_file.upper_bounds().get(&0),
1742            Some(Datum::decimal(decimal_max).unwrap()).as_ref()
1743        );
1744        assert_eq!(
1745            data_file.lower_bounds().get(&0),
1746            Some(Datum::decimal(decimal_min).unwrap()).as_ref()
1747        );
1748
1749        // test max and min for scale 38
1750        // # TODO
1751        // Readd this case after resolve https://github.com/apache/iceberg-rust/issues/669
1752        // let schema = Arc::new(
1753        //     Schema::builder()
1754        //         .with_fields(vec![NestedField::optional(
1755        //             0,
1756        //             "decimal",
1757        //             Type::Primitive(PrimitiveType::Decimal {
1758        //                 precision: 38,
1759        //                 scale: 0,
1760        //             }),
1761        //         )
1762        //         .into()])
1763        //         .build()
1764        //         .unwrap(),
1765        // );
1766        // let arrow_schema: ArrowSchemaRef = Arc::new(schema_to_arrow_schema(&schema).unwrap());
1767        // let mut pw = ParquetWriterBuilder::new(
1768        //     WriterProperties::builder().build(),
1769        //     schema,
1770        //     file_io.clone(),
1771        //     loccation_gen,
1772        //     file_name_gen,
1773        // )
1774        // .build()
1775        // .await?;
1776        // let col0 = Arc::new(
1777        //     Decimal128Array::from(vec![
1778        //         Some(99999999999999999999999999999999999999_i128),
1779        //         Some(-99999999999999999999999999999999999999_i128),
1780        //     ])
1781        //     .with_data_type(DataType::Decimal128(38, 0)),
1782        // ) as ArrayRef;
1783        // let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![col0]).unwrap();
1784        // pw.write(&to_write).await?;
1785        // let res = pw.close().await?;
1786        // assert_eq!(res.len(), 1);
1787        // let data_file = res
1788        //     .into_iter()
1789        //     .next()
1790        //     .unwrap()
1791        //     .content(crate::spec::DataContentType::Data)
1792        //     .partition(Struct::empty())
1793        //     .build()
1794        //     .unwrap();
1795        // assert_eq!(
1796        //     data_file.upper_bounds().get(&0),
1797        //     Some(Datum::new(
1798        //         PrimitiveType::Decimal {
1799        //             precision: 38,
1800        //             scale: 0
1801        //         },
1802        //         PrimitiveLiteral::Int128(99999999999999999999999999999999999999_i128)
1803        //     ))
1804        //     .as_ref()
1805        // );
1806        // assert_eq!(
1807        //     data_file.lower_bounds().get(&0),
1808        //     Some(Datum::new(
1809        //         PrimitiveType::Decimal {
1810        //             precision: 38,
1811        //             scale: 0
1812        //         },
1813        //         PrimitiveLiteral::Int128(-99999999999999999999999999999999999999_i128)
1814        //     ))
1815        //     .as_ref()
1816        // );
1817
1818        Ok(())
1819    }
1820
1821    #[tokio::test]
1822    async fn test_empty_write() -> Result<()> {
1823        let temp_dir = TempDir::new().unwrap();
1824        let file_io = FileIO::new_with_fs();
1825        let location_gen = DefaultLocationGenerator::with_data_location(
1826            temp_dir.path().to_str().unwrap().to_string(),
1827        );
1828        let file_name_gen =
1829            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1830
1831        // Test that file will create if data to write
1832        let schema = {
1833            let fields =
1834                vec![
1835                    Field::new("col", DataType::Int64, true).with_metadata(HashMap::from([(
1836                        PARQUET_FIELD_ID_META_KEY.to_string(),
1837                        "0".to_string(),
1838                    )])),
1839                ];
1840            Arc::new(arrow_schema::Schema::new(fields))
1841        };
1842        let col = Arc::new(Int64Array::from_iter_values(0..1024)) as ArrayRef;
1843        let to_write = RecordBatch::try_new(schema.clone(), vec![col]).unwrap();
1844        let file_path = location_gen.generate_location(None, &file_name_gen.generate_file_name());
1845        let output_file = file_io.new_output(&file_path)?;
1846        let mut pw = ParquetWriterBuilder::new(
1847            WriterProperties::builder().build(),
1848            Arc::new(to_write.schema().as_ref().try_into().unwrap()),
1849        )
1850        .build(output_file)
1851        .await?;
1852        pw.write(&to_write).await?;
1853        pw.close().await.unwrap();
1854        assert!(file_io.exists(&file_path).await.unwrap());
1855
1856        // Test that file will not create if no data to write
1857        let file_name_gen =
1858            DefaultFileNameGenerator::new("test_empty".to_string(), None, DataFileFormat::Parquet);
1859        let file_path = location_gen.generate_location(None, &file_name_gen.generate_file_name());
1860        let output_file = file_io.new_output(&file_path)?;
1861        let pw = ParquetWriterBuilder::new(
1862            WriterProperties::builder().build(),
1863            Arc::new(to_write.schema().as_ref().try_into().unwrap()),
1864        )
1865        .build(output_file)
1866        .await?;
1867        pw.close().await.unwrap();
1868        assert!(!file_io.exists(&file_path).await.unwrap());
1869
1870        Ok(())
1871    }
1872
1873    #[tokio::test]
1874    async fn test_nan_val_cnts_primitive_type() -> Result<()> {
1875        let temp_dir = TempDir::new().unwrap();
1876        let file_io = FileIO::new_with_fs();
1877        let location_gen = DefaultLocationGenerator::with_data_location(
1878            temp_dir.path().to_str().unwrap().to_string(),
1879        );
1880        let file_name_gen =
1881            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1882
1883        // prepare data
1884        let arrow_schema = {
1885            let fields = vec![
1886                Field::new("col", DataType::Float32, false).with_metadata(HashMap::from([(
1887                    PARQUET_FIELD_ID_META_KEY.to_string(),
1888                    "0".to_string(),
1889                )])),
1890                Field::new("col2", DataType::Float64, false).with_metadata(HashMap::from([(
1891                    PARQUET_FIELD_ID_META_KEY.to_string(),
1892                    "1".to_string(),
1893                )])),
1894            ];
1895            Arc::new(arrow_schema::Schema::new(fields))
1896        };
1897
1898        let float_32_col = Arc::new(Float32Array::from_iter_values_with_nulls(
1899            [1.0_f32, f32::NAN, 2.0, 2.0].into_iter(),
1900            None,
1901        )) as ArrayRef;
1902
1903        let float_64_col = Arc::new(Float64Array::from_iter_values_with_nulls(
1904            [1.0_f64, f64::NAN, 2.0, 2.0].into_iter(),
1905            None,
1906        )) as ArrayRef;
1907
1908        let to_write =
1909            RecordBatch::try_new(arrow_schema.clone(), vec![float_32_col, float_64_col]).unwrap();
1910        let output_file = file_io.new_output(
1911            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
1912        )?;
1913
1914        // write data
1915        let mut pw = ParquetWriterBuilder::new(
1916            WriterProperties::builder().build(),
1917            Arc::new(to_write.schema().as_ref().try_into().unwrap()),
1918        )
1919        .build(output_file)
1920        .await?;
1921
1922        pw.write(&to_write).await?;
1923        let res = pw.close().await?;
1924        assert_eq!(res.len(), 1);
1925        let data_file = res
1926            .into_iter()
1927            .next()
1928            .unwrap()
1929            // Put dummy field for build successfully.
1930            .content(DataContentType::Data)
1931            .partition(Struct::empty())
1932            .partition_spec_id(0)
1933            .build()
1934            .unwrap();
1935
1936        // check data file
1937        assert_eq!(data_file.record_count(), 4);
1938        assert_eq!(*data_file.value_counts(), HashMap::from([(0, 4), (1, 4)]));
1939        assert_eq!(
1940            *data_file.lower_bounds(),
1941            HashMap::from([(0, Datum::float(1.0_f32)), (1, Datum::double(1.0)),])
1942        );
1943        assert_eq!(
1944            *data_file.upper_bounds(),
1945            HashMap::from([(0, Datum::float(2.0_f32)), (1, Datum::double(2.0)),])
1946        );
1947        assert_eq!(
1948            *data_file.null_value_counts(),
1949            HashMap::from([(0, 0), (1, 0)])
1950        );
1951        assert_eq!(
1952            *data_file.nan_value_counts(),
1953            HashMap::from([(0, 1), (1, 1)])
1954        );
1955
1956        // check the written file
1957        let expect_batch = concat_batches(&arrow_schema, vec![&to_write]).unwrap();
1958        check_parquet_data_file(&file_io, &data_file, &expect_batch).await;
1959
1960        Ok(())
1961    }
1962
1963    #[tokio::test]
1964    async fn test_nan_val_cnts_struct_type() -> Result<()> {
1965        let temp_dir = TempDir::new().unwrap();
1966        let file_io = FileIO::new_with_fs();
1967        let location_gen = DefaultLocationGenerator::with_data_location(
1968            temp_dir.path().to_str().unwrap().to_string(),
1969        );
1970        let file_name_gen =
1971            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
1972
1973        let schema_struct_float_fields = Fields::from(vec![
1974            Field::new("col4", DataType::Float32, false).with_metadata(HashMap::from([(
1975                PARQUET_FIELD_ID_META_KEY.to_string(),
1976                "4".to_string(),
1977            )])),
1978        ]);
1979
1980        let schema_struct_nested_float_fields = Fields::from(vec![
1981            Field::new("col7", DataType::Float32, false).with_metadata(HashMap::from([(
1982                PARQUET_FIELD_ID_META_KEY.to_string(),
1983                "7".to_string(),
1984            )])),
1985        ]);
1986
1987        let schema_struct_nested_fields = Fields::from(vec![
1988            Field::new(
1989                "col6",
1990                DataType::Struct(schema_struct_nested_float_fields.clone()),
1991                false,
1992            )
1993            .with_metadata(HashMap::from([(
1994                PARQUET_FIELD_ID_META_KEY.to_string(),
1995                "6".to_string(),
1996            )])),
1997        ]);
1998
1999        // prepare data
2000        let arrow_schema = {
2001            let fields = vec![
2002                Field::new(
2003                    "col3",
2004                    DataType::Struct(schema_struct_float_fields.clone()),
2005                    false,
2006                )
2007                .with_metadata(HashMap::from([(
2008                    PARQUET_FIELD_ID_META_KEY.to_string(),
2009                    "3".to_string(),
2010                )])),
2011                Field::new(
2012                    "col5",
2013                    DataType::Struct(schema_struct_nested_fields.clone()),
2014                    false,
2015                )
2016                .with_metadata(HashMap::from([(
2017                    PARQUET_FIELD_ID_META_KEY.to_string(),
2018                    "5".to_string(),
2019                )])),
2020            ];
2021            Arc::new(arrow_schema::Schema::new(fields))
2022        };
2023
2024        let float_32_col = Arc::new(Float32Array::from_iter_values_with_nulls(
2025            [1.0_f32, f32::NAN, 2.0, 2.0].into_iter(),
2026            None,
2027        )) as ArrayRef;
2028
2029        let struct_float_field_col = Arc::new(StructArray::new(
2030            schema_struct_float_fields,
2031            vec![float_32_col.clone()],
2032            None,
2033        )) as ArrayRef;
2034
2035        let struct_nested_float_field_col = Arc::new(StructArray::new(
2036            schema_struct_nested_fields,
2037            vec![Arc::new(StructArray::new(
2038                schema_struct_nested_float_fields,
2039                vec![float_32_col.clone()],
2040                None,
2041            )) as ArrayRef],
2042            None,
2043        )) as ArrayRef;
2044
2045        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![
2046            struct_float_field_col,
2047            struct_nested_float_field_col,
2048        ])
2049        .unwrap();
2050        let output_file = file_io.new_output(
2051            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
2052        )?;
2053
2054        // write data
2055        let mut pw = ParquetWriterBuilder::new(
2056            WriterProperties::builder().build(),
2057            Arc::new(to_write.schema().as_ref().try_into().unwrap()),
2058        )
2059        .build(output_file)
2060        .await?;
2061
2062        pw.write(&to_write).await?;
2063        let res = pw.close().await?;
2064        assert_eq!(res.len(), 1);
2065        let data_file = res
2066            .into_iter()
2067            .next()
2068            .unwrap()
2069            // Put dummy field for build successfully.
2070            .content(DataContentType::Data)
2071            .partition(Struct::empty())
2072            .partition_spec_id(0)
2073            .build()
2074            .unwrap();
2075
2076        // check data file
2077        assert_eq!(data_file.record_count(), 4);
2078        assert_eq!(*data_file.value_counts(), HashMap::from([(4, 4), (7, 4)]));
2079        assert_eq!(
2080            *data_file.lower_bounds(),
2081            HashMap::from([(4, Datum::float(1.0_f32)), (7, Datum::float(1.0_f32)),])
2082        );
2083        assert_eq!(
2084            *data_file.upper_bounds(),
2085            HashMap::from([(4, Datum::float(2.0_f32)), (7, Datum::float(2.0_f32)),])
2086        );
2087        assert_eq!(
2088            *data_file.null_value_counts(),
2089            HashMap::from([(4, 0), (7, 0)])
2090        );
2091        assert_eq!(
2092            *data_file.nan_value_counts(),
2093            HashMap::from([(4, 1), (7, 1)])
2094        );
2095
2096        // check the written file
2097        let expect_batch = concat_batches(&arrow_schema, vec![&to_write]).unwrap();
2098        check_parquet_data_file(&file_io, &data_file, &expect_batch).await;
2099
2100        Ok(())
2101    }
2102
2103    #[tokio::test]
2104    async fn test_nan_val_cnts_list_type() -> Result<()> {
2105        let temp_dir = TempDir::new().unwrap();
2106        let file_io = FileIO::new_with_fs();
2107        let location_gen = DefaultLocationGenerator::with_data_location(
2108            temp_dir.path().to_str().unwrap().to_string(),
2109        );
2110        let file_name_gen =
2111            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
2112
2113        let schema_list_float_field = Field::new("element", DataType::Float32, true).with_metadata(
2114            HashMap::from([(PARQUET_FIELD_ID_META_KEY.to_string(), "1".to_string())]),
2115        );
2116
2117        let schema_struct_list_float_field = Field::new("element", DataType::Float32, true)
2118            .with_metadata(HashMap::from([(
2119                PARQUET_FIELD_ID_META_KEY.to_string(),
2120                "4".to_string(),
2121            )]));
2122
2123        let schema_struct_list_field = Fields::from(vec![
2124            Field::new_list("col2", schema_struct_list_float_field.clone(), true).with_metadata(
2125                HashMap::from([(PARQUET_FIELD_ID_META_KEY.to_string(), "3".to_string())]),
2126            ),
2127        ]);
2128
2129        let arrow_schema = {
2130            let fields = vec![
2131                Field::new_list("col0", schema_list_float_field.clone(), true).with_metadata(
2132                    HashMap::from([(PARQUET_FIELD_ID_META_KEY.to_string(), "0".to_string())]),
2133                ),
2134                Field::new_struct("col1", schema_struct_list_field.clone(), true)
2135                    .with_metadata(HashMap::from([(
2136                        PARQUET_FIELD_ID_META_KEY.to_string(),
2137                        "2".to_string(),
2138                    )]))
2139                    .clone(),
2140                // Field::new_large_list("col3", schema_large_list_float_field.clone(), true).with_metadata(
2141                //     HashMap::from([(PARQUET_FIELD_ID_META_KEY.to_string(), "5".to_string())]),
2142                // ).clone(),
2143            ];
2144            Arc::new(arrow_schema::Schema::new(fields))
2145        };
2146
2147        let list_parts = ListArray::from_iter_primitive::<Float32Type, _, _>(vec![Some(vec![
2148            Some(1.0_f32),
2149            Some(f32::NAN),
2150            Some(2.0),
2151            Some(2.0),
2152        ])])
2153        .into_parts();
2154
2155        let list_float_field_col = Arc::new({
2156            let list_parts = list_parts.clone();
2157            ListArray::new(
2158                {
2159                    if let DataType::List(field) = arrow_schema.field(0).data_type() {
2160                        field.clone()
2161                    } else {
2162                        unreachable!()
2163                    }
2164                },
2165                list_parts.1,
2166                list_parts.2,
2167                list_parts.3,
2168            )
2169        }) as ArrayRef;
2170
2171        let struct_list_fields_schema =
2172            if let DataType::Struct(fields) = arrow_schema.field(1).data_type() {
2173                fields.clone()
2174            } else {
2175                unreachable!()
2176            };
2177
2178        let struct_list_float_field_col = Arc::new({
2179            ListArray::new(
2180                {
2181                    if let DataType::List(field) = struct_list_fields_schema
2182                        .first()
2183                        .expect("could not find first list field")
2184                        .data_type()
2185                    {
2186                        field.clone()
2187                    } else {
2188                        unreachable!()
2189                    }
2190                },
2191                list_parts.1,
2192                list_parts.2,
2193                list_parts.3,
2194            )
2195        }) as ArrayRef;
2196
2197        let struct_list_float_field_col = Arc::new(StructArray::new(
2198            struct_list_fields_schema,
2199            vec![struct_list_float_field_col.clone()],
2200            None,
2201        )) as ArrayRef;
2202
2203        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![
2204            list_float_field_col,
2205            struct_list_float_field_col,
2206            // large_list_float_field_col,
2207        ])
2208        .expect("Could not form record batch");
2209        let output_file = file_io.new_output(
2210            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
2211        )?;
2212
2213        // write data
2214        let mut pw = ParquetWriterBuilder::new(
2215            WriterProperties::builder().build(),
2216            Arc::new(
2217                to_write
2218                    .schema()
2219                    .as_ref()
2220                    .try_into()
2221                    .expect("Could not convert iceberg schema"),
2222            ),
2223        )
2224        .build(output_file)
2225        .await?;
2226
2227        pw.write(&to_write).await?;
2228        let res = pw.close().await?;
2229        assert_eq!(res.len(), 1);
2230        let data_file = res
2231            .into_iter()
2232            .next()
2233            .unwrap()
2234            .content(DataContentType::Data)
2235            .partition(Struct::empty())
2236            .partition_spec_id(0)
2237            .build()
2238            .unwrap();
2239
2240        // check data file
2241        assert_eq!(data_file.record_count(), 1);
2242        assert_eq!(*data_file.value_counts(), HashMap::from([(1, 4), (4, 4)]));
2243        assert_eq!(
2244            *data_file.lower_bounds(),
2245            HashMap::from([(1, Datum::float(1.0_f32)), (4, Datum::float(1.0_f32))])
2246        );
2247        assert_eq!(
2248            *data_file.upper_bounds(),
2249            HashMap::from([(1, Datum::float(2.0_f32)), (4, Datum::float(2.0_f32))])
2250        );
2251        assert_eq!(
2252            *data_file.null_value_counts(),
2253            HashMap::from([(1, 0), (4, 0)])
2254        );
2255        assert_eq!(
2256            *data_file.nan_value_counts(),
2257            HashMap::from([(1, 1), (4, 1)])
2258        );
2259
2260        // check the written file
2261        let expect_batch = concat_batches(&arrow_schema, vec![&to_write]).unwrap();
2262        check_parquet_data_file(&file_io, &data_file, &expect_batch).await;
2263
2264        Ok(())
2265    }
2266
2267    macro_rules! construct_map_arr {
2268        ($map_key_field_schema:ident, $map_value_field_schema:ident) => {{
2269            let int_builder = Int32Builder::new();
2270            let float_builder = Float32Builder::with_capacity(4);
2271            let mut builder = MapBuilder::new(None, int_builder, float_builder);
2272            builder.keys().append_value(1);
2273            builder.values().append_value(1.0_f32);
2274            builder.append(true).unwrap();
2275            builder.keys().append_value(2);
2276            builder.values().append_value(f32::NAN);
2277            builder.append(true).unwrap();
2278            builder.keys().append_value(3);
2279            builder.values().append_value(2.0);
2280            builder.append(true).unwrap();
2281            builder.keys().append_value(4);
2282            builder.values().append_value(2.0);
2283            builder.append(true).unwrap();
2284            let array = builder.finish();
2285
2286            let (_field, offsets, entries, nulls, ordered) = array.into_parts();
2287            let new_struct_fields_schema =
2288                Fields::from(vec![$map_key_field_schema, $map_value_field_schema]);
2289
2290            let entries = {
2291                let (_, arrays, nulls) = entries.into_parts();
2292                StructArray::new(new_struct_fields_schema.clone(), arrays, nulls)
2293            };
2294
2295            let field = Arc::new(Field::new(
2296                DEFAULT_MAP_FIELD_NAME,
2297                DataType::Struct(new_struct_fields_schema),
2298                false,
2299            ));
2300
2301            Arc::new(MapArray::new(field, offsets, entries, nulls, ordered))
2302        }};
2303    }
2304
2305    #[tokio::test]
2306    async fn test_nan_val_cnts_map_type() -> Result<()> {
2307        let temp_dir = TempDir::new().unwrap();
2308        let file_io = FileIO::new_with_fs();
2309        let location_gen = DefaultLocationGenerator::with_data_location(
2310            temp_dir.path().to_str().unwrap().to_string(),
2311        );
2312        let file_name_gen =
2313            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
2314
2315        let map_key_field_schema =
2316            Field::new(MAP_KEY_FIELD_NAME, DataType::Int32, false).with_metadata(HashMap::from([
2317                (PARQUET_FIELD_ID_META_KEY.to_string(), "1".to_string()),
2318            ]));
2319
2320        let map_value_field_schema =
2321            Field::new(MAP_VALUE_FIELD_NAME, DataType::Float32, true).with_metadata(HashMap::from(
2322                [(PARQUET_FIELD_ID_META_KEY.to_string(), "2".to_string())],
2323            ));
2324
2325        let struct_map_key_field_schema =
2326            Field::new(MAP_KEY_FIELD_NAME, DataType::Int32, false).with_metadata(HashMap::from([
2327                (PARQUET_FIELD_ID_META_KEY.to_string(), "6".to_string()),
2328            ]));
2329
2330        let struct_map_value_field_schema =
2331            Field::new(MAP_VALUE_FIELD_NAME, DataType::Float32, true).with_metadata(HashMap::from(
2332                [(PARQUET_FIELD_ID_META_KEY.to_string(), "7".to_string())],
2333            ));
2334
2335        let schema_struct_map_field = Fields::from(vec![
2336            Field::new_map(
2337                "col3",
2338                DEFAULT_MAP_FIELD_NAME,
2339                struct_map_key_field_schema.clone(),
2340                struct_map_value_field_schema.clone(),
2341                false,
2342                false,
2343            )
2344            .with_metadata(HashMap::from([(
2345                PARQUET_FIELD_ID_META_KEY.to_string(),
2346                "5".to_string(),
2347            )])),
2348        ]);
2349
2350        let arrow_schema = {
2351            let fields = vec![
2352                Field::new_map(
2353                    "col0",
2354                    DEFAULT_MAP_FIELD_NAME,
2355                    map_key_field_schema.clone(),
2356                    map_value_field_schema.clone(),
2357                    false,
2358                    false,
2359                )
2360                .with_metadata(HashMap::from([(
2361                    PARQUET_FIELD_ID_META_KEY.to_string(),
2362                    "0".to_string(),
2363                )])),
2364                Field::new_struct("col1", schema_struct_map_field.clone(), true)
2365                    .with_metadata(HashMap::from([(
2366                        PARQUET_FIELD_ID_META_KEY.to_string(),
2367                        "3".to_string(),
2368                    )]))
2369                    .clone(),
2370            ];
2371            Arc::new(arrow_schema::Schema::new(fields))
2372        };
2373
2374        let map_array = construct_map_arr!(map_key_field_schema, map_value_field_schema);
2375
2376        let struct_map_arr =
2377            construct_map_arr!(struct_map_key_field_schema, struct_map_value_field_schema);
2378
2379        let struct_list_float_field_col = Arc::new(StructArray::new(
2380            schema_struct_map_field,
2381            vec![struct_map_arr],
2382            None,
2383        )) as ArrayRef;
2384
2385        let to_write = RecordBatch::try_new(arrow_schema.clone(), vec![
2386            map_array,
2387            struct_list_float_field_col,
2388        ])
2389        .expect("Could not form record batch");
2390        let output_file = file_io.new_output(
2391            location_gen.generate_location(None, &file_name_gen.generate_file_name()),
2392        )?;
2393
2394        // write data
2395        let mut pw = ParquetWriterBuilder::new(
2396            WriterProperties::builder().build(),
2397            Arc::new(
2398                to_write
2399                    .schema()
2400                    .as_ref()
2401                    .try_into()
2402                    .expect("Could not convert iceberg schema"),
2403            ),
2404        )
2405        .build(output_file)
2406        .await?;
2407
2408        pw.write(&to_write).await?;
2409        let res = pw.close().await?;
2410        assert_eq!(res.len(), 1);
2411        let data_file = res
2412            .into_iter()
2413            .next()
2414            .unwrap()
2415            .content(DataContentType::Data)
2416            .partition(Struct::empty())
2417            .partition_spec_id(0)
2418            .build()
2419            .unwrap();
2420
2421        // check data file
2422        assert_eq!(data_file.record_count(), 4);
2423        assert_eq!(
2424            *data_file.value_counts(),
2425            HashMap::from([(1, 4), (2, 4), (6, 4), (7, 4)])
2426        );
2427        assert_eq!(
2428            *data_file.lower_bounds(),
2429            HashMap::from([
2430                (1, Datum::int(1)),
2431                (2, Datum::float(1.0_f32)),
2432                (6, Datum::int(1)),
2433                (7, Datum::float(1.0_f32))
2434            ])
2435        );
2436        assert_eq!(
2437            *data_file.upper_bounds(),
2438            HashMap::from([
2439                (1, Datum::int(4)),
2440                (2, Datum::float(2.0_f32)),
2441                (6, Datum::int(4)),
2442                (7, Datum::float(2.0_f32))
2443            ])
2444        );
2445        assert_eq!(
2446            *data_file.null_value_counts(),
2447            HashMap::from([(1, 0), (2, 0), (6, 0), (7, 0)])
2448        );
2449        assert_eq!(
2450            *data_file.nan_value_counts(),
2451            HashMap::from([(2, 1), (7, 1)])
2452        );
2453
2454        // check the written file
2455        let expect_batch = concat_batches(&arrow_schema, vec![&to_write]).unwrap();
2456        check_parquet_data_file(&file_io, &data_file, &expect_batch).await;
2457
2458        Ok(())
2459    }
2460
2461    #[tokio::test]
2462    async fn test_write_empty_parquet_file() {
2463        let temp_dir = TempDir::new().unwrap();
2464        let file_io = FileIO::new_with_fs();
2465        let location_gen = DefaultLocationGenerator::with_data_location(
2466            temp_dir.path().to_str().unwrap().to_string(),
2467        );
2468        let file_name_gen =
2469            DefaultFileNameGenerator::new("test".to_string(), None, DataFileFormat::Parquet);
2470        let output_file = file_io
2471            .new_output(location_gen.generate_location(None, &file_name_gen.generate_file_name()))
2472            .unwrap();
2473
2474        // write data
2475        let pw = ParquetWriterBuilder::new(
2476            WriterProperties::builder().build(),
2477            Arc::new(
2478                Schema::builder()
2479                    .with_schema_id(1)
2480                    .with_fields(vec![
2481                        NestedField::required(0, "col", Type::Primitive(PrimitiveType::Long))
2482                            .with_id(0)
2483                            .into(),
2484                    ])
2485                    .build()
2486                    .expect("Failed to create schema"),
2487            ),
2488        )
2489        .build(output_file)
2490        .await
2491        .unwrap();
2492
2493        let res = pw.close().await.unwrap();
2494        assert_eq!(res.len(), 0);
2495
2496        // Check that file should have been deleted.
2497        assert_eq!(std::fs::read_dir(temp_dir.path()).unwrap().count(), 0);
2498    }
2499
2500    #[test]
2501    fn test_min_max_aggregator() {
2502        let schema = Arc::new(
2503            Schema::builder()
2504                .with_schema_id(1)
2505                .with_fields(vec![
2506                    NestedField::required(0, "col", Type::Primitive(PrimitiveType::Int))
2507                        .with_id(0)
2508                        .into(),
2509                ])
2510                .build()
2511                .expect("Failed to create schema"),
2512        );
2513        let mut min_max_agg = MinMaxColAggregator::new(schema);
2514        let create_statistics =
2515            |min, max| Statistics::Int32(ValueStatistics::new(min, max, None, None, false));
2516        min_max_agg
2517            .update(0, &create_statistics(None, Some(42)))
2518            .unwrap();
2519        min_max_agg
2520            .update(0, &create_statistics(Some(0), Some(i32::MAX)))
2521            .unwrap();
2522        min_max_agg
2523            .update(0, &create_statistics(Some(i32::MIN), None))
2524            .unwrap();
2525        min_max_agg
2526            .update(0, &create_statistics(None, None))
2527            .unwrap();
2528
2529        let (lower_bounds, upper_bounds) = min_max_agg.produce();
2530
2531        assert_eq!(lower_bounds, HashMap::from([(0, Datum::int(i32::MIN))]));
2532        assert_eq!(upper_bounds, HashMap::from([(0, Datum::int(i32::MAX))]));
2533    }
2534
2535    fn cdc_test_schema() -> SchemaRef {
2536        Arc::new(
2537            Schema::builder()
2538                .with_schema_id(1)
2539                .with_fields(vec![
2540                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
2541                    NestedField::required(2, "payload", Type::Primitive(PrimitiveType::String))
2542                        .into(),
2543                ])
2544                .build()
2545                .unwrap(),
2546        )
2547    }
2548
2549    #[test]
2550    fn test_from_table_properties_no_cdc_by_default() {
2551        let raw_properties = HashMap::new();
2552        let tp = TableProperties::new(&raw_properties);
2553        let builder = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema()).unwrap();
2554        assert!(builder.props.content_defined_chunking().is_none());
2555    }
2556
2557    #[tokio::test]
2558    async fn test_from_table_properties_without_encryption_writes_plaintext() {
2559        let raw_properties = HashMap::new();
2560        let tp = TableProperties::new(&raw_properties);
2561        let tmp = TempDir::new().unwrap();
2562        let output = FileIO::new_with_fs()
2563            .new_output(format!("{}/plain.parquet", tmp.path().to_str().unwrap()))
2564            .unwrap();
2565        let writer = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema())
2566            .unwrap()
2567            .build(output)
2568            .await
2569            .unwrap();
2570
2571        assert!(writer.key_metadata.is_none());
2572        assert!(
2573            writer
2574                .writer_properties
2575                .file_encryption_properties()
2576                .is_none()
2577        );
2578    }
2579
2580    #[tokio::test]
2581    async fn test_from_table_properties_propagate_to_writer() {
2582        // `build()` must carry the translated `WriterProperties` through to the
2583        // `ParquetWriter` unchanged — otherwise the `write.parquet.*` settings
2584        // derived in `from_table_properties` would never reach parquet-rs.
2585        //
2586        // Asserting on the writer's `WriterProperties` (rather than re-reading a
2587        // written file) keeps this a direct propagation check: every future
2588        // `write.parquet.*` option just adds an assertion on its corresponding
2589        // `WriterProperties` getter here.
2590        let raw_properties = HashMap::from([
2591            (
2592                TableProperties::PROPERTY_PARQUET_CDC_ENABLED.to_string(),
2593                "true".to_string(),
2594            ),
2595            (
2596                TableProperties::PROPERTY_PARQUET_CDC_MIN_CHUNK_SIZE.to_string(),
2597                "4096".to_string(),
2598            ),
2599            (
2600                TableProperties::PROPERTY_PARQUET_CDC_MAX_CHUNK_SIZE.to_string(),
2601                "8192".to_string(),
2602            ),
2603            (
2604                TableProperties::PROPERTY_PARQUET_CDC_NORM_LEVEL.to_string(),
2605                "2".to_string(),
2606            ),
2607        ]);
2608        let tp = TableProperties::new(&raw_properties);
2609
2610        let tmp = TempDir::new().unwrap();
2611        let output = FileIO::new_with_fs()
2612            .new_output(format!("{}/cdc.parquet", tmp.path().to_str().unwrap()))
2613            .unwrap();
2614        let writer = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema())
2615            .unwrap()
2616            .build(output)
2617            .await
2618            .unwrap();
2619
2620        let cdc = writer
2621            .writer_properties
2622            .content_defined_chunking()
2623            .copied()
2624            .expect("CDC should be enabled on the built writer");
2625        assert_eq!(cdc.min_chunk_size, 4096);
2626        assert_eq!(cdc.max_chunk_size, 8192);
2627        assert_eq!(cdc.norm_level, 2);
2628    }
2629
2630    #[test]
2631    fn test_from_table_properties_sizing_defaults() {
2632        // With no properties set, the writer must use Iceberg's defaults (which
2633        // differ from parquet-rs's own defaults), not parquet-rs's.
2634        let raw_properties = HashMap::new();
2635        let tp = TableProperties::new(&raw_properties);
2636        let props = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema())
2637            .unwrap()
2638            .props;
2639
2640        assert_eq!(
2641            props.max_row_group_bytes(),
2642            Some(TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT)
2643        );
2644        assert_eq!(
2645            props.data_page_size_limit(),
2646            TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES_DEFAULT
2647        );
2648        assert_eq!(
2649            props.data_page_row_count_limit(),
2650            TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT_DEFAULT
2651        );
2652        assert_eq!(
2653            props.dictionary_page_size_limit(),
2654            TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES_DEFAULT
2655        );
2656        // Default codec is zstd at the Java-aligned default level (3).
2657        assert_eq!(
2658            props.compression(&ColumnPath::from("id")),
2659            Compression::ZSTD(ZstdLevel::try_new(3).unwrap())
2660        );
2661    }
2662
2663    #[test]
2664    fn test_from_table_properties_sizing_and_compression_overrides() {
2665        let raw_properties = HashMap::from([
2666            (
2667                TableProperties::PROPERTY_PARQUET_ROW_GROUP_SIZE_BYTES.to_string(),
2668                "1048576".to_string(),
2669            ),
2670            (
2671                TableProperties::PROPERTY_PARQUET_PAGE_SIZE_BYTES.to_string(),
2672                "65536".to_string(),
2673            ),
2674            (
2675                TableProperties::PROPERTY_PARQUET_PAGE_ROW_LIMIT.to_string(),
2676                "5000".to_string(),
2677            ),
2678            (
2679                TableProperties::PROPERTY_PARQUET_DICT_SIZE_BYTES.to_string(),
2680                "131072".to_string(),
2681            ),
2682            (
2683                TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(),
2684                "gzip".to_string(),
2685            ),
2686            (
2687                TableProperties::PROPERTY_PARQUET_COMPRESSION_LEVEL.to_string(),
2688                "9".to_string(),
2689            ),
2690        ]);
2691        let tp = TableProperties::new(&raw_properties);
2692        let props = ParquetWriterBuilder::from_table_properties(&tp, cdc_test_schema())
2693            .unwrap()
2694            .props;
2695
2696        assert_eq!(props.max_row_group_bytes(), Some(1048576));
2697        assert_eq!(props.data_page_size_limit(), 65536);
2698        assert_eq!(props.data_page_row_count_limit(), 5000);
2699        assert_eq!(props.dictionary_page_size_limit(), 131072);
2700        assert_eq!(
2701            props.compression(&ColumnPath::from("id")),
2702            Compression::GZIP(GzipLevel::try_new(9).unwrap())
2703        );
2704    }
2705
2706    #[test]
2707    fn test_from_table_properties_invalid_codec_errors() {
2708        let entries = HashMap::from([(
2709            TableProperties::PROPERTY_PARQUET_COMPRESSION_CODEC.to_string(),
2710            "bogus".to_string(),
2711        )]);
2712        let err = ParquetWriterBuilder::from_table_properties(
2713            &TableProperties::new(&entries),
2714            cdc_test_schema(),
2715        )
2716        .unwrap_err();
2717        assert_eq!(err.kind(), ErrorKind::DataInvalid);
2718        assert!(err.to_string().contains("bogus"));
2719    }
2720
2721    #[test]
2722    fn test_parquet_compression_mapping() {
2723        // Codecs without a level.
2724        assert_eq!(
2725            parquet_compression(CompressionCodec::None).unwrap(),
2726            Compression::UNCOMPRESSED
2727        );
2728        assert_eq!(
2729            parquet_compression(CompressionCodec::Snappy).unwrap(),
2730            Compression::SNAPPY
2731        );
2732        assert_eq!(
2733            parquet_compression(CompressionCodec::Lz4).unwrap(),
2734            Compression::LZ4
2735        );
2736        assert_eq!(
2737            parquet_compression(CompressionCodec::Lz4Raw).unwrap(),
2738            Compression::LZ4_RAW
2739        );
2740        assert_eq!(
2741            parquet_compression(CompressionCodec::Lzo).unwrap(),
2742            Compression::LZO
2743        );
2744
2745        // Level-carrying codecs at their defaults.
2746        assert_eq!(
2747            parquet_compression(CompressionCodec::zstd_default()).unwrap(),
2748            Compression::ZSTD(ZstdLevel::try_new(3).unwrap())
2749        );
2750        assert_eq!(
2751            parquet_compression(CompressionCodec::gzip_default()).unwrap(),
2752            Compression::GZIP(GzipLevel::try_new(6).unwrap())
2753        );
2754        assert_eq!(
2755            parquet_compression(CompressionCodec::brotli_default()).unwrap(),
2756            Compression::BROTLI(BrotliLevel::try_new(1).unwrap())
2757        );
2758
2759        // Explicit levels are honored.
2760        assert_eq!(
2761            parquet_compression(CompressionCodec::Zstd(10)).unwrap(),
2762            Compression::ZSTD(ZstdLevel::try_new(10).unwrap())
2763        );
2764    }
2765
2766    #[test]
2767    fn test_parquet_compression_invalid_level() {
2768        // zstd valid range is 1..=22; 99 is out of range.
2769        let err = parquet_compression(CompressionCodec::Zstd(99)).unwrap_err();
2770        assert_eq!(err.kind(), ErrorKind::DataInvalid);
2771        assert!(err.to_string().contains("zstd"));
2772    }
2773}