1use 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#[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 pub fn new(props: WriterProperties, schema: SchemaRef) -> Self {
71 Self::new_with_match_mode(props, schema, FieldMatchMode::Id)
72 }
73
74 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 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 pub fn with_match_mode(mut self, match_mode: FieldMatchMode) -> Self {
121 self.match_mode = match_mode;
122 self
123 }
124
125 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
184struct IndexByParquetPathName {
186 name_to_id: HashMap<String, i32>,
187
188 field_names: Vec<String>,
189
190 field_id: i32,
191}
192
193impl IndexByParquetPathName {
194 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 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
306pub 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
317struct MinMaxColAggregator {
319 lower_bounds: HashMap<i32, Datum>,
320 upper_bounds: HashMap<i32, Datum>,
321 schema: SchemaRef,
322}
323
324impl MinMaxColAggregator {
325 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 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 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 fn produce(self) -> (HashMap<i32, Datum>, HashMap<i32, Datum>) {
401 (self.lower_bounds, self.upper_bounds)
402 }
403}
404
405impl ParquetWriter {
406 #[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 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 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 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 .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 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 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.bytes_written() + inner.in_progress_size()
702 } else {
703 0
705 }
706 }
707}
708
709struct AsyncFileWriter(Box<dyn FileWrite>);
715
716impl AsyncFileWriter {
717 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 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 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 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 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 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 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 .content(DataContentType::Data)
997 .partition(Struct::empty())
998 .partition_spec_id(0)
999 .build()
1000 .unwrap();
1001
1002 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 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 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 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 .content(DataContentType::Data)
1282 .partition(Struct::empty())
1283 .partition_spec_id(0)
1284 .build()
1285 .unwrap();
1286
1287 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_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 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 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 .content(DataContentType::Data)
1472 .partition(Struct::empty())
1473 .partition_spec_id(0)
1474 .build()
1475 .unwrap();
1476
1477 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_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 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 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 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 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 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 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 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 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 .content(DataContentType::Data)
1931 .partition(Struct::empty())
1932 .partition_spec_id(0)
1933 .build()
1934 .unwrap();
1935
1936 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 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 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 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 .content(DataContentType::Data)
2071 .partition(Struct::empty())
2072 .partition_spec_id(0)
2073 .build()
2074 .unwrap();
2075
2076 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 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 ];
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 ])
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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}