1use crate::bloom_filter::Sbbf;
21use crate::file::metadata::thrift::PageHeader;
22use crate::file::page_index::column_index::ColumnIndexMetaData;
23use crate::file::page_index::offset_index::OffsetIndexMetaData;
24use crate::parquet_thrift::{ThriftCompactOutputProtocol, WriteThrift};
25#[cfg(feature = "arrow")]
26use bytes::Bytes;
27use std::fmt::Debug;
28use std::io::{BufWriter, IoSlice, Read};
29use std::{io::Write, sync::Arc};
30
31use crate::column::page_encryption::PageEncryptor;
32use crate::column::writer::{ColumnCloseResult, ColumnWriterImpl, get_typed_column_writer_mut};
33use crate::column::{
34 page::{CompressedPage, PageWriteSpec, PageWriter},
35 writer::{ColumnWriter, get_column_writer},
36};
37use crate::data_type::DataType;
38#[cfg(feature = "encryption")]
39use crate::encryption::encrypt::{
40 FileEncryptionProperties, FileEncryptor, get_column_crypto_metadata,
41};
42use crate::errors::{ParquetError, Result};
43#[cfg(feature = "encryption")]
44use crate::file::PARQUET_MAGIC_ENCR_FOOTER;
45use crate::file::properties::{BloomFilterPosition, WriterPropertiesPtr};
46use crate::file::reader::ChunkReader;
47use crate::file::{PARQUET_MAGIC, metadata::*};
48use crate::schema::types::{ColumnDescPtr, SchemaDescPtr, SchemaDescriptor, TypePtr};
49
50pub struct TrackedWrite<W: Write> {
54 inner: BufWriter<W>,
55 bytes_written: usize,
56}
57
58impl<W: Write> TrackedWrite<W> {
59 pub fn new(inner: W) -> Self {
61 let buf_write = BufWriter::new(inner);
62 Self {
63 inner: buf_write,
64 bytes_written: 0,
65 }
66 }
67
68 pub fn bytes_written(&self) -> usize {
70 self.bytes_written
71 }
72
73 pub fn inner(&self) -> &W {
75 self.inner.get_ref()
76 }
77
78 pub fn inner_mut(&mut self) -> &mut W {
83 self.inner.get_mut()
84 }
85
86 pub fn into_inner(self) -> Result<W> {
88 self.inner.into_inner().map_err(|err| {
89 ParquetError::General(format!("fail to get inner writer: {:?}", err.to_string()))
90 })
91 }
92}
93
94impl<W: Write> Write for TrackedWrite<W> {
95 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
96 let bytes = self.inner.write(buf)?;
97 self.bytes_written += bytes;
98 Ok(bytes)
99 }
100
101 fn write_vectored(&mut self, bufs: &[IoSlice<'_>]) -> std::io::Result<usize> {
102 let bytes = self.inner.write_vectored(bufs)?;
103 self.bytes_written += bytes;
104 Ok(bytes)
105 }
106
107 fn write_all(&mut self, buf: &[u8]) -> std::io::Result<()> {
108 self.inner.write_all(buf)?;
109 self.bytes_written += buf.len();
110
111 Ok(())
112 }
113
114 fn flush(&mut self) -> std::io::Result<()> {
115 self.inner.flush()
116 }
117}
118
119pub type OnCloseColumnChunk<'a> = Box<dyn FnOnce(ColumnCloseResult) -> Result<()> + 'a>;
121
122pub type OnCloseRowGroup<'a, W> = Box<
128 dyn FnOnce(
129 &'a mut TrackedWrite<W>,
130 RowGroupMetaData,
131 Vec<Option<Sbbf>>,
132 Vec<Option<ColumnIndexMetaData>>,
133 Vec<Option<OffsetIndexMetaData>>,
134 ) -> Result<()>
135 + 'a
136 + Send,
137>;
138
139pub struct SerializedFileWriter<W: Write> {
159 buf: TrackedWrite<W>,
160 descr: SchemaDescPtr,
161 props: WriterPropertiesPtr,
162 row_groups: Vec<RowGroupMetaData>,
163 bloom_filters: Vec<Vec<Option<Sbbf>>>,
164 column_indexes: Vec<Vec<Option<ColumnIndexMetaData>>>,
165 offset_indexes: Vec<Vec<Option<OffsetIndexMetaData>>>,
166 row_group_index: usize,
167 kv_metadatas: Vec<KeyValue>,
169 finished: bool,
170 #[cfg(feature = "encryption")]
171 file_encryptor: Option<Arc<FileEncryptor>>,
172}
173
174impl<W: Write> Debug for SerializedFileWriter<W> {
175 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
176 f.debug_struct("SerializedFileWriter")
179 .field("descr", &self.descr)
180 .field("row_group_index", &self.row_group_index)
181 .field("kv_metadatas", &self.kv_metadatas)
182 .finish_non_exhaustive()
183 }
184}
185
186impl<W: Write + Send> SerializedFileWriter<W> {
187 pub fn new(buf: W, schema: TypePtr, properties: WriterPropertiesPtr) -> Result<Self> {
189 let mut buf = TrackedWrite::new(buf);
190
191 let schema_descriptor = SchemaDescriptor::new(schema.clone());
192
193 #[cfg(feature = "encryption")]
194 let file_encryptor = Self::get_file_encryptor(&properties, &schema_descriptor)?;
195
196 Self::start_file(&properties, &mut buf)?;
197 Ok(Self {
198 buf,
199 descr: Arc::new(schema_descriptor),
200 props: properties,
201 row_groups: vec![],
202 bloom_filters: vec![],
203 column_indexes: Vec::new(),
204 offset_indexes: Vec::new(),
205 row_group_index: 0,
206 kv_metadatas: Vec::new(),
207 finished: false,
208 #[cfg(feature = "encryption")]
209 file_encryptor,
210 })
211 }
212
213 #[cfg(feature = "encryption")]
214 fn get_file_encryptor(
215 properties: &WriterPropertiesPtr,
216 schema_descriptor: &SchemaDescriptor,
217 ) -> Result<Option<Arc<FileEncryptor>>> {
218 if let Some(file_encryption_properties) = properties.file_encryption_properties() {
219 file_encryption_properties.validate_encrypted_column_names(schema_descriptor)?;
220
221 Ok(Some(Arc::new(FileEncryptor::new(Arc::clone(
222 file_encryption_properties,
223 ))?)))
224 } else {
225 Ok(None)
226 }
227 }
228
229 pub fn next_row_group(&mut self) -> Result<SerializedRowGroupWriter<'_, W>> {
238 self.assert_previous_writer_closed()?;
239 let ordinal = self.row_group_index;
240
241 let ordinal: i16 = ordinal.try_into().map_err(|_| {
242 ParquetError::General(format!(
243 "Parquet does not support more than {} row groups per file (currently: {})",
244 i16::MAX,
245 ordinal
246 ))
247 })?;
248
249 self.row_group_index = self
250 .row_group_index
251 .checked_add(1)
252 .expect("SerializedFileWriter::row_group_index overflowed");
253
254 let bloom_filter_position = self.properties().bloom_filter_position();
255 let row_groups = &mut self.row_groups;
256 let row_bloom_filters = &mut self.bloom_filters;
257 let row_column_indexes = &mut self.column_indexes;
258 let row_offset_indexes = &mut self.offset_indexes;
259 let on_close = move |buf,
260 mut metadata,
261 row_group_bloom_filter,
262 row_group_column_index,
263 row_group_offset_index| {
264 row_bloom_filters.push(row_group_bloom_filter);
265 row_column_indexes.push(row_group_column_index);
266 row_offset_indexes.push(row_group_offset_index);
267 match bloom_filter_position {
269 BloomFilterPosition::AfterRowGroup => {
270 write_bloom_filters(buf, row_bloom_filters, &mut metadata)?
271 }
272 BloomFilterPosition::End => (),
273 };
274 row_groups.push(metadata);
275 Ok(())
276 };
277
278 let row_group_writer = SerializedRowGroupWriter::new(
279 self.descr.clone(),
280 self.props.clone(),
281 &mut self.buf,
282 ordinal,
283 Some(Box::new(on_close)),
284 );
285 #[cfg(feature = "encryption")]
286 let row_group_writer = row_group_writer.with_file_encryptor(self.file_encryptor.clone());
287
288 Ok(row_group_writer)
289 }
290
291 pub fn flushed_row_groups(&self) -> &[RowGroupMetaData] {
293 &self.row_groups
294 }
295
296 pub fn finish(&mut self) -> Result<ParquetMetaData> {
302 self.assert_previous_writer_closed()?;
303 let metadata = self.write_metadata()?;
304 self.buf.flush()?;
305 Ok(metadata)
306 }
307
308 pub fn close(mut self) -> Result<ParquetMetaData> {
310 self.finish()
311 }
312
313 #[cfg(not(feature = "encryption"))]
315 fn start_file(_properties: &WriterPropertiesPtr, buf: &mut TrackedWrite<W>) -> Result<()> {
316 buf.write_all(get_file_magic())?;
317 Ok(())
318 }
319
320 #[cfg(feature = "encryption")]
322 fn start_file(properties: &WriterPropertiesPtr, buf: &mut TrackedWrite<W>) -> Result<()> {
323 let magic = get_file_magic(properties.file_encryption_properties.as_ref());
324
325 buf.write_all(magic)?;
326 Ok(())
327 }
328
329 fn write_metadata(&mut self) -> Result<ParquetMetaData> {
332 self.finished = true;
333
334 for row_group in &mut self.row_groups {
336 write_bloom_filters(&mut self.buf, &mut self.bloom_filters, row_group)?;
337 }
338
339 let key_value_metadata = match self.props.key_value_metadata() {
340 Some(kv) => Some(kv.iter().chain(&self.kv_metadatas).cloned().collect()),
341 None if self.kv_metadatas.is_empty() => None,
342 None => Some(self.kv_metadatas.clone()),
343 };
344
345 let row_groups = std::mem::take(&mut self.row_groups);
347 let column_indexes = std::mem::take(&mut self.column_indexes);
348 let offset_indexes = std::mem::take(&mut self.offset_indexes);
349
350 let write_path_in_schema = self.props.write_path_in_schema();
351 let mut encoder = ThriftMetadataWriter::new(
352 &mut self.buf,
353 &self.descr,
354 row_groups,
355 Some(self.props.created_by().to_string()),
356 self.props.writer_version().as_num(),
357 write_path_in_schema,
358 );
359
360 #[cfg(feature = "encryption")]
361 {
362 encoder = encoder.with_file_encryptor(self.file_encryptor.clone());
363 }
364
365 if let Some(key_value_metadata) = key_value_metadata {
366 encoder = encoder.with_key_value_metadata(key_value_metadata)
367 }
368
369 encoder = encoder.with_column_indexes(column_indexes);
370 if !self.props.offset_index_disabled() {
371 encoder = encoder.with_offset_indexes(offset_indexes);
372 }
373 encoder.finish()
374 }
375
376 #[inline]
377 fn assert_previous_writer_closed(&self) -> Result<()> {
378 if self.finished {
379 return Err(general_err!("SerializedFileWriter already finished"));
380 }
381
382 if self.row_group_index != self.row_groups.len() {
383 Err(general_err!("Previous row group writer was not closed"))
384 } else {
385 Ok(())
386 }
387 }
388
389 pub fn append_key_value_metadata(&mut self, kv_metadata: KeyValue) {
391 self.kv_metadatas.push(kv_metadata);
392 }
393
394 pub fn schema_descr(&self) -> &SchemaDescriptor {
396 &self.descr
397 }
398
399 #[cfg(feature = "arrow")]
401 pub(crate) fn schema_descr_ptr(&self) -> &SchemaDescPtr {
402 &self.descr
403 }
404
405 pub fn properties(&self) -> &WriterPropertiesPtr {
407 &self.props
408 }
409
410 pub fn inner(&self) -> &W {
412 self.buf.inner()
413 }
414
415 pub fn write_all(&mut self, buf: &[u8]) -> std::io::Result<()> {
424 self.buf.write_all(buf)
425 }
426
427 pub fn flush(&mut self) -> std::io::Result<()> {
429 self.buf.flush()
430 }
431
432 pub fn inner_mut(&mut self) -> &mut W {
441 self.buf.inner_mut()
442 }
443
444 pub fn into_inner(mut self) -> Result<W> {
446 self.assert_previous_writer_closed()?;
447 let _ = self.write_metadata()?;
448
449 self.buf.into_inner()
450 }
451
452 pub fn bytes_written(&self) -> usize {
454 self.buf.bytes_written()
455 }
456
457 #[cfg(feature = "encryption")]
459 pub(crate) fn file_encryptor(&self) -> Option<Arc<FileEncryptor>> {
460 self.file_encryptor.clone()
461 }
462}
463
464fn write_bloom_filters<W: Write + Send>(
467 buf: &mut TrackedWrite<W>,
468 bloom_filters: &mut [Vec<Option<Sbbf>>],
469 row_group: &mut RowGroupMetaData,
470) -> Result<()> {
471 let row_group_idx: u16 = row_group
476 .ordinal()
477 .expect("Missing row group ordinal")
478 .try_into()
479 .map_err(|_| {
480 ParquetError::General(format!(
481 "Negative row group ordinal: {})",
482 row_group.ordinal().unwrap()
483 ))
484 })?;
485 let row_group_idx = row_group_idx as usize;
486 for (column_idx, column_chunk) in row_group.columns_mut().iter_mut().enumerate() {
487 if let Some(bloom_filter) = bloom_filters[row_group_idx][column_idx].take() {
488 let start_offset = buf.bytes_written();
489 bloom_filter.write(&mut *buf)?;
490 let end_offset = buf.bytes_written();
491 *column_chunk = column_chunk
493 .clone()
494 .into_builder()
495 .set_bloom_filter_offset(Some(start_offset as i64))
496 .set_bloom_filter_length(Some((end_offset - start_offset) as i32))
497 .build()?;
498 }
499 }
500 Ok(())
501}
502
503pub struct SerializedRowGroupWriter<'a, W: Write> {
516 descr: SchemaDescPtr,
517 props: WriterPropertiesPtr,
518 buf: &'a mut TrackedWrite<W>,
519 total_rows_written: Option<u64>,
520 total_bytes_written: u64,
521 total_uncompressed_bytes: i64,
522 column_index: usize,
523 row_group_metadata: Option<RowGroupMetaDataPtr>,
524 column_chunks: Vec<ColumnChunkMetaData>,
525 bloom_filters: Vec<Option<Sbbf>>,
526 column_indexes: Vec<Option<ColumnIndexMetaData>>,
527 offset_indexes: Vec<Option<OffsetIndexMetaData>>,
528 row_group_index: i16,
529 file_offset: i64,
530 on_close: Option<OnCloseRowGroup<'a, W>>,
531 #[cfg(feature = "encryption")]
532 file_encryptor: Option<Arc<FileEncryptor>>,
533}
534
535impl<'a, W: Write + Send> SerializedRowGroupWriter<'a, W> {
536 pub fn new(
545 schema_descr: SchemaDescPtr,
546 properties: WriterPropertiesPtr,
547 buf: &'a mut TrackedWrite<W>,
548 row_group_index: i16,
549 on_close: Option<OnCloseRowGroup<'a, W>>,
550 ) -> Self {
551 let num_columns = schema_descr.num_columns();
552 let file_offset = buf.bytes_written() as i64;
553 Self {
554 buf,
555 row_group_index,
556 file_offset,
557 on_close,
558 total_rows_written: None,
559 descr: schema_descr,
560 props: properties,
561 column_index: 0,
562 row_group_metadata: None,
563 column_chunks: Vec::with_capacity(num_columns),
564 bloom_filters: Vec::with_capacity(num_columns),
565 column_indexes: Vec::with_capacity(num_columns),
566 offset_indexes: Vec::with_capacity(num_columns),
567 total_bytes_written: 0,
568 total_uncompressed_bytes: 0,
569 #[cfg(feature = "encryption")]
570 file_encryptor: None,
571 }
572 }
573
574 #[cfg(feature = "encryption")]
575 pub(crate) fn with_file_encryptor(
577 mut self,
578 file_encryptor: Option<Arc<FileEncryptor>>,
579 ) -> Self {
580 self.file_encryptor = file_encryptor;
581 self
582 }
583
584 fn next_column_desc(&mut self) -> Option<ColumnDescPtr> {
586 let ret = self.descr.columns().get(self.column_index)?.clone();
587 self.column_index += 1;
588 Some(ret)
589 }
590
591 fn get_on_close(&mut self) -> (&mut TrackedWrite<W>, OnCloseColumnChunk<'_>) {
593 let total_bytes_written = &mut self.total_bytes_written;
594 let total_uncompressed_bytes = &mut self.total_uncompressed_bytes;
595 let total_rows_written = &mut self.total_rows_written;
596 let column_chunks = &mut self.column_chunks;
597 let column_indexes = &mut self.column_indexes;
598 let offset_indexes = &mut self.offset_indexes;
599 let bloom_filters = &mut self.bloom_filters;
600
601 let on_close = |r: ColumnCloseResult| {
602 *total_bytes_written += r.bytes_written;
604 *total_uncompressed_bytes += r.metadata.uncompressed_size();
605 column_chunks.push(r.metadata);
606 bloom_filters.push(r.bloom_filter);
607 column_indexes.push(r.column_index);
608 offset_indexes.push(r.offset_index);
609
610 if let Some(rows) = *total_rows_written {
611 if rows != r.rows_written {
612 return Err(general_err!(
613 "Incorrect number of rows, expected {} != {} rows",
614 rows,
615 r.rows_written
616 ));
617 }
618 } else {
619 *total_rows_written = Some(r.rows_written);
620 }
621
622 Ok(())
623 };
624 (self.buf, Box::new(on_close))
625 }
626
627 pub(crate) fn next_column_with_factory<'b, F, C>(&'b mut self, factory: F) -> Result<Option<C>>
630 where
631 F: FnOnce(
632 ColumnDescPtr,
633 WriterPropertiesPtr,
634 Box<dyn PageWriter + 'b>,
635 OnCloseColumnChunk<'b>,
636 ) -> Result<C>,
637 {
638 self.assert_previous_writer_closed()?;
639
640 let encryptor_context = self.get_page_encryptor_context();
641
642 Ok(match self.next_column_desc() {
643 Some(column) => {
644 let props = self.props.clone();
645 let (buf, on_close) = self.get_on_close();
646
647 let page_writer = SerializedPageWriter::new(buf);
648 let page_writer =
649 Self::set_page_writer_encryptor(&column, encryptor_context, page_writer)?;
650
651 Some(factory(
652 column,
653 props,
654 Box::new(page_writer),
655 Box::new(on_close),
656 )?)
657 }
658 None => None,
659 })
660 }
661
662 pub fn next_column(&mut self) -> Result<Option<SerializedColumnWriter<'_>>> {
666 self.next_column_with_factory(|descr, props, page_writer, on_close| {
667 let column_writer = get_column_writer(descr, props, page_writer);
668 Ok(SerializedColumnWriter::new(column_writer, Some(on_close)))
669 })
670 }
671
672 pub fn append_column<R: ChunkReader>(
687 &mut self,
688 reader: &R,
689 close: ColumnCloseResult,
690 ) -> Result<()> {
691 let metadata = &close.metadata;
694 let src_offset = metadata
695 .dictionary_page_offset()
696 .unwrap_or_else(|| metadata.data_page_offset());
697 let read = reader.get_read(src_offset as _)?;
698 self.append_column_from_read(read, close)
699 }
700
701 pub(crate) fn append_column_from_read<R: Read>(
711 &mut self,
712 read: R,
713 close: ColumnCloseResult,
714 ) -> Result<()> {
715 let (src_offset, src_length, write_offset) = self.begin_appended_column(&close)?;
716
717 let mut read = read.take(src_length as _);
718 let write_length = std::io::copy(&mut read, &mut self.buf)?;
719
720 if src_length as u64 != write_length {
721 return Err(general_err!(
722 "Failed to splice column data, expected {src_length} got {write_length}"
723 ));
724 }
725
726 self.finish_appended_column(close, src_offset, write_offset)
727 }
728
729 #[cfg(feature = "arrow")]
742 pub(crate) fn append_column_from_pages<I>(
743 &mut self,
744 pages: I,
745 close: ColumnCloseResult,
746 ) -> Result<()>
747 where
748 I: IntoIterator<Item = Result<Bytes>>,
749 {
750 let (src_offset, src_length, write_offset) = self.begin_appended_column(&close)?;
751
752 let mut write_length = 0u64;
753 for page in pages {
754 let page = page?;
755 self.buf.write_all(&page)?;
756 write_length += page.len() as u64;
757 }
758
759 if src_length as u64 != write_length {
760 return Err(general_err!(
761 "Failed to splice column data, expected {src_length} got {write_length}"
762 ));
763 }
764
765 self.finish_appended_column(close, src_offset, write_offset)
766 }
767
768 fn begin_appended_column(&mut self, close: &ColumnCloseResult) -> Result<(i64, i64, usize)> {
776 self.assert_previous_writer_closed()?;
777 let desc = self
778 .next_column_desc()
779 .ok_or_else(|| general_err!("exhausted columns in SerializedRowGroupWriter"))?;
780
781 let metadata = &close.metadata;
782
783 if metadata.column_descr() != desc.as_ref() {
784 return Err(general_err!(
785 "column descriptor mismatch, expected {:?} got {:?}",
786 desc,
787 metadata.column_descr()
788 ));
789 }
790
791 let src_offset = metadata
792 .dictionary_page_offset()
793 .unwrap_or_else(|| metadata.data_page_offset());
794 let src_length = metadata.compressed_size();
795 let write_offset = self.buf.bytes_written();
796 Ok((src_offset, src_length, write_offset))
797 }
798
799 fn finish_appended_column(
803 &mut self,
804 mut close: ColumnCloseResult,
805 src_offset: i64,
806 write_offset: usize,
807 ) -> Result<()> {
808 let metadata = close.metadata;
809 let src_dictionary_offset = metadata.dictionary_page_offset();
810 let src_data_offset = metadata.data_page_offset();
811
812 let map_offset = |x| x - src_offset + write_offset as i64;
813 let mut builder = ColumnChunkMetaData::builder(metadata.column_descr_ptr())
814 .set_compression_codec(metadata.compression_codec())
815 .set_encodings_mask(*metadata.encodings_mask())
816 .set_total_compressed_size(metadata.compressed_size())
817 .set_total_uncompressed_size(metadata.uncompressed_size())
818 .set_num_values(metadata.num_values())
819 .set_data_page_offset(map_offset(src_data_offset))
820 .set_dictionary_page_offset(src_dictionary_offset.map(map_offset))
821 .set_unencoded_byte_array_data_bytes(metadata.unencoded_byte_array_data_bytes());
822
823 if let Some(rep_hist) = metadata.repetition_level_histogram() {
824 builder = builder.set_repetition_level_histogram(Some(rep_hist.clone()))
825 }
826 if let Some(def_hist) = metadata.definition_level_histogram() {
827 builder = builder.set_definition_level_histogram(Some(def_hist.clone()))
828 }
829 if let Some(statistics) = metadata.statistics() {
830 builder = builder.set_statistics(statistics.clone())
831 }
832 if let Some(geo_statistics) = metadata.geo_statistics() {
833 builder = builder.set_geo_statistics(Box::new(geo_statistics.clone()))
834 }
835 if let Some(page_encoding_stats) = metadata.page_encoding_stats() {
836 builder = builder.set_page_encoding_stats(page_encoding_stats.clone())
837 }
838 builder = self.set_column_crypto_metadata(builder, &metadata);
839 close.metadata = builder.build()?;
840
841 if let Some(offsets) = close.offset_index.as_mut() {
842 for location in &mut offsets.page_locations {
843 location.offset = map_offset(location.offset)
844 }
845 }
846
847 let (_, on_close) = self.get_on_close();
848 on_close(close)
849 }
850
851 pub fn close(mut self) -> Result<RowGroupMetaDataPtr> {
853 if self.row_group_metadata.is_none() {
854 self.assert_previous_writer_closed()?;
855
856 let column_chunks = std::mem::take(&mut self.column_chunks);
857 let row_group_metadata = RowGroupMetaData::builder(self.descr.clone())
858 .set_column_metadata(column_chunks)
859 .set_total_byte_size(self.total_uncompressed_bytes)
860 .set_num_rows(self.total_rows_written.unwrap_or(0) as i64)
861 .set_sorting_columns(self.props.sorting_columns().cloned())
862 .set_ordinal(self.row_group_index)
863 .set_file_offset(self.file_offset)
864 .build()?;
865
866 self.row_group_metadata = Some(Arc::new(row_group_metadata.clone()));
867
868 if let Some(on_close) = self.on_close.take() {
869 on_close(
870 self.buf,
871 row_group_metadata,
872 self.bloom_filters,
873 self.column_indexes,
874 self.offset_indexes,
875 )?
876 }
877 }
878
879 let metadata = self.row_group_metadata.as_ref().unwrap().clone();
880 Ok(metadata)
881 }
882
883 #[cfg(feature = "encryption")]
885 fn set_column_crypto_metadata(
886 &self,
887 builder: ColumnChunkMetaDataBuilder,
888 metadata: &ColumnChunkMetaData,
889 ) -> ColumnChunkMetaDataBuilder {
890 if let Some(file_encryptor) = self.file_encryptor.as_ref() {
891 builder.set_column_crypto_metadata(get_column_crypto_metadata(
892 file_encryptor.properties(),
893 &metadata.column_descr_ptr(),
894 ))
895 } else {
896 builder
897 }
898 }
899
900 #[cfg(feature = "encryption")]
902 fn get_page_encryptor_context(&self) -> PageEncryptorContext {
903 PageEncryptorContext {
904 file_encryptor: self.file_encryptor.clone(),
905 row_group_index: self.row_group_index as usize,
906 column_index: self.column_index,
907 }
908 }
909
910 #[cfg(feature = "encryption")]
912 fn set_page_writer_encryptor<'b>(
913 column: &ColumnDescPtr,
914 context: PageEncryptorContext,
915 page_writer: SerializedPageWriter<'b, W>,
916 ) -> Result<SerializedPageWriter<'b, W>> {
917 let page_encryptor = PageEncryptor::create_if_column_encrypted(
918 &context.file_encryptor,
919 context.row_group_index,
920 context.column_index,
921 &column.path().string(),
922 )?;
923
924 Ok(page_writer.with_page_encryptor(page_encryptor))
925 }
926
927 #[cfg(not(feature = "encryption"))]
929 fn set_column_crypto_metadata(
930 &self,
931 builder: ColumnChunkMetaDataBuilder,
932 _metadata: &ColumnChunkMetaData,
933 ) -> ColumnChunkMetaDataBuilder {
934 builder
935 }
936
937 #[cfg(not(feature = "encryption"))]
938 fn get_page_encryptor_context(&self) -> PageEncryptorContext {
939 PageEncryptorContext {}
940 }
941
942 #[cfg(not(feature = "encryption"))]
944 fn set_page_writer_encryptor<'b>(
945 _column: &ColumnDescPtr,
946 _context: PageEncryptorContext,
947 page_writer: SerializedPageWriter<'b, W>,
948 ) -> Result<SerializedPageWriter<'b, W>> {
949 Ok(page_writer)
950 }
951
952 #[inline]
953 fn assert_previous_writer_closed(&self) -> Result<()> {
954 if self.column_index != self.column_chunks.len() {
955 Err(general_err!("Previous column writer was not closed"))
956 } else {
957 Ok(())
958 }
959 }
960}
961
962#[cfg(feature = "encryption")]
964struct PageEncryptorContext {
965 file_encryptor: Option<Arc<FileEncryptor>>,
966 row_group_index: usize,
967 column_index: usize,
968}
969
970#[cfg(not(feature = "encryption"))]
971struct PageEncryptorContext {}
972
973pub struct SerializedColumnWriter<'a> {
975 inner: ColumnWriter<'a>,
976 on_close: Option<OnCloseColumnChunk<'a>>,
977}
978
979impl<'a> SerializedColumnWriter<'a> {
980 pub fn new(inner: ColumnWriter<'a>, on_close: Option<OnCloseColumnChunk<'a>>) -> Self {
983 Self { inner, on_close }
984 }
985
986 pub fn untyped(&mut self) -> &mut ColumnWriter<'a> {
988 &mut self.inner
989 }
990
991 pub fn typed<T: DataType>(&mut self) -> &mut ColumnWriterImpl<'a, T> {
993 get_typed_column_writer_mut(&mut self.inner)
994 }
995
996 pub fn close(mut self) -> Result<()> {
998 let r = self.inner.close()?;
999 if let Some(on_close) = self.on_close.take() {
1000 on_close(r)?
1001 }
1002
1003 Ok(())
1004 }
1005}
1006
1007pub struct SerializedPageWriter<'a, W: Write> {
1012 sink: &'a mut TrackedWrite<W>,
1013 #[cfg(feature = "encryption")]
1014 page_encryptor: Option<PageEncryptor>,
1015}
1016
1017impl<'a, W: Write> SerializedPageWriter<'a, W> {
1018 pub fn new(sink: &'a mut TrackedWrite<W>) -> Self {
1020 Self {
1021 sink,
1022 #[cfg(feature = "encryption")]
1023 page_encryptor: None,
1024 }
1025 }
1026
1027 #[inline]
1030 fn serialize_page_header(&mut self, header: PageHeader) -> Result<usize> {
1031 let start_pos = self.sink.bytes_written();
1032 match self.page_encryptor_and_sink_mut() {
1033 Some((page_encryptor, sink)) => {
1034 page_encryptor.encrypt_page_header(&header, sink)?;
1035 }
1036 None => {
1037 let mut protocol = ThriftCompactOutputProtocol::new(&mut self.sink);
1038 header.write_thrift(&mut protocol)?;
1039 }
1040 }
1041 Ok(self.sink.bytes_written() - start_pos)
1042 }
1043}
1044
1045#[cfg(feature = "encryption")]
1046impl<'a, W: Write> SerializedPageWriter<'a, W> {
1047 fn with_page_encryptor(mut self, page_encryptor: Option<PageEncryptor>) -> Self {
1049 self.page_encryptor = page_encryptor;
1050 self
1051 }
1052
1053 fn page_encryptor_mut(&mut self) -> Option<&mut PageEncryptor> {
1054 self.page_encryptor.as_mut()
1055 }
1056
1057 fn page_encryptor_and_sink_mut(
1058 &mut self,
1059 ) -> Option<(&mut PageEncryptor, &mut &'a mut TrackedWrite<W>)> {
1060 self.page_encryptor.as_mut().map(|pe| (pe, &mut self.sink))
1061 }
1062}
1063
1064#[cfg(not(feature = "encryption"))]
1065impl<'a, W: Write> SerializedPageWriter<'a, W> {
1066 fn page_encryptor_mut(&mut self) -> Option<&mut PageEncryptor> {
1067 None
1068 }
1069
1070 fn page_encryptor_and_sink_mut(
1071 &mut self,
1072 ) -> Option<(&mut PageEncryptor, &mut &'a mut TrackedWrite<W>)> {
1073 None
1074 }
1075}
1076
1077impl<W: Write + Send> PageWriter for SerializedPageWriter<'_, W> {
1078 fn write_page(&mut self, page: CompressedPage) -> Result<PageWriteSpec> {
1079 let page = match self.page_encryptor_mut() {
1080 Some(page_encryptor) => page_encryptor.encrypt_compressed_page(page)?,
1081 None => page,
1082 };
1083
1084 let page_type = page.page_type();
1085 let start_pos = self.sink.bytes_written() as u64;
1086
1087 let page_header = page.to_thrift_header()?;
1088 let header_size = self.serialize_page_header(page_header)?;
1089
1090 self.sink.write_all(page.data())?;
1091
1092 let mut spec = PageWriteSpec::new();
1093 spec.page_type = page_type;
1094 spec.uncompressed_size = page.uncompressed_size() + header_size;
1095 spec.compressed_size = page.compressed_size() + header_size;
1096 spec.offset = start_pos;
1097 spec.bytes_written = self.sink.bytes_written() as u64 - start_pos;
1098 spec.num_values = page.num_values();
1099
1100 if let Some(page_encryptor) = self.page_encryptor_mut() {
1101 if page.compressed_page().is_data_page() {
1102 page_encryptor.increment_page();
1103 }
1104 }
1105 Ok(spec)
1106 }
1107
1108 fn close(&mut self) -> Result<()> {
1109 self.sink.flush()?;
1110 Ok(())
1111 }
1112}
1113
1114#[cfg(feature = "encryption")]
1117pub(crate) fn get_file_magic(
1118 file_encryption_properties: Option<&Arc<FileEncryptionProperties>>,
1119) -> &'static [u8; 4] {
1120 match file_encryption_properties.as_ref() {
1121 Some(encryption_properties) if encryption_properties.encrypt_footer() => {
1122 &PARQUET_MAGIC_ENCR_FOOTER
1123 }
1124 _ => &PARQUET_MAGIC,
1125 }
1126}
1127
1128#[cfg(not(feature = "encryption"))]
1129pub(crate) fn get_file_magic() -> &'static [u8; 4] {
1130 &PARQUET_MAGIC
1131}
1132
1133#[cfg(test)]
1134mod tests {
1135 use super::*;
1136
1137 #[cfg(feature = "arrow")]
1138 use arrow_array::RecordBatchReader;
1139 use bytes::Bytes;
1140 use std::fs::File;
1141
1142 #[cfg(feature = "arrow")]
1143 use crate::arrow::ArrowWriter;
1144 #[cfg(feature = "arrow")]
1145 use crate::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
1146 use crate::basic::{
1147 ColumnOrder, Compression, ConvertedType, Encoding, LogicalType, Repetition, SortOrder, Type,
1148 };
1149 use crate::column::page::{Page, PageReader};
1150 use crate::column::reader::get_typed_column_reader;
1151 use crate::compression::{Codec, CodecOptionsBuilder, create_codec};
1152 use crate::data_type::{BoolType, ByteArrayType, Int32Type};
1153 use crate::file::page_index::column_index::ColumnIndexMetaData;
1154 use crate::file::properties::EnabledStatistics;
1155 use crate::file::serialized_reader::ReadOptionsBuilder;
1156 use crate::file::statistics::{from_thrift_page_stats, page_stats_to_thrift};
1157 use crate::file::{
1158 properties::{ReaderProperties, WriterProperties, WriterVersion},
1159 reader::{FileReader, SerializedFileReader, SerializedPageReader},
1160 statistics::Statistics,
1161 };
1162 use crate::record::{Row, RowAccessor};
1163 use crate::schema::parser::parse_message_type;
1164 use crate::schema::types;
1165 use crate::schema::types::{ColumnDescriptor, ColumnPath};
1166 use crate::util::test_common::file_util::get_test_file;
1167 use crate::util::test_common::rand_gen::RandGen;
1168
1169 #[test]
1170 fn test_row_group_writer_error_not_all_columns_written() {
1171 let file = tempfile::tempfile().unwrap();
1172 let schema = Arc::new(
1173 types::Type::group_type_builder("schema")
1174 .with_fields(vec![Arc::new(
1175 types::Type::primitive_type_builder("col1", Type::INT32)
1176 .build()
1177 .unwrap(),
1178 )])
1179 .build()
1180 .unwrap(),
1181 );
1182 let props = Default::default();
1183 let mut writer = SerializedFileWriter::new(file, schema, props).unwrap();
1184 let row_group_writer = writer.next_row_group().unwrap();
1185 let res = row_group_writer.close();
1186 assert!(res.is_err());
1187 if let Err(err) = res {
1188 assert_eq!(
1189 format!("{err}"),
1190 "Parquet error: Column length mismatch: 1 != 0"
1191 );
1192 }
1193 }
1194
1195 #[test]
1196 fn test_row_group_writer_num_records_mismatch() {
1197 let file = tempfile::tempfile().unwrap();
1198 let schema = Arc::new(
1199 types::Type::group_type_builder("schema")
1200 .with_fields(vec![
1201 Arc::new(
1202 types::Type::primitive_type_builder("col1", Type::INT32)
1203 .with_repetition(Repetition::REQUIRED)
1204 .build()
1205 .unwrap(),
1206 ),
1207 Arc::new(
1208 types::Type::primitive_type_builder("col2", Type::INT32)
1209 .with_repetition(Repetition::REQUIRED)
1210 .build()
1211 .unwrap(),
1212 ),
1213 ])
1214 .build()
1215 .unwrap(),
1216 );
1217 let props = Default::default();
1218 let mut writer = SerializedFileWriter::new(file, schema, props).unwrap();
1219 let mut row_group_writer = writer.next_row_group().unwrap();
1220
1221 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
1222 col_writer
1223 .typed::<Int32Type>()
1224 .write_batch(&[1, 2, 3], None, None)
1225 .unwrap();
1226 col_writer.close().unwrap();
1227
1228 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
1229 col_writer
1230 .typed::<Int32Type>()
1231 .write_batch(&[1, 2], None, None)
1232 .unwrap();
1233
1234 let err = col_writer.close().unwrap_err();
1235 assert_eq!(
1236 err.to_string(),
1237 "Parquet error: Incorrect number of rows, expected 3 != 2 rows"
1238 );
1239 }
1240
1241 #[test]
1242 fn test_file_writer_empty_file() {
1243 let file = tempfile::tempfile().unwrap();
1244
1245 let schema = Arc::new(
1246 types::Type::group_type_builder("schema")
1247 .with_fields(vec![Arc::new(
1248 types::Type::primitive_type_builder("col1", Type::INT32)
1249 .build()
1250 .unwrap(),
1251 )])
1252 .build()
1253 .unwrap(),
1254 );
1255 let props = Default::default();
1256 let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1257 writer.close().unwrap();
1258
1259 let reader = SerializedFileReader::new(file).unwrap();
1260 assert_eq!(reader.get_row_iter(None).unwrap().count(), 0);
1261 }
1262
1263 #[test]
1264 fn test_file_writer_column_orders_populated() {
1265 let file = tempfile::tempfile().unwrap();
1266
1267 let schema = Arc::new(
1268 types::Type::group_type_builder("schema")
1269 .with_fields(vec![
1270 Arc::new(
1271 types::Type::primitive_type_builder("col1", Type::INT32)
1272 .build()
1273 .unwrap(),
1274 ),
1275 Arc::new(
1276 types::Type::primitive_type_builder("col2", Type::FIXED_LEN_BYTE_ARRAY)
1277 .with_converted_type(ConvertedType::INTERVAL)
1278 .with_length(12)
1279 .build()
1280 .unwrap(),
1281 ),
1282 Arc::new(
1283 types::Type::group_type_builder("nested")
1284 .with_repetition(Repetition::REQUIRED)
1285 .with_fields(vec![
1286 Arc::new(
1287 types::Type::primitive_type_builder(
1288 "col3",
1289 Type::FIXED_LEN_BYTE_ARRAY,
1290 )
1291 .with_logical_type(Some(LogicalType::Float16))
1292 .with_length(2)
1293 .build()
1294 .unwrap(),
1295 ),
1296 Arc::new(
1297 types::Type::primitive_type_builder("col4", Type::BYTE_ARRAY)
1298 .with_logical_type(Some(LogicalType::String))
1299 .build()
1300 .unwrap(),
1301 ),
1302 ])
1303 .build()
1304 .unwrap(),
1305 ),
1306 ])
1307 .build()
1308 .unwrap(),
1309 );
1310
1311 let props = Default::default();
1312 let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1313 writer.close().unwrap();
1314
1315 let reader = SerializedFileReader::new(file).unwrap();
1316
1317 let expected = vec![
1319 ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED),
1321 ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNDEFINED),
1323 ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED),
1325 ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::UNSIGNED),
1327 ];
1328 let actual = reader.metadata().file_metadata().column_orders();
1329
1330 assert!(actual.is_some());
1331 let actual = actual.unwrap();
1332 assert_eq!(*actual, expected);
1333 }
1334
1335 #[test]
1336 fn test_file_writer_with_metadata() {
1337 let file = tempfile::tempfile().unwrap();
1338
1339 let schema = Arc::new(
1340 types::Type::group_type_builder("schema")
1341 .with_fields(vec![Arc::new(
1342 types::Type::primitive_type_builder("col1", Type::INT32)
1343 .build()
1344 .unwrap(),
1345 )])
1346 .build()
1347 .unwrap(),
1348 );
1349 let props = Arc::new(
1350 WriterProperties::builder()
1351 .set_key_value_metadata(Some(vec![KeyValue::new(
1352 "key".to_string(),
1353 "value".to_string(),
1354 )]))
1355 .build(),
1356 );
1357 let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1358 writer.close().unwrap();
1359
1360 let reader = SerializedFileReader::new(file).unwrap();
1361 assert_eq!(
1362 reader
1363 .metadata()
1364 .file_metadata()
1365 .key_value_metadata()
1366 .to_owned()
1367 .unwrap()
1368 .len(),
1369 1
1370 );
1371 }
1372
1373 #[test]
1374 fn test_file_writer_v2_with_metadata() {
1375 let file = tempfile::tempfile().unwrap();
1376 let field_logical_type = Some(LogicalType::integer(8, false));
1377 let field = Arc::new(
1378 types::Type::primitive_type_builder("col1", Type::INT32)
1379 .with_logical_type(field_logical_type.clone())
1380 .with_converted_type(field_logical_type.into())
1381 .build()
1382 .unwrap(),
1383 );
1384 let schema = Arc::new(
1385 types::Type::group_type_builder("schema")
1386 .with_fields(vec![field.clone()])
1387 .build()
1388 .unwrap(),
1389 );
1390 let props = Arc::new(
1391 WriterProperties::builder()
1392 .set_key_value_metadata(Some(vec![KeyValue::new(
1393 "key".to_string(),
1394 "value".to_string(),
1395 )]))
1396 .set_writer_version(WriterVersion::PARQUET_2_0)
1397 .build(),
1398 );
1399 let writer = SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1400 writer.close().unwrap();
1401
1402 let reader = SerializedFileReader::new(file).unwrap();
1403
1404 assert_eq!(
1405 reader
1406 .metadata()
1407 .file_metadata()
1408 .key_value_metadata()
1409 .to_owned()
1410 .unwrap()
1411 .len(),
1412 1
1413 );
1414
1415 let fields = reader.metadata().file_metadata().schema().get_fields();
1417 assert_eq!(fields.len(), 1);
1418 assert_eq!(fields[0], field);
1419 }
1420
1421 #[test]
1422 fn test_file_writer_with_sorting_columns_metadata() {
1423 let file = tempfile::tempfile().unwrap();
1424
1425 let schema = Arc::new(
1426 types::Type::group_type_builder("schema")
1427 .with_fields(vec![
1428 Arc::new(
1429 types::Type::primitive_type_builder("col1", Type::INT32)
1430 .build()
1431 .unwrap(),
1432 ),
1433 Arc::new(
1434 types::Type::primitive_type_builder("col2", Type::INT32)
1435 .build()
1436 .unwrap(),
1437 ),
1438 ])
1439 .build()
1440 .unwrap(),
1441 );
1442 let expected_result = Some(vec![SortingColumn {
1443 column_idx: 0,
1444 descending: false,
1445 nulls_first: true,
1446 }]);
1447 let props = Arc::new(
1448 WriterProperties::builder()
1449 .set_key_value_metadata(Some(vec![KeyValue::new(
1450 "key".to_string(),
1451 "value".to_string(),
1452 )]))
1453 .set_sorting_columns(expected_result.clone())
1454 .build(),
1455 );
1456 let mut writer =
1457 SerializedFileWriter::new(file.try_clone().unwrap(), schema, props).unwrap();
1458 let mut row_group_writer = writer.next_row_group().expect("get row group writer");
1459
1460 let col_writer = row_group_writer.next_column().unwrap().unwrap();
1461 col_writer.close().unwrap();
1462
1463 let col_writer = row_group_writer.next_column().unwrap().unwrap();
1464 col_writer.close().unwrap();
1465
1466 row_group_writer.close().unwrap();
1467 writer.close().unwrap();
1468
1469 let reader = SerializedFileReader::new(file).unwrap();
1470 let result: Vec<Option<&Vec<SortingColumn>>> = reader
1471 .metadata()
1472 .row_groups()
1473 .iter()
1474 .map(|f| f.sorting_columns())
1475 .collect();
1476 assert_eq!(expected_result.as_ref(), result[0]);
1478 }
1479
1480 #[test]
1481 fn test_file_writer_empty_row_groups() {
1482 let file = tempfile::tempfile().unwrap();
1483 test_file_roundtrip(file, vec![]);
1484 }
1485
1486 #[test]
1487 fn test_file_writer_single_row_group() {
1488 let file = tempfile::tempfile().unwrap();
1489 test_file_roundtrip(file, vec![vec![1, 2, 3, 4, 5]]);
1490 }
1491
1492 #[test]
1493 fn test_file_writer_multiple_row_groups() {
1494 let file = tempfile::tempfile().unwrap();
1495 test_file_roundtrip(
1496 file,
1497 vec![
1498 vec![1, 2, 3, 4, 5],
1499 vec![1, 2, 3],
1500 vec![1],
1501 vec![1, 2, 3, 4, 5, 6],
1502 ],
1503 );
1504 }
1505
1506 #[test]
1507 fn test_file_writer_multiple_large_row_groups() {
1508 let file = tempfile::tempfile().unwrap();
1509 test_file_roundtrip(
1510 file,
1511 vec![vec![123; 1024], vec![124; 1000], vec![125; 15], vec![]],
1512 );
1513 }
1514
1515 #[test]
1516 fn test_page_writer_data_pages() {
1517 let pages = [
1518 Page::DataPage {
1519 buf: Bytes::from(vec![1, 2, 3, 4, 5, 6, 7, 8]),
1520 num_values: 10,
1521 encoding: Encoding::DELTA_BINARY_PACKED,
1522 def_level_encoding: Encoding::RLE,
1523 rep_level_encoding: Encoding::RLE,
1524 statistics: Some(Statistics::int32(Some(1), Some(3), None, Some(7), true)),
1525 },
1526 Page::DataPageV2 {
1527 buf: Bytes::from(vec![4; 128]),
1528 num_values: 10,
1529 encoding: Encoding::DELTA_BINARY_PACKED,
1530 num_nulls: 2,
1531 num_rows: 12,
1532 def_levels_byte_len: 24,
1533 rep_levels_byte_len: 32,
1534 is_compressed: false,
1535 statistics: Some(Statistics::int32(Some(1), Some(3), None, Some(7), true)),
1536 },
1537 ];
1538
1539 test_page_roundtrip(&pages[..], Compression::SNAPPY, Type::INT32);
1540 test_page_roundtrip(&pages[..], Compression::UNCOMPRESSED, Type::INT32);
1541 }
1542
1543 #[test]
1544 fn test_page_writer_dict_pages() {
1545 let pages = [
1546 Page::DictionaryPage {
1547 buf: Bytes::from(vec![1, 2, 3, 4, 5]),
1548 num_values: 5,
1549 encoding: Encoding::RLE_DICTIONARY,
1550 is_sorted: false,
1551 },
1552 Page::DataPage {
1553 buf: Bytes::from(vec![1, 2, 3, 4, 5, 6, 7, 8]),
1554 num_values: 10,
1555 encoding: Encoding::DELTA_BINARY_PACKED,
1556 def_level_encoding: Encoding::RLE,
1557 rep_level_encoding: Encoding::RLE,
1558 statistics: Some(Statistics::int32(Some(1), Some(3), None, Some(7), true)),
1559 },
1560 Page::DataPageV2 {
1561 buf: Bytes::from(vec![4; 128]),
1562 num_values: 10,
1563 encoding: Encoding::DELTA_BINARY_PACKED,
1564 num_nulls: 2,
1565 num_rows: 12,
1566 def_levels_byte_len: 24,
1567 rep_levels_byte_len: 32,
1568 is_compressed: false,
1569 statistics: None,
1570 },
1571 ];
1572
1573 test_page_roundtrip(&pages[..], Compression::SNAPPY, Type::INT32);
1574 test_page_roundtrip(&pages[..], Compression::UNCOMPRESSED, Type::INT32);
1575 }
1576
1577 fn test_page_roundtrip(pages: &[Page], codec: Compression, physical_type: Type) {
1581 let mut compressed_pages = vec![];
1582 let mut total_num_values = 0i64;
1583 let codec_options = CodecOptionsBuilder::default()
1584 .set_backward_compatible_lz4(false)
1585 .build();
1586 let mut compressor = create_codec(codec, &codec_options).unwrap();
1587
1588 for page in pages {
1589 let uncompressed_len = page.buffer().len();
1590
1591 let compressed_page = match *page {
1592 Page::DataPage {
1593 ref buf,
1594 num_values,
1595 encoding,
1596 def_level_encoding,
1597 rep_level_encoding,
1598 ref statistics,
1599 } => {
1600 total_num_values += num_values as i64;
1601 let output_buf = compress_helper(compressor.as_mut(), buf);
1602
1603 Page::DataPage {
1604 buf: Bytes::from(output_buf),
1605 num_values,
1606 encoding,
1607 def_level_encoding,
1608 rep_level_encoding,
1609 statistics: from_thrift_page_stats(
1610 physical_type,
1611 page_stats_to_thrift(statistics.as_ref()),
1612 )
1613 .unwrap(),
1614 }
1615 }
1616 Page::DataPageV2 {
1617 ref buf,
1618 num_values,
1619 encoding,
1620 num_nulls,
1621 num_rows,
1622 def_levels_byte_len,
1623 rep_levels_byte_len,
1624 ref statistics,
1625 ..
1626 } => {
1627 total_num_values += num_values as i64;
1628 let offset = (def_levels_byte_len + rep_levels_byte_len) as usize;
1629 let cmp_buf = compress_helper(compressor.as_mut(), &buf[offset..]);
1630 let mut output_buf = Vec::from(&buf[..offset]);
1631 output_buf.extend_from_slice(&cmp_buf[..]);
1632
1633 Page::DataPageV2 {
1634 buf: Bytes::from(output_buf),
1635 num_values,
1636 encoding,
1637 num_nulls,
1638 num_rows,
1639 def_levels_byte_len,
1640 rep_levels_byte_len,
1641 is_compressed: compressor.is_some(),
1642 statistics: from_thrift_page_stats(
1643 physical_type,
1644 page_stats_to_thrift(statistics.as_ref()),
1645 )
1646 .unwrap(),
1647 }
1648 }
1649 Page::DictionaryPage {
1650 ref buf,
1651 num_values,
1652 encoding,
1653 is_sorted,
1654 } => {
1655 let output_buf = compress_helper(compressor.as_mut(), buf);
1656
1657 Page::DictionaryPage {
1658 buf: Bytes::from(output_buf),
1659 num_values,
1660 encoding,
1661 is_sorted,
1662 }
1663 }
1664 };
1665
1666 let compressed_page = CompressedPage::new(compressed_page, uncompressed_len);
1667 compressed_pages.push(compressed_page);
1668 }
1669
1670 let mut buffer: Vec<u8> = vec![];
1671 let mut result_pages: Vec<Page> = vec![];
1672 {
1673 let mut writer = TrackedWrite::new(&mut buffer);
1674 let mut page_writer = SerializedPageWriter::new(&mut writer);
1675
1676 for page in compressed_pages {
1677 page_writer.write_page(page).unwrap();
1678 }
1679 page_writer.close().unwrap();
1680 }
1681 {
1682 let reader = bytes::Bytes::from(buffer);
1683
1684 let t = types::Type::primitive_type_builder("t", physical_type)
1685 .build()
1686 .unwrap();
1687
1688 let desc = ColumnDescriptor::new(Arc::new(t), 0, 0, ColumnPath::new(vec![]));
1689 let meta = ColumnChunkMetaData::builder(Arc::new(desc))
1690 .set_compression_codec(codec.into())
1691 .set_total_compressed_size(reader.len() as i64)
1692 .set_num_values(total_num_values)
1693 .build()
1694 .unwrap();
1695
1696 let props = ReaderProperties::builder()
1697 .set_backward_compatible_lz4(false)
1698 .set_read_page_statistics(true)
1699 .build();
1700 let mut page_reader = SerializedPageReader::new_with_properties(
1701 Arc::new(reader),
1702 &meta,
1703 total_num_values as usize,
1704 None,
1705 Arc::new(props),
1706 )
1707 .unwrap();
1708
1709 while let Some(page) = page_reader.get_next_page().unwrap() {
1710 result_pages.push(page);
1711 }
1712 }
1713
1714 assert_eq!(result_pages.len(), pages.len());
1715 for i in 0..result_pages.len() {
1716 assert_page(&result_pages[i], &pages[i]);
1717 }
1718 }
1719
1720 fn compress_helper(compressor: Option<&mut Box<dyn Codec>>, data: &[u8]) -> Vec<u8> {
1722 let mut output_buf = vec![];
1723 if let Some(cmpr) = compressor {
1724 cmpr.compress(data, &mut output_buf).unwrap();
1725 } else {
1726 output_buf.extend_from_slice(data);
1727 }
1728 output_buf
1729 }
1730
1731 fn assert_page(left: &Page, right: &Page) {
1733 assert_eq!(left.page_type(), right.page_type());
1734 assert_eq!(&left.buffer(), &right.buffer());
1735 assert_eq!(left.num_values(), right.num_values());
1736 assert_eq!(left.encoding(), right.encoding());
1737 assert_eq!(
1738 page_stats_to_thrift(left.statistics()),
1739 page_stats_to_thrift(right.statistics())
1740 );
1741 }
1742
1743 fn test_roundtrip_i32<W, R>(
1745 file: W,
1746 data: Vec<Vec<i32>>,
1747 compression: Compression,
1748 ) -> ParquetMetaData
1749 where
1750 W: Write + Send,
1751 R: ChunkReader + From<W> + 'static,
1752 {
1753 test_roundtrip::<W, R, Int32Type, _>(file, data, |r| r.get_int(0).unwrap(), compression)
1754 }
1755
1756 fn test_roundtrip<W, R, D, F>(
1759 mut file: W,
1760 data: Vec<Vec<D::T>>,
1761 value: F,
1762 compression: Compression,
1763 ) -> ParquetMetaData
1764 where
1765 W: Write + Send,
1766 R: ChunkReader + From<W> + 'static,
1767 D: DataType,
1768 F: Fn(Row) -> D::T,
1769 {
1770 let schema = Arc::new(
1771 types::Type::group_type_builder("schema")
1772 .with_fields(vec![Arc::new(
1773 types::Type::primitive_type_builder("col1", D::get_physical_type())
1774 .with_repetition(Repetition::REQUIRED)
1775 .build()
1776 .unwrap(),
1777 )])
1778 .build()
1779 .unwrap(),
1780 );
1781 let props = Arc::new(
1782 WriterProperties::builder()
1783 .set_compression(compression)
1784 .build(),
1785 );
1786 let mut file_writer = SerializedFileWriter::new(&mut file, schema, props).unwrap();
1787 let mut rows: i64 = 0;
1788
1789 for (idx, subset) in data.iter().enumerate() {
1790 let row_group_file_offset = file_writer.buf.bytes_written();
1791 let mut row_group_writer = file_writer.next_row_group().unwrap();
1792 if let Some(mut writer) = row_group_writer.next_column().unwrap() {
1793 rows += writer
1794 .typed::<D>()
1795 .write_batch(&subset[..], None, None)
1796 .unwrap() as i64;
1797 writer.close().unwrap();
1798 }
1799 let last_group = row_group_writer.close().unwrap();
1800 let flushed = file_writer.flushed_row_groups();
1801 assert_eq!(flushed.len(), idx + 1);
1802 assert_eq!(Some(idx as i16), last_group.ordinal());
1803 assert_eq!(Some(row_group_file_offset as i64), last_group.file_offset());
1804 assert_eq!(&flushed[idx], last_group.as_ref());
1805 }
1806 let file_metadata = file_writer.close().unwrap();
1807
1808 let reader = SerializedFileReader::new(R::from(file)).unwrap();
1809 assert_eq!(reader.num_row_groups(), data.len());
1810 assert_eq!(
1811 reader.metadata().file_metadata().num_rows(),
1812 rows,
1813 "row count in metadata not equal to number of rows written"
1814 );
1815 for (i, item) in data.iter().enumerate().take(reader.num_row_groups()) {
1816 let row_group_reader = reader.get_row_group(i).unwrap();
1817 let iter = row_group_reader.get_row_iter(None).unwrap();
1818 let res: Vec<_> = iter.map(|row| row.unwrap()).map(&value).collect();
1819 let row_group_size = row_group_reader.metadata().total_byte_size();
1820 let uncompressed_size: i64 = row_group_reader
1821 .metadata()
1822 .columns()
1823 .iter()
1824 .map(|v| v.uncompressed_size())
1825 .sum();
1826 assert_eq!(row_group_size, uncompressed_size);
1827 assert_eq!(res, *item);
1828 }
1829 file_metadata
1830 }
1831
1832 fn test_file_roundtrip(file: File, data: Vec<Vec<i32>>) -> ParquetMetaData {
1835 test_roundtrip_i32::<File, File>(file, data, Compression::UNCOMPRESSED)
1836 }
1837
1838 #[test]
1839 fn test_bytes_writer_empty_row_groups() {
1840 test_bytes_roundtrip(vec![], Compression::UNCOMPRESSED);
1841 }
1842
1843 #[test]
1844 fn test_bytes_writer_single_row_group() {
1845 test_bytes_roundtrip(vec![vec![1, 2, 3, 4, 5]], Compression::UNCOMPRESSED);
1846 }
1847
1848 #[test]
1849 fn test_bytes_writer_multiple_row_groups() {
1850 test_bytes_roundtrip(
1851 vec![
1852 vec![1, 2, 3, 4, 5],
1853 vec![1, 2, 3],
1854 vec![1],
1855 vec![1, 2, 3, 4, 5, 6],
1856 ],
1857 Compression::UNCOMPRESSED,
1858 );
1859 }
1860
1861 #[test]
1862 fn test_bytes_writer_single_row_group_compressed() {
1863 test_bytes_roundtrip(vec![vec![1, 2, 3, 4, 5]], Compression::SNAPPY);
1864 }
1865
1866 #[test]
1867 fn test_bytes_writer_multiple_row_groups_compressed() {
1868 test_bytes_roundtrip(
1869 vec![
1870 vec![1, 2, 3, 4, 5],
1871 vec![1, 2, 3],
1872 vec![1],
1873 vec![1, 2, 3, 4, 5, 6],
1874 ],
1875 Compression::SNAPPY,
1876 );
1877 }
1878
1879 fn test_bytes_roundtrip(data: Vec<Vec<i32>>, compression: Compression) {
1880 test_roundtrip_i32::<Vec<u8>, Bytes>(Vec::with_capacity(1024), data, compression);
1881 }
1882
1883 #[test]
1884 fn test_boolean_roundtrip() {
1885 let my_bool_values: Vec<_> = (0..2049).map(|idx| idx % 2 == 0).collect();
1886 test_roundtrip::<Vec<u8>, Bytes, BoolType, _>(
1887 Vec::with_capacity(1024),
1888 vec![my_bool_values],
1889 |r| r.get_bool(0).unwrap(),
1890 Compression::UNCOMPRESSED,
1891 );
1892 }
1893
1894 #[test]
1895 fn test_boolean_compressed_roundtrip() {
1896 let my_bool_values: Vec<_> = (0..2049).map(|idx| idx % 2 == 0).collect();
1897 test_roundtrip::<Vec<u8>, Bytes, BoolType, _>(
1898 Vec::with_capacity(1024),
1899 vec![my_bool_values],
1900 |r| r.get_bool(0).unwrap(),
1901 Compression::SNAPPY,
1902 );
1903 }
1904
1905 #[test]
1906 fn test_column_offset_index_file() {
1907 let file = tempfile::tempfile().unwrap();
1908 let file_metadata = test_file_roundtrip(file, vec![vec![1, 2, 3, 4, 5]]);
1909 file_metadata.row_groups().iter().for_each(|row_group| {
1910 row_group.columns().iter().for_each(|column_chunk| {
1911 assert!(column_chunk.column_index_offset().is_some());
1912 assert!(column_chunk.column_index_length().is_some());
1913 assert!(column_chunk.offset_index_offset().is_some());
1914 assert!(column_chunk.offset_index_length().is_some());
1915 })
1916 });
1917 }
1918
1919 fn test_kv_metadata(initial_kv: Option<Vec<KeyValue>>, final_kv: Option<Vec<KeyValue>>) {
1920 let schema = Arc::new(
1921 types::Type::group_type_builder("schema")
1922 .with_fields(vec![Arc::new(
1923 types::Type::primitive_type_builder("col1", Type::INT32)
1924 .with_repetition(Repetition::REQUIRED)
1925 .build()
1926 .unwrap(),
1927 )])
1928 .build()
1929 .unwrap(),
1930 );
1931 let mut out = Vec::with_capacity(1024);
1932 let props = Arc::new(
1933 WriterProperties::builder()
1934 .set_key_value_metadata(initial_kv.clone())
1935 .build(),
1936 );
1937 let mut writer = SerializedFileWriter::new(&mut out, schema, props).unwrap();
1938 let mut row_group_writer = writer.next_row_group().unwrap();
1939 let column = row_group_writer.next_column().unwrap().unwrap();
1940 column.close().unwrap();
1941 row_group_writer.close().unwrap();
1942 if let Some(kvs) = &final_kv {
1943 for kv in kvs {
1944 writer.append_key_value_metadata(kv.clone())
1945 }
1946 }
1947 writer.close().unwrap();
1948
1949 let reader = SerializedFileReader::new(Bytes::from(out)).unwrap();
1950 let metadata = reader.metadata().file_metadata();
1951 let keys = metadata.key_value_metadata();
1952
1953 match (initial_kv, final_kv) {
1954 (Some(a), Some(b)) => {
1955 let keys = keys.unwrap();
1956 assert_eq!(keys.len(), a.len() + b.len());
1957 assert_eq!(&keys[..a.len()], a.as_slice());
1958 assert_eq!(&keys[a.len()..], b.as_slice());
1959 }
1960 (Some(v), None) => assert_eq!(keys.unwrap(), &v),
1961 (None, Some(v)) if !v.is_empty() => assert_eq!(keys.unwrap(), &v),
1962 _ => assert!(keys.is_none()),
1963 }
1964 }
1965
1966 #[test]
1967 fn test_append_metadata() {
1968 let kv1 = KeyValue::new("cupcakes".to_string(), "awesome".to_string());
1969 let kv2 = KeyValue::new("bingo".to_string(), "bongo".to_string());
1970
1971 test_kv_metadata(None, None);
1972 test_kv_metadata(Some(vec![kv1.clone()]), None);
1973 test_kv_metadata(None, Some(vec![kv2.clone()]));
1974 test_kv_metadata(Some(vec![kv1.clone()]), Some(vec![kv2.clone()]));
1975 test_kv_metadata(Some(vec![]), Some(vec![kv2]));
1976 test_kv_metadata(Some(vec![]), Some(vec![]));
1977 test_kv_metadata(Some(vec![kv1]), Some(vec![]));
1978 test_kv_metadata(None, Some(vec![]));
1979 }
1980
1981 #[test]
1982 fn test_backwards_compatible_statistics() {
1983 let message_type = "
1984 message test_schema {
1985 REQUIRED INT32 decimal1 (DECIMAL(8,2));
1986 REQUIRED INT32 i32 (INTEGER(32,true));
1987 REQUIRED INT32 u32 (INTEGER(32,false));
1988 }
1989 ";
1990
1991 let schema = Arc::new(parse_message_type(message_type).unwrap());
1992 let props = Default::default();
1993 let mut writer = SerializedFileWriter::new(vec![], schema, props).unwrap();
1994 let mut row_group_writer = writer.next_row_group().unwrap();
1995
1996 for _ in 0..3 {
1997 let mut writer = row_group_writer.next_column().unwrap().unwrap();
1998 writer
1999 .typed::<Int32Type>()
2000 .write_batch(&[1, 2, 3], None, None)
2001 .unwrap();
2002 writer.close().unwrap();
2003 }
2004 let metadata = row_group_writer.close().unwrap();
2005 writer.close().unwrap();
2006
2007 let s = page_stats_to_thrift(metadata.column(0).statistics()).unwrap();
2009 assert_eq!(s.min.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2010 assert_eq!(s.max.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2011 assert_eq!(s.min_value.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2012 assert_eq!(s.max_value.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2013
2014 let s = page_stats_to_thrift(metadata.column(1).statistics()).unwrap();
2016 assert_eq!(s.min.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2017 assert_eq!(s.max.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2018 assert_eq!(s.min_value.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2019 assert_eq!(s.max_value.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2020
2021 let s = page_stats_to_thrift(metadata.column(2).statistics()).unwrap();
2023 assert_eq!(s.min.as_deref(), None);
2024 assert_eq!(s.max.as_deref(), None);
2025 assert_eq!(s.min_value.as_deref(), Some(1_i32.to_le_bytes().as_ref()));
2026 assert_eq!(s.max_value.as_deref(), Some(3_i32.to_le_bytes().as_ref()));
2027 }
2028
2029 #[test]
2030 fn test_spliced_write() {
2031 let message_type = "
2032 message test_schema {
2033 REQUIRED INT32 i32 (INTEGER(32,true));
2034 REQUIRED INT32 u32 (INTEGER(32,false));
2035 }
2036 ";
2037 let schema = Arc::new(parse_message_type(message_type).unwrap());
2038 let props = Arc::new(WriterProperties::builder().build());
2039
2040 let mut file = Vec::with_capacity(1024);
2041 let mut file_writer = SerializedFileWriter::new(&mut file, schema, props.clone()).unwrap();
2042
2043 let columns = file_writer.descr.columns();
2044 let mut column_state: Vec<(_, Option<ColumnCloseResult>)> = columns
2045 .iter()
2046 .map(|_| (TrackedWrite::new(Vec::with_capacity(1024)), None))
2047 .collect();
2048
2049 let mut column_state_slice = column_state.as_mut_slice();
2050 let mut column_writers = Vec::with_capacity(columns.len());
2051 for c in columns {
2052 let ((buf, out), tail) = column_state_slice.split_first_mut().unwrap();
2053 column_state_slice = tail;
2054
2055 let page_writer = Box::new(SerializedPageWriter::new(buf));
2056 let col_writer = get_column_writer(c.clone(), props.clone(), page_writer);
2057 column_writers.push(SerializedColumnWriter::new(
2058 col_writer,
2059 Some(Box::new(|on_close| {
2060 *out = Some(on_close);
2061 Ok(())
2062 })),
2063 ));
2064 }
2065
2066 let column_data = [[1, 2, 3, 4], [7, 3, 7, 3]];
2067
2068 for (writer, batch) in column_writers.iter_mut().zip(column_data) {
2070 let writer = writer.typed::<Int32Type>();
2071 writer.write_batch(&batch, None, None).unwrap();
2072 }
2073
2074 for writer in column_writers {
2076 writer.close().unwrap()
2077 }
2078
2079 let mut row_group_writer = file_writer.next_row_group().unwrap();
2081 for (write, close) in column_state {
2082 let buf = Bytes::from(write.into_inner().unwrap());
2083 row_group_writer
2084 .append_column(&buf, close.unwrap())
2085 .unwrap();
2086 }
2087 row_group_writer.close().unwrap();
2088 file_writer.close().unwrap();
2089
2090 let file = Bytes::from(file);
2092 let test_read = |reader: SerializedFileReader<Bytes>| {
2093 let row_group = reader.get_row_group(0).unwrap();
2094
2095 let mut out = Vec::with_capacity(4);
2096 let c1 = row_group.get_column_reader(0).unwrap();
2097 let mut c1 = get_typed_column_reader::<Int32Type>(c1);
2098 c1.read_records(4, None, None, &mut out).unwrap();
2099 assert_eq!(out, column_data[0]);
2100
2101 out.clear();
2102
2103 let c2 = row_group.get_column_reader(1).unwrap();
2104 let mut c2 = get_typed_column_reader::<Int32Type>(c2);
2105 c2.read_records(4, None, None, &mut out).unwrap();
2106 assert_eq!(out, column_data[1]);
2107 };
2108
2109 let reader = SerializedFileReader::new(file.clone()).unwrap();
2110 test_read(reader);
2111
2112 let options = ReadOptionsBuilder::new().with_page_index().build();
2113 let reader = SerializedFileReader::new_with_options(file, options).unwrap();
2114 test_read(reader);
2115 }
2116
2117 #[test]
2118 fn test_disabled_statistics() {
2119 let message_type = "
2120 message test_schema {
2121 REQUIRED INT32 a;
2122 REQUIRED INT32 b;
2123 }
2124 ";
2125 let schema = Arc::new(parse_message_type(message_type).unwrap());
2126 let props = WriterProperties::builder()
2127 .set_statistics_enabled(EnabledStatistics::None)
2128 .set_column_statistics_enabled("a".into(), EnabledStatistics::Page)
2129 .set_offset_index_disabled(true) .build();
2131 let mut file = Vec::with_capacity(1024);
2132 let mut file_writer =
2133 SerializedFileWriter::new(&mut file, schema, Arc::new(props)).unwrap();
2134
2135 let mut row_group_writer = file_writer.next_row_group().unwrap();
2136 let mut a_writer = row_group_writer.next_column().unwrap().unwrap();
2137 let col_writer = a_writer.typed::<Int32Type>();
2138 col_writer.write_batch(&[1, 2, 3], None, None).unwrap();
2139 a_writer.close().unwrap();
2140
2141 let mut b_writer = row_group_writer.next_column().unwrap().unwrap();
2142 let col_writer = b_writer.typed::<Int32Type>();
2143 col_writer.write_batch(&[4, 5, 6], None, None).unwrap();
2144 b_writer.close().unwrap();
2145 row_group_writer.close().unwrap();
2146
2147 let metadata = file_writer.finish().unwrap();
2148 assert_eq!(metadata.num_row_groups(), 1);
2149 let row_group = metadata.row_group(0);
2150 assert_eq!(row_group.num_columns(), 2);
2151 assert!(row_group.column(0).offset_index_offset().is_some());
2153 assert!(row_group.column(0).column_index_offset().is_some());
2154 assert!(row_group.column(1).offset_index_offset().is_some());
2156 assert!(row_group.column(1).column_index_offset().is_none());
2157
2158 let err = file_writer.next_row_group().err().unwrap().to_string();
2159 assert_eq!(err, "Parquet error: SerializedFileWriter already finished");
2160
2161 drop(file_writer);
2162
2163 let options = ReadOptionsBuilder::new().with_page_index().build();
2164 let reader = SerializedFileReader::new_with_options(Bytes::from(file), options).unwrap();
2165
2166 let offset_index = reader.metadata().offset_index().unwrap();
2167 assert_eq!(offset_index.len(), 1); assert_eq!(offset_index[0].len(), 2); let column_index = reader.metadata().column_index().unwrap();
2171 assert_eq!(column_index.len(), 1); assert_eq!(column_index[0].len(), 2); let a_idx = &column_index[0][0];
2175 assert!(matches!(a_idx, ColumnIndexMetaData::INT32(_)), "{a_idx:?}");
2176 let b_idx = &column_index[0][1];
2177 assert!(matches!(b_idx, ColumnIndexMetaData::NONE), "{b_idx:?}");
2178 }
2179
2180 #[test]
2181 fn test_byte_array_size_statistics() {
2182 let message_type = "
2183 message test_schema {
2184 OPTIONAL BYTE_ARRAY a (UTF8);
2185 }
2186 ";
2187 let schema = Arc::new(parse_message_type(message_type).unwrap());
2188 let data = ByteArrayType::gen_vec(32, 7);
2189 let def_levels = [1, 1, 1, 1, 0, 1, 0, 1, 0, 1];
2190 let unenc_size: i64 = data.iter().map(|x| x.len() as i64).sum();
2191 let file: File = tempfile::tempfile().unwrap();
2192 let props = Arc::new(
2193 WriterProperties::builder()
2194 .set_statistics_enabled(EnabledStatistics::Page)
2195 .build(),
2196 );
2197
2198 let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
2199 let mut row_group_writer = writer.next_row_group().unwrap();
2200
2201 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
2202 col_writer
2203 .typed::<ByteArrayType>()
2204 .write_batch(&data, Some(&def_levels), None)
2205 .unwrap();
2206 col_writer.close().unwrap();
2207 row_group_writer.close().unwrap();
2208 let file_metadata = writer.close().unwrap();
2209
2210 assert_eq!(file_metadata.num_row_groups(), 1);
2211 assert_eq!(file_metadata.row_group(0).num_columns(), 1);
2212
2213 let check_def_hist = |def_hist: &[i64]| {
2214 assert_eq!(def_hist.len(), 2);
2215 assert_eq!(def_hist[0], 3);
2216 assert_eq!(def_hist[1], 7);
2217 };
2218
2219 let meta_data = file_metadata.row_group(0).column(0);
2220
2221 assert!(meta_data.repetition_level_histogram().is_none());
2222 assert!(meta_data.definition_level_histogram().is_some());
2223 assert!(meta_data.unencoded_byte_array_data_bytes().is_some());
2224 assert_eq!(
2225 unenc_size,
2226 meta_data.unencoded_byte_array_data_bytes().unwrap()
2227 );
2228 check_def_hist(meta_data.definition_level_histogram().unwrap().values());
2229
2230 let options = ReadOptionsBuilder::new().with_page_index().build();
2232 let reader = SerializedFileReader::new_with_options(file, options).unwrap();
2233
2234 let rfile_metadata = reader.metadata().file_metadata();
2235 assert_eq!(
2236 rfile_metadata.num_rows(),
2237 file_metadata.file_metadata().num_rows()
2238 );
2239 assert_eq!(reader.num_row_groups(), 1);
2240 let rowgroup = reader.get_row_group(0).unwrap();
2241 assert_eq!(rowgroup.num_columns(), 1);
2242 let column = rowgroup.metadata().column(0);
2243 assert!(column.definition_level_histogram().is_some());
2244 assert!(column.repetition_level_histogram().is_none());
2245 assert!(column.unencoded_byte_array_data_bytes().is_some());
2246 check_def_hist(column.definition_level_histogram().unwrap().values());
2247 assert_eq!(
2248 unenc_size,
2249 column.unencoded_byte_array_data_bytes().unwrap()
2250 );
2251
2252 assert!(reader.metadata().column_index().is_some());
2254 let column_index = reader.metadata().column_index().unwrap();
2255 assert_eq!(column_index.len(), 1);
2256 assert_eq!(column_index[0].len(), 1);
2257 let col_idx = if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][0] {
2258 assert_eq!(index.num_pages(), 1);
2259 index
2260 } else {
2261 unreachable!()
2262 };
2263
2264 assert!(col_idx.repetition_level_histogram(0).is_none());
2265 assert!(col_idx.definition_level_histogram(0).is_some());
2266 check_def_hist(col_idx.definition_level_histogram(0).unwrap());
2267
2268 assert!(reader.metadata().offset_index().is_some());
2269 let offset_index = reader.metadata().offset_index().unwrap();
2270 assert_eq!(offset_index.len(), 1);
2271 assert_eq!(offset_index[0].len(), 1);
2272 assert!(offset_index[0][0].unencoded_byte_array_data_bytes.is_some());
2273 let page_sizes = offset_index[0][0]
2274 .unencoded_byte_array_data_bytes
2275 .as_ref()
2276 .unwrap();
2277 assert_eq!(page_sizes.len(), 1);
2278 assert_eq!(page_sizes[0], unenc_size);
2279 }
2280
2281 #[test]
2282 fn test_too_many_rowgroups() {
2283 let message_type = "
2284 message test_schema {
2285 REQUIRED BYTE_ARRAY a (UTF8);
2286 }
2287 ";
2288 let schema = Arc::new(parse_message_type(message_type).unwrap());
2289 let file: File = tempfile::tempfile().unwrap();
2290 let props = Arc::new(
2291 WriterProperties::builder()
2292 .set_statistics_enabled(EnabledStatistics::None)
2293 .set_max_row_group_row_count(Some(1))
2294 .build(),
2295 );
2296 let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
2297
2298 for i in 0..0x8001 {
2300 match writer.next_row_group() {
2301 Ok(mut row_group_writer) => {
2302 assert_ne!(i, 0x8000);
2303 let col_writer = row_group_writer.next_column().unwrap().unwrap();
2304 col_writer.close().unwrap();
2305 row_group_writer.close().unwrap();
2306 }
2307 Err(e) => {
2308 assert_eq!(i, 0x8000);
2309 assert_eq!(
2310 e.to_string(),
2311 "Parquet error: Parquet does not support more than 32767 row groups per file (currently: 32768)"
2312 );
2313 }
2314 }
2315 }
2316 writer.close().unwrap();
2317 }
2318
2319 #[test]
2320 fn test_size_statistics_with_repetition_and_nulls() {
2321 let message_type = "
2322 message test_schema {
2323 OPTIONAL group i32_list (LIST) {
2324 REPEATED group list {
2325 OPTIONAL INT32 element;
2326 }
2327 }
2328 }
2329 ";
2330 let schema = Arc::new(parse_message_type(message_type).unwrap());
2337 let data = [1, 2, 4, 7, 8, 9, 10];
2338 let def_levels = [3, 3, 0, 3, 2, 1, 3, 3, 3, 3];
2339 let rep_levels = [0, 1, 0, 0, 1, 0, 0, 1, 1, 1];
2340 let file = tempfile::tempfile().unwrap();
2341 let props = Arc::new(
2342 WriterProperties::builder()
2343 .set_statistics_enabled(EnabledStatistics::Page)
2344 .build(),
2345 );
2346 let mut writer = SerializedFileWriter::new(&file, schema, props).unwrap();
2347 let mut row_group_writer = writer.next_row_group().unwrap();
2348
2349 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
2350 col_writer
2351 .typed::<Int32Type>()
2352 .write_batch(&data, Some(&def_levels), Some(&rep_levels))
2353 .unwrap();
2354 col_writer.close().unwrap();
2355 row_group_writer.close().unwrap();
2356 let file_metadata = writer.close().unwrap();
2357
2358 assert_eq!(file_metadata.num_row_groups(), 1);
2359 assert_eq!(file_metadata.row_group(0).num_columns(), 1);
2360
2361 let check_def_hist = |def_hist: &[i64]| {
2362 assert_eq!(def_hist.len(), 4);
2363 assert_eq!(def_hist[0], 1);
2364 assert_eq!(def_hist[1], 1);
2365 assert_eq!(def_hist[2], 1);
2366 assert_eq!(def_hist[3], 7);
2367 };
2368
2369 let check_rep_hist = |rep_hist: &[i64]| {
2370 assert_eq!(rep_hist.len(), 2);
2371 assert_eq!(rep_hist[0], 5);
2372 assert_eq!(rep_hist[1], 5);
2373 };
2374
2375 let meta_data = file_metadata.row_group(0).column(0);
2378 assert!(meta_data.repetition_level_histogram().is_some());
2379 assert!(meta_data.definition_level_histogram().is_some());
2380 assert!(meta_data.unencoded_byte_array_data_bytes().is_none());
2381 check_def_hist(meta_data.definition_level_histogram().unwrap().values());
2382 check_rep_hist(meta_data.repetition_level_histogram().unwrap().values());
2383
2384 let options = ReadOptionsBuilder::new().with_page_index().build();
2386 let reader = SerializedFileReader::new_with_options(file, options).unwrap();
2387
2388 let rfile_metadata = reader.metadata().file_metadata();
2389 assert_eq!(
2390 rfile_metadata.num_rows(),
2391 file_metadata.file_metadata().num_rows()
2392 );
2393 assert_eq!(reader.num_row_groups(), 1);
2394 let rowgroup = reader.get_row_group(0).unwrap();
2395 assert_eq!(rowgroup.num_columns(), 1);
2396 let column = rowgroup.metadata().column(0);
2397 assert!(column.definition_level_histogram().is_some());
2398 assert!(column.repetition_level_histogram().is_some());
2399 assert!(column.unencoded_byte_array_data_bytes().is_none());
2400 check_def_hist(column.definition_level_histogram().unwrap().values());
2401 check_rep_hist(column.repetition_level_histogram().unwrap().values());
2402
2403 assert!(reader.metadata().column_index().is_some());
2405 let column_index = reader.metadata().column_index().unwrap();
2406 assert_eq!(column_index.len(), 1);
2407 assert_eq!(column_index[0].len(), 1);
2408 let col_idx = if let ColumnIndexMetaData::INT32(index) = &column_index[0][0] {
2409 assert_eq!(index.num_pages(), 1);
2410 index
2411 } else {
2412 unreachable!()
2413 };
2414
2415 check_def_hist(col_idx.definition_level_histogram(0).unwrap());
2416 check_rep_hist(col_idx.repetition_level_histogram(0).unwrap());
2417
2418 assert!(reader.metadata().offset_index().is_some());
2419 let offset_index = reader.metadata().offset_index().unwrap();
2420 assert_eq!(offset_index.len(), 1);
2421 assert_eq!(offset_index[0].len(), 1);
2422 assert!(offset_index[0][0].unencoded_byte_array_data_bytes.is_none());
2423 }
2424
2425 #[test]
2426 #[cfg(feature = "arrow")]
2427 fn test_byte_stream_split_extended_roundtrip() {
2428 let path = format!(
2429 "{}/byte_stream_split_extended.gzip.parquet",
2430 arrow::util::test_util::parquet_test_data(),
2431 );
2432 let file = File::open(path).unwrap();
2433
2434 let parquet_reader = ParquetRecordBatchReaderBuilder::try_new(file)
2436 .expect("parquet open")
2437 .build()
2438 .expect("parquet open");
2439
2440 let file = tempfile::tempfile().unwrap();
2441 let props = WriterProperties::builder()
2442 .set_dictionary_enabled(false)
2443 .set_column_encoding(
2444 ColumnPath::from("float16_byte_stream_split"),
2445 Encoding::BYTE_STREAM_SPLIT,
2446 )
2447 .set_column_encoding(
2448 ColumnPath::from("float_byte_stream_split"),
2449 Encoding::BYTE_STREAM_SPLIT,
2450 )
2451 .set_column_encoding(
2452 ColumnPath::from("double_byte_stream_split"),
2453 Encoding::BYTE_STREAM_SPLIT,
2454 )
2455 .set_column_encoding(
2456 ColumnPath::from("int32_byte_stream_split"),
2457 Encoding::BYTE_STREAM_SPLIT,
2458 )
2459 .set_column_encoding(
2460 ColumnPath::from("int64_byte_stream_split"),
2461 Encoding::BYTE_STREAM_SPLIT,
2462 )
2463 .set_column_encoding(
2464 ColumnPath::from("flba5_byte_stream_split"),
2465 Encoding::BYTE_STREAM_SPLIT,
2466 )
2467 .set_column_encoding(
2468 ColumnPath::from("decimal_byte_stream_split"),
2469 Encoding::BYTE_STREAM_SPLIT,
2470 )
2471 .build();
2472
2473 let mut parquet_writer = ArrowWriter::try_new(
2474 file.try_clone().expect("cannot open file"),
2475 parquet_reader.schema(),
2476 Some(props),
2477 )
2478 .expect("create arrow writer");
2479
2480 for maybe_batch in parquet_reader {
2481 let batch = maybe_batch.expect("reading batch");
2482 parquet_writer.write(&batch).expect("writing data");
2483 }
2484
2485 parquet_writer.close().expect("finalizing file");
2486
2487 let reader = SerializedFileReader::new(file).expect("Failed to create reader");
2488 let filemeta = reader.metadata();
2489
2490 let check_encoding = |x: usize, filemeta: &ParquetMetaData| {
2492 assert!(
2493 filemeta
2494 .row_group(0)
2495 .column(x)
2496 .encodings()
2497 .collect::<Vec<_>>()
2498 .contains(&Encoding::BYTE_STREAM_SPLIT)
2499 );
2500 };
2501
2502 check_encoding(1, filemeta);
2503 check_encoding(3, filemeta);
2504 check_encoding(5, filemeta);
2505 check_encoding(7, filemeta);
2506 check_encoding(9, filemeta);
2507 check_encoding(11, filemeta);
2508 check_encoding(13, filemeta);
2509
2510 let mut iter = reader
2512 .get_row_iter(None)
2513 .expect("Failed to create row iterator");
2514
2515 let mut start = 0;
2516 let end = reader.metadata().file_metadata().num_rows();
2517
2518 let check_row = |row: Result<Row, ParquetError>| {
2519 assert!(row.is_ok());
2520 let r = row.unwrap();
2521 assert_eq!(r.get_float16(0).unwrap(), r.get_float16(1).unwrap());
2522 assert_eq!(r.get_float(2).unwrap(), r.get_float(3).unwrap());
2523 assert_eq!(r.get_double(4).unwrap(), r.get_double(5).unwrap());
2524 assert_eq!(r.get_int(6).unwrap(), r.get_int(7).unwrap());
2525 assert_eq!(r.get_long(8).unwrap(), r.get_long(9).unwrap());
2526 assert_eq!(r.get_bytes(10).unwrap(), r.get_bytes(11).unwrap());
2527 assert_eq!(r.get_decimal(12).unwrap(), r.get_decimal(13).unwrap());
2528 };
2529
2530 while start < end {
2531 match iter.next() {
2532 Some(row) => check_row(row),
2533 None => break,
2534 };
2535 start += 1;
2536 }
2537 }
2538
2539 #[test]
2540 fn test_rewrite_no_page_indexes() {
2541 let file = get_test_file("alltypes_tiny_pages.parquet");
2542 let metadata = ParquetMetaDataReader::new()
2543 .with_page_index_policy(PageIndexPolicy::Optional)
2544 .parse_and_finish(&file)
2545 .unwrap();
2546
2547 let props = Arc::new(WriterProperties::builder().build());
2548 let schema = metadata.file_metadata().schema_descr().root_schema_ptr();
2549 let output = Vec::<u8>::new();
2550 let mut writer = SerializedFileWriter::new(output, schema, props).unwrap();
2551
2552 for rg in metadata.row_groups() {
2553 let mut rg_out = writer.next_row_group().unwrap();
2554 for column in rg.columns() {
2555 let result = ColumnCloseResult {
2556 bytes_written: column.compressed_size() as _,
2557 rows_written: rg.num_rows() as _,
2558 metadata: column.clone(),
2559 bloom_filter: None,
2560 column_index: None,
2561 offset_index: None,
2562 };
2563 rg_out.append_column(&file, result).unwrap();
2564 }
2565 rg_out.close().unwrap();
2566 }
2567 writer.close().unwrap();
2568 }
2569
2570 #[test]
2571 fn test_rewrite_missing_column_index() {
2572 let file = get_test_file("alltypes_tiny_pages.parquet");
2574 let metadata = ParquetMetaDataReader::new()
2575 .with_page_index_policy(PageIndexPolicy::Optional)
2576 .parse_and_finish(&file)
2577 .unwrap();
2578
2579 let props = Arc::new(WriterProperties::builder().build());
2580 let schema = metadata.file_metadata().schema_descr().root_schema_ptr();
2581 let output = Vec::<u8>::new();
2582 let mut writer = SerializedFileWriter::new(output, schema, props).unwrap();
2583
2584 let column_indexes = metadata.column_index();
2585 let offset_indexes = metadata.offset_index();
2586
2587 for (rg_idx, rg) in metadata.row_groups().iter().enumerate() {
2588 let rg_column_indexes = column_indexes.and_then(|ci| ci.get(rg_idx));
2589 let rg_offset_indexes = offset_indexes.and_then(|oi| oi.get(rg_idx));
2590 let mut rg_out = writer.next_row_group().unwrap();
2591 for (col_idx, column) in rg.columns().iter().enumerate() {
2592 let column_index = rg_column_indexes.and_then(|row| row.get(col_idx)).cloned();
2593 let offset_index = rg_offset_indexes.and_then(|row| row.get(col_idx)).cloned();
2594 let result = ColumnCloseResult {
2595 bytes_written: column.compressed_size() as _,
2596 rows_written: rg.num_rows() as _,
2597 metadata: column.clone(),
2598 bloom_filter: None,
2599 column_index,
2600 offset_index,
2601 };
2602 rg_out.append_column(&file, result).unwrap();
2603 }
2604 rg_out.close().unwrap();
2605 }
2606 writer.close().unwrap();
2607 }
2608}