1use crate::basic::{PageType, Type};
22use crate::bloom_filter::Sbbf;
23use crate::column::page::{Page, PageMetadata, PageReader};
24use crate::compression::{Codec, create_codec};
25#[cfg(feature = "encryption")]
26use crate::encryption::decrypt::{CryptoContext, read_and_decrypt};
27use crate::errors::{ParquetError, Result};
28use crate::file::metadata::page_index::RowGroupPageIndex;
29use crate::file::metadata::thrift::PageHeader;
30use crate::file::page_index::offset_index::PageLocation;
31use crate::file::statistics;
32use crate::file::{
33 metadata::*,
34 properties::{ReaderProperties, ReaderPropertiesPtr},
35 reader::*,
36};
37#[cfg(feature = "encryption")]
38use crate::parquet_thrift::ThriftSliceInputProtocol;
39use crate::parquet_thrift::{ReadThrift, ThriftReadInputProtocol};
40use crate::record::Row;
41use crate::record::reader::RowIter;
42use crate::schema::types::{SchemaDescPtr, Type as SchemaType};
43use bytes::Bytes;
44use std::collections::VecDeque;
45use std::{fs::File, io::Read, path::Path, sync::Arc};
46
47impl TryFrom<File> for SerializedFileReader<File> {
48 type Error = ParquetError;
49
50 fn try_from(file: File) -> Result<Self> {
51 Self::new(file)
52 }
53}
54
55impl TryFrom<&Path> for SerializedFileReader<File> {
56 type Error = ParquetError;
57
58 fn try_from(path: &Path) -> Result<Self> {
59 let file = File::open(path)?;
60 Self::try_from(file)
61 }
62}
63
64impl TryFrom<String> for SerializedFileReader<File> {
65 type Error = ParquetError;
66
67 fn try_from(path: String) -> Result<Self> {
68 Self::try_from(Path::new(&path))
69 }
70}
71
72impl TryFrom<&str> for SerializedFileReader<File> {
73 type Error = ParquetError;
74
75 fn try_from(path: &str) -> Result<Self> {
76 Self::try_from(Path::new(&path))
77 }
78}
79
80impl IntoIterator for SerializedFileReader<File> {
83 type Item = Result<Row>;
84 type IntoIter = RowIter<'static>;
85
86 fn into_iter(self) -> Self::IntoIter {
87 RowIter::from_file_into(Box::new(self))
88 }
89}
90
91pub struct SerializedFileReader<R: ChunkReader> {
96 chunk_reader: Arc<R>,
97 metadata: Arc<ParquetMetaData>,
98 props: ReaderPropertiesPtr,
99}
100
101pub type ReadGroupPredicate = Box<dyn FnMut(&RowGroupMetaData, usize) -> bool>;
105
106#[derive(Default)]
110pub struct ReadOptionsBuilder {
111 predicates: Vec<ReadGroupPredicate>,
112 enable_page_index: bool,
113 props: Option<ReaderProperties>,
114 metadata_options: ParquetMetaDataOptions,
115}
116
117impl ReadOptionsBuilder {
118 pub fn new() -> Self {
120 Self::default()
121 }
122
123 pub fn with_predicate(mut self, predicate: ReadGroupPredicate) -> Self {
126 self.predicates.push(predicate);
127 self
128 }
129
130 pub fn with_range(mut self, start: i64, end: i64) -> Self {
137 assert!(start < end);
138 let predicate = move |rg: &RowGroupMetaData, _: usize| {
139 let mid = get_midpoint_offset(rg);
140 mid >= start && mid < end
141 };
142 self.predicates.push(Box::new(predicate));
143 self
144 }
145
146 pub fn with_page_index(mut self) -> Self {
151 self.enable_page_index = true;
152 self
153 }
154
155 pub fn with_reader_properties(mut self, properties: ReaderProperties) -> Self {
157 self.props = Some(properties);
158 self
159 }
160
161 pub fn with_parquet_schema(mut self, schema: SchemaDescPtr) -> Self {
164 self.metadata_options.set_schema(schema);
165 self
166 }
167
168 pub fn with_encoding_stats_as_mask(mut self, val: bool) -> Self {
177 self.metadata_options.set_encoding_stats_as_mask(val);
178 self
179 }
180
181 pub fn with_encoding_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
186 self.metadata_options.set_encoding_stats_policy(policy);
187 self
188 }
189
190 pub fn with_column_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
195 self.metadata_options.set_column_stats_policy(policy);
196 self
197 }
198
199 pub fn with_size_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
204 self.metadata_options.set_size_stats_policy(policy);
205 self
206 }
207
208 pub fn build(self) -> ReadOptions {
210 let props = self
211 .props
212 .unwrap_or_else(|| ReaderProperties::builder().build());
213 ReadOptions {
214 predicates: self.predicates,
215 enable_page_index: self.enable_page_index,
216 props,
217 metadata_options: self.metadata_options,
218 }
219 }
220}
221
222pub struct ReadOptions {
227 predicates: Vec<ReadGroupPredicate>,
228 enable_page_index: bool,
229 props: ReaderProperties,
230 metadata_options: ParquetMetaDataOptions,
231}
232
233impl<R: 'static + ChunkReader> SerializedFileReader<R> {
234 pub fn new(chunk_reader: R) -> Result<Self> {
237 let metadata = ParquetMetaDataReader::new().parse_and_finish(&chunk_reader)?;
238 let props = Arc::new(ReaderProperties::builder().build());
239 Ok(Self {
240 chunk_reader: Arc::new(chunk_reader),
241 metadata: Arc::new(metadata),
242 props,
243 })
244 }
245
246 pub fn new_with_options(chunk_reader: R, options: ReadOptions) -> Result<Self> {
249 let mut metadata_builder = ParquetMetaDataReader::new()
250 .with_metadata_options(Some(options.metadata_options.clone()))
251 .parse_and_finish(&chunk_reader)?
252 .into_builder();
253 let mut predicates = options.predicates;
254
255 for (i, rg_meta) in metadata_builder.take_row_groups().into_iter().enumerate() {
257 let mut keep = true;
258 for predicate in &mut predicates {
259 if !predicate(&rg_meta, i) {
260 keep = false;
261 break;
262 }
263 }
264 if keep {
265 metadata_builder = metadata_builder.add_row_group(rg_meta);
266 }
267 }
268
269 let mut metadata = metadata_builder.build();
270
271 if options.enable_page_index {
273 let mut reader = ParquetMetaDataReader::new_with_metadata(metadata)
274 .with_page_index_policy(PageIndexPolicy::Required);
275 reader.read_page_indexes(&chunk_reader)?;
276 metadata = reader.finish()?;
277 }
278
279 Ok(Self {
280 chunk_reader: Arc::new(chunk_reader),
281 metadata: Arc::new(metadata),
282 props: Arc::new(options.props),
283 })
284 }
285}
286
287fn get_midpoint_offset(meta: &RowGroupMetaData) -> i64 {
289 let col = meta.column(0);
290 let mut offset = col.data_page_offset();
291 if let Some(dic_offset) = col.dictionary_page_offset()
292 && offset > dic_offset
293 {
294 offset = dic_offset
295 }
296 offset + meta.compressed_size() / 2
297}
298
299impl<R: 'static + ChunkReader> FileReader for SerializedFileReader<R> {
300 fn metadata(&self) -> &ParquetMetaData {
301 &self.metadata
302 }
303
304 fn num_row_groups(&self) -> usize {
305 self.metadata.num_row_groups()
306 }
307
308 fn get_row_group(&self, i: usize) -> Result<Box<dyn RowGroupReader + '_>> {
309 let row_group_metadata = self.metadata.row_group(i);
310 let props = Arc::clone(&self.props);
312 let f = Arc::clone(&self.chunk_reader);
313 let page_index = self.metadata.page_index_for_row_group(i);
314 Ok(Box::new(SerializedRowGroupReader::new(
315 f,
316 row_group_metadata,
317 page_index,
318 props,
319 )?))
320 }
321
322 fn get_row_iter(&self, projection: Option<SchemaType>) -> Result<RowIter<'_>> {
323 RowIter::from_file(projection, self)
324 }
325}
326
327pub struct SerializedRowGroupReader<'a, R: ChunkReader> {
329 chunk_reader: Arc<R>,
330 metadata: &'a RowGroupMetaData,
331 page_index: RowGroupPageIndex,
332 props: ReaderPropertiesPtr,
333 bloom_filters: Vec<Option<Sbbf>>,
334}
335
336impl<'a, R: ChunkReader> SerializedRowGroupReader<'a, R> {
337 pub fn new(
339 chunk_reader: Arc<R>,
340 metadata: &'a RowGroupMetaData,
341 page_index: RowGroupPageIndex,
342 props: ReaderPropertiesPtr,
343 ) -> Result<Self> {
344 let bloom_filters = if props.read_bloom_filter() {
345 metadata
346 .columns()
347 .iter()
348 .map(|col| Sbbf::read_from_column_chunk(col, &*chunk_reader))
349 .collect::<Result<Vec<_>>>()?
350 } else {
351 std::iter::repeat_n(None, metadata.columns().len()).collect()
352 };
353 Ok(Self {
354 chunk_reader,
355 metadata,
356 page_index,
357 props,
358 bloom_filters,
359 })
360 }
361}
362
363impl<R: 'static + ChunkReader> RowGroupReader for SerializedRowGroupReader<'_, R> {
364 fn metadata(&self) -> &RowGroupMetaData {
365 self.metadata
366 }
367
368 fn num_columns(&self) -> usize {
369 self.metadata.num_columns()
370 }
371
372 fn get_column_page_reader(&self, i: usize) -> Result<Box<dyn PageReader>> {
374 let col = self.metadata.column(i);
375
376 let page_locations = self.page_index.page_locations(i).cloned();
378
379 let props = Arc::clone(&self.props);
380 Ok(Box::new(SerializedPageReader::new_with_properties(
381 Arc::clone(&self.chunk_reader),
382 col,
383 usize::try_from(self.metadata.num_rows())?,
384 page_locations,
385 props,
386 )?))
387 }
388
389 fn get_column_bloom_filter(&self, i: usize) -> Option<&Sbbf> {
391 self.bloom_filters[i].as_ref()
392 }
393
394 fn get_row_iter(&self, projection: Option<SchemaType>) -> Result<RowIter<'_>> {
395 RowIter::from_row_group(projection, self)
396 }
397}
398
399pub(crate) fn decode_page(
401 page_header: PageHeader,
402 buffer: Bytes,
403 physical_type: Type,
404 decompressor: Option<&mut Box<dyn Codec>>,
405) -> Result<Page> {
406 #[cfg(feature = "crc")]
408 if let Some(expected_crc) = page_header.crc {
409 let crc = crc32fast::hash(&buffer);
410 if crc != expected_crc as u32 {
411 return Err(general_err!("Page CRC checksum mismatch"));
412 }
413 }
414
415 let (offset, can_decompress): (usize, bool) = match page_header.data_page_header_v2 {
422 Some(ref header_v2) => {
423 if header_v2.definition_levels_byte_length < 0
424 || header_v2.repetition_levels_byte_length < 0
425 || header_v2.definition_levels_byte_length + header_v2.repetition_levels_byte_length
426 > page_header.uncompressed_page_size
427 {
428 return Err(general_err!(
429 "DataPage v2 header contains implausible values \
430 for definition_levels_byte_length ({}) \
431 and repetition_levels_byte_length ({}) \
432 given DataPage header provides uncompressed_page_size ({})",
433 header_v2.definition_levels_byte_length,
434 header_v2.repetition_levels_byte_length,
435 page_header.uncompressed_page_size
436 ));
437 }
438 (
439 usize::try_from(
440 header_v2.definition_levels_byte_length
441 + header_v2.repetition_levels_byte_length,
442 )?,
443 header_v2.is_compressed.unwrap_or(true),
445 )
446 }
447 None => (0, true),
448 };
449
450 let buffer = match decompressor {
451 Some(decompressor) if can_decompress => {
452 let uncompressed_page_size = usize::try_from(page_header.uncompressed_page_size)?;
453 if offset > buffer.len() || offset > uncompressed_page_size {
454 return Err(general_err!("Invalid page header"));
455 }
456 let decompressed_size = uncompressed_page_size - offset;
457 let mut decompressed = Vec::with_capacity(uncompressed_page_size);
458 decompressed.extend_from_slice(&buffer[..offset]);
459 if decompressed_size > 0 {
462 let compressed = &buffer[offset..];
463 decompressor.decompress(compressed, &mut decompressed, Some(decompressed_size))?;
464 }
465
466 if decompressed.len() != uncompressed_page_size {
467 return Err(general_err!(
468 "Actual decompressed size doesn't match the expected one ({} vs {})",
469 decompressed.len(),
470 uncompressed_page_size
471 ));
472 }
473
474 Bytes::from(decompressed)
475 }
476 _ => buffer,
477 };
478
479 let result = match page_header.r#type {
480 PageType::DICTIONARY_PAGE => {
481 let dict_header = page_header.dictionary_page_header.as_ref().ok_or_else(|| {
482 ParquetError::General("Missing dictionary page header".to_string())
483 })?;
484 let is_sorted = dict_header.is_sorted.unwrap_or(false);
485 Page::DictionaryPage {
486 buf: buffer,
487 num_values: dict_header.num_values.try_into()?,
488 encoding: dict_header.encoding,
489 is_sorted,
490 }
491 }
492 PageType::DATA_PAGE => {
493 let header = page_header
494 .data_page_header
495 .ok_or_else(|| ParquetError::General("Missing V1 data page header".to_string()))?;
496 Page::DataPage {
497 buf: buffer,
498 num_values: header.num_values.try_into()?,
499 encoding: header.encoding,
500 def_level_encoding: header.definition_level_encoding,
501 rep_level_encoding: header.repetition_level_encoding,
502 statistics: statistics::from_thrift_page_stats(physical_type, header.statistics)?,
503 }
504 }
505 PageType::DATA_PAGE_V2 => {
506 let header = page_header
507 .data_page_header_v2
508 .ok_or_else(|| ParquetError::General("Missing V2 data page header".to_string()))?;
509 let is_compressed = header.is_compressed.unwrap_or(true);
510 Page::DataPageV2 {
511 buf: buffer,
512 num_values: header.num_values.try_into()?,
513 encoding: header.encoding,
514 num_nulls: header.num_nulls.try_into()?,
515 num_rows: header.num_rows.try_into()?,
516 def_levels_byte_len: header.definition_levels_byte_length.try_into()?,
517 rep_levels_byte_len: header.repetition_levels_byte_length.try_into()?,
518 is_compressed,
519 statistics: statistics::from_thrift_page_stats(physical_type, header.statistics)?,
520 }
521 }
522 PageType::INDEX_PAGE => {
523 return Err(general_err!(
525 "Page type {:?} is not supported",
526 page_header.r#type
527 ));
528 }
529 };
530
531 Ok(result)
532}
533
534enum SerializedPageReaderState {
535 Values {
536 offset: u64,
539
540 remaining_bytes: u64,
543
544 next_page_header: Option<Box<PageHeader>>,
546
547 page_index: usize,
549
550 require_dictionary: bool,
552 },
553 Pages {
554 page_locations: VecDeque<PageLocation>,
556 dictionary_page: Option<PageLocation>,
558 total_rows: usize,
560 page_index: usize,
562 },
563}
564
565#[derive(Default)]
566struct SerializedPageReaderContext {
567 read_stats: bool,
569 #[cfg(feature = "encryption")]
571 crypto_context: Option<Arc<CryptoContext>>,
572}
573
574pub struct SerializedPageReader<R: ChunkReader> {
576 reader: Arc<R>,
578
579 decompressor: Option<Box<dyn Codec>>,
581
582 physical_type: Type,
584
585 state: SerializedPageReaderState,
586
587 context: SerializedPageReaderContext,
588}
589
590impl<R: ChunkReader> SerializedPageReader<R> {
591 pub fn new(
593 reader: Arc<R>,
594 column_chunk_metadata: &ColumnChunkMetaData,
595 total_rows: usize,
596 page_locations: Option<Vec<PageLocation>>,
597 ) -> Result<Self> {
598 let props = Arc::new(ReaderProperties::builder().build());
599 SerializedPageReader::new_with_properties(
600 reader,
601 column_chunk_metadata,
602 total_rows,
603 page_locations,
604 props,
605 )
606 }
607
608 #[cfg(all(feature = "arrow", not(feature = "encryption")))]
610 pub(crate) fn add_crypto_context(
611 self,
612 _rg_idx: usize,
613 _column_idx: usize,
614 _parquet_meta_data: &ParquetMetaData,
615 _column_chunk_metadata: &ColumnChunkMetaData,
616 ) -> Result<SerializedPageReader<R>> {
617 Ok(self)
618 }
619
620 #[cfg(feature = "encryption")]
622 pub(crate) fn add_crypto_context(
623 mut self,
624 rg_idx: usize,
625 column_idx: usize,
626 parquet_meta_data: &ParquetMetaData,
627 column_chunk_metadata: &ColumnChunkMetaData,
628 ) -> Result<SerializedPageReader<R>> {
629 let Some(file_decryptor) = parquet_meta_data.file_decryptor() else {
630 return Ok(self);
631 };
632 let Some(crypto_metadata) = column_chunk_metadata.crypto_metadata() else {
633 return Ok(self);
634 };
635 let crypto_context =
636 CryptoContext::for_column(file_decryptor, crypto_metadata, rg_idx, column_idx)?;
637 self.context.crypto_context = Some(Arc::new(crypto_context));
638 Ok(self)
639 }
640
641 pub fn new_with_properties(
643 reader: Arc<R>,
644 meta: &ColumnChunkMetaData,
645 total_rows: usize,
646 page_locations: Option<Vec<PageLocation>>,
647 props: ReaderPropertiesPtr,
648 ) -> Result<Self> {
649 let decompressor = create_codec(meta.compression(), props.codec_options())?;
650 let (start, len) = meta.byte_range();
651
652 let state = match page_locations {
653 Some(locations) => {
654 let dictionary_page = match locations.first() {
657 Some(dict_offset) if dict_offset.offset as u64 != start => Some(PageLocation {
658 offset: start as i64,
659 compressed_page_size: (dict_offset.offset as u64 - start) as i32,
660 first_row_index: 0,
661 }),
662 _ => None,
663 };
664
665 SerializedPageReaderState::Pages {
666 page_locations: locations.into(),
667 dictionary_page,
668 total_rows,
669 page_index: 0,
670 }
671 }
672 None => SerializedPageReaderState::Values {
673 offset: start,
674 remaining_bytes: len,
675 next_page_header: None,
676 page_index: 0,
677 require_dictionary: meta.dictionary_page_offset().is_some(),
678 },
679 };
680 let mut context = SerializedPageReaderContext::default();
681 if props.read_page_stats() {
682 context.read_stats = true;
683 }
684 Ok(Self {
685 reader,
686 decompressor,
687 state,
688 physical_type: meta.column_type(),
689 context,
690 })
691 }
692
693 #[cfg(test)]
699 fn peek_next_page_offset(&mut self) -> Result<Option<u64>> {
700 match &mut self.state {
701 SerializedPageReaderState::Values {
702 offset,
703 remaining_bytes,
704 next_page_header,
705 page_index,
706 require_dictionary,
707 } => {
708 loop {
709 if *remaining_bytes == 0 {
710 return Ok(None);
711 }
712 return if let Some(header) = next_page_header.as_ref() {
713 if let Ok(_page_meta) = PageMetadata::try_from(&**header) {
714 Ok(Some(*offset))
715 } else {
716 *next_page_header = None;
718 continue;
719 }
720 } else {
721 let mut read = self.reader.get_read(*offset)?;
722 let (header_len, header) = Self::read_page_header_len(
723 &self.context,
724 &mut read,
725 *page_index,
726 *require_dictionary,
727 )?;
728 *offset += header_len as u64;
729 *remaining_bytes -= header_len as u64;
730 let page_meta = if let Ok(_page_meta) = PageMetadata::try_from(&header) {
731 Ok(Some(*offset))
732 } else {
733 continue;
735 };
736 *next_page_header = Some(Box::new(header));
737 page_meta
738 };
739 }
740 }
741 SerializedPageReaderState::Pages {
742 page_locations,
743 dictionary_page,
744 ..
745 } => {
746 if let Some(page) = dictionary_page {
747 Ok(Some(page.offset as u64))
748 } else if let Some(page) = page_locations.front() {
749 Ok(Some(page.offset as u64))
750 } else {
751 Ok(None)
752 }
753 }
754 }
755 }
756
757 fn read_page_header_len<T: Read>(
758 context: &SerializedPageReaderContext,
759 input: &mut T,
760 page_index: usize,
761 dictionary_page: bool,
762 ) -> Result<(usize, PageHeader)> {
763 struct TrackedRead<R> {
765 inner: R,
766 bytes_read: usize,
767 }
768
769 impl<R: Read> Read for TrackedRead<R> {
770 fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
771 let v = self.inner.read(buf)?;
772 self.bytes_read += v;
773 Ok(v)
774 }
775 }
776
777 let mut tracked = TrackedRead {
778 inner: input,
779 bytes_read: 0,
780 };
781 let header = context.read_page_header(&mut tracked, page_index, dictionary_page)?;
782 Ok((tracked.bytes_read, header))
783 }
784
785 fn read_page_header_len_from_bytes(
786 context: &SerializedPageReaderContext,
787 buffer: &[u8],
788 page_index: usize,
789 dictionary_page: bool,
790 ) -> Result<(usize, PageHeader)> {
791 let mut input = std::io::Cursor::new(buffer);
792 let header = context.read_page_header(&mut input, page_index, dictionary_page)?;
793 let header_len = input.position() as usize;
794 Ok((header_len, header))
795 }
796}
797
798#[cfg(not(feature = "encryption"))]
799impl SerializedPageReaderContext {
800 fn read_page_header<T: Read>(
801 &self,
802 input: &mut T,
803 _page_index: usize,
804 _dictionary_page: bool,
805 ) -> Result<PageHeader> {
806 let mut prot = ThriftReadInputProtocol::new(input);
807 if self.read_stats {
808 Ok(PageHeader::read_thrift(&mut prot)?)
809 } else {
810 Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
811 }
812 }
813
814 fn decrypt_page_data<T>(
815 &self,
816 buffer: T,
817 _page_index: usize,
818 _dictionary_page: bool,
819 ) -> Result<T> {
820 Ok(buffer)
821 }
822}
823
824#[cfg(feature = "encryption")]
825impl SerializedPageReaderContext {
826 fn read_page_header<T: Read>(
827 &self,
828 input: &mut T,
829 page_index: usize,
830 dictionary_page: bool,
831 ) -> Result<PageHeader> {
832 match self.page_crypto_context(page_index, dictionary_page) {
833 None => {
834 let mut prot = ThriftReadInputProtocol::new(input);
835 if self.read_stats {
836 Ok(PageHeader::read_thrift(&mut prot)?)
837 } else {
838 use crate::file::metadata::thrift::PageHeader;
839
840 Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
841 }
842 }
843 Some(page_crypto_context) => {
844 let data_decryptor = page_crypto_context.data_decryptor();
845 let aad = page_crypto_context.create_page_header_aad()?;
846
847 let buf = read_and_decrypt(data_decryptor, input, aad.as_ref()).map_err(|_| {
848 ParquetError::General(format!(
849 "Error decrypting page header for column {}, decryption key may be wrong",
850 page_crypto_context.column_ordinal
851 ))
852 })?;
853
854 let mut prot = ThriftSliceInputProtocol::new(buf.as_slice());
855 if self.read_stats {
856 Ok(PageHeader::read_thrift(&mut prot)?)
857 } else {
858 Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
859 }
860 }
861 }
862 }
863
864 fn decrypt_page_data<T>(&self, buffer: T, page_index: usize, dictionary_page: bool) -> Result<T>
865 where
866 T: AsRef<[u8]>,
867 T: From<Vec<u8>>,
868 {
869 let page_crypto_context = self.page_crypto_context(page_index, dictionary_page);
870 if let Some(page_crypto_context) = page_crypto_context {
871 let decryptor = page_crypto_context.data_decryptor();
872 let aad = page_crypto_context.create_page_aad()?;
873 let decrypted = decryptor.decrypt(buffer.as_ref(), &aad)?;
874 Ok(T::from(decrypted))
875 } else {
876 Ok(buffer)
877 }
878 }
879
880 fn page_crypto_context(
881 &self,
882 page_index: usize,
883 dictionary_page: bool,
884 ) -> Option<Arc<CryptoContext>> {
885 self.crypto_context.as_ref().map(|c| {
886 Arc::new(if dictionary_page {
887 c.for_dictionary_page()
888 } else {
889 c.with_page_ordinal(page_index)
890 })
891 })
892 }
893}
894
895impl<R: ChunkReader> Iterator for SerializedPageReader<R> {
896 type Item = Result<Page>;
897
898 fn next(&mut self) -> Option<Self::Item> {
899 self.get_next_page().transpose()
900 }
901}
902
903fn verify_page_header_len(header_len: usize, remaining_bytes: u64) -> Result<()> {
904 if header_len as u64 > remaining_bytes {
905 return Err(eof_err!("Invalid page header"));
906 }
907 Ok(())
908}
909
910fn verify_page_size(
911 compressed_size: i32,
912 uncompressed_size: i32,
913 remaining_bytes: u64,
914) -> Result<()> {
915 if compressed_size < 0 || compressed_size as u64 > remaining_bytes || uncompressed_size < 0 {
919 return Err(eof_err!("Invalid page header"));
920 }
921 Ok(())
922}
923
924impl<R: ChunkReader> PageReader for SerializedPageReader<R> {
925 fn get_next_page(&mut self) -> Result<Option<Page>> {
926 loop {
927 let page = match &mut self.state {
928 SerializedPageReaderState::Values {
929 offset,
930 remaining_bytes: remaining,
931 next_page_header,
932 page_index,
933 require_dictionary,
934 } => {
935 if *remaining == 0 {
936 return Ok(None);
937 }
938
939 let mut read = self.reader.get_read(*offset)?;
940 let header = if let Some(header) = next_page_header.take() {
941 *header
942 } else {
943 let (header_len, header) = Self::read_page_header_len(
944 &self.context,
945 &mut read,
946 *page_index,
947 *require_dictionary,
948 )?;
949 verify_page_header_len(header_len, *remaining)?;
950 *offset += header_len as u64;
951 *remaining -= header_len as u64;
952 header
953 };
954 verify_page_size(
955 header.compressed_page_size,
956 header.uncompressed_page_size,
957 *remaining,
958 )?;
959 let data_len = header.compressed_page_size as usize;
960 let data_start = *offset;
961 *offset += data_len as u64;
962 *remaining -= data_len as u64;
963
964 if header.r#type == PageType::INDEX_PAGE {
965 continue;
966 }
967
968 let buffer = self.reader.get_bytes(data_start, data_len)?;
969
970 let buffer =
971 self.context
972 .decrypt_page_data(buffer, *page_index, *require_dictionary)?;
973
974 let page = decode_page(
975 header,
976 buffer,
977 self.physical_type,
978 self.decompressor.as_mut(),
979 )?;
980 if page.is_data_page() {
981 *page_index += 1;
982 } else if page.is_dictionary_page() {
983 *require_dictionary = false;
984 }
985 page
986 }
987 SerializedPageReaderState::Pages {
988 page_locations,
989 dictionary_page,
990 page_index,
991 ..
992 } => {
993 let (front, is_dictionary_page) = match dictionary_page.take() {
994 Some(front) => (front, true),
995 None => match page_locations.pop_front() {
996 Some(front) => (front, false),
997 None => return Ok(None),
998 },
999 };
1000
1001 let page_len = usize::try_from(front.compressed_page_size)?;
1002 let buffer = self.reader.get_bytes(front.offset as u64, page_len)?;
1003
1004 let (offset, header) = Self::read_page_header_len_from_bytes(
1005 &self.context,
1006 buffer.as_ref(),
1007 *page_index,
1008 is_dictionary_page,
1009 )?;
1010 let bytes = buffer.slice(offset..);
1011 let bytes =
1012 self.context
1013 .decrypt_page_data(bytes, *page_index, is_dictionary_page)?;
1014
1015 if !is_dictionary_page {
1016 *page_index += 1;
1017 }
1018 decode_page(
1019 header,
1020 bytes,
1021 self.physical_type,
1022 self.decompressor.as_mut(),
1023 )?
1024 }
1025 };
1026
1027 return Ok(Some(page));
1028 }
1029 }
1030
1031 fn peek_next_page(&mut self) -> Result<Option<PageMetadata>> {
1032 match &mut self.state {
1033 SerializedPageReaderState::Values {
1034 offset,
1035 remaining_bytes,
1036 next_page_header,
1037 page_index,
1038 require_dictionary,
1039 } => {
1040 loop {
1041 if *remaining_bytes == 0 {
1042 return Ok(None);
1043 }
1044 return if let Some(header) = next_page_header.as_ref() {
1045 if let Ok(page_meta) = (&**header).try_into() {
1046 Ok(Some(page_meta))
1047 } else {
1048 *next_page_header = None;
1050 continue;
1051 }
1052 } else {
1053 let mut read = self.reader.get_read(*offset)?;
1054 let (header_len, header) = Self::read_page_header_len(
1055 &self.context,
1056 &mut read,
1057 *page_index,
1058 *require_dictionary,
1059 )?;
1060 verify_page_header_len(header_len, *remaining_bytes)?;
1061 *offset += header_len as u64;
1062 *remaining_bytes -= header_len as u64;
1063 let page_meta = if let Ok(page_meta) = (&header).try_into() {
1064 Ok(Some(page_meta))
1065 } else {
1066 continue;
1068 };
1069 *next_page_header = Some(Box::new(header));
1070 page_meta
1071 };
1072 }
1073 }
1074 SerializedPageReaderState::Pages {
1075 page_locations,
1076 dictionary_page,
1077 total_rows,
1078 page_index: _,
1079 } => {
1080 if dictionary_page.is_some() {
1081 Ok(Some(PageMetadata {
1082 num_rows: None,
1083 num_levels: None,
1084 is_dict: true,
1085 }))
1086 } else if let Some(page) = page_locations.front() {
1087 let next_rows = page_locations
1088 .get(1)
1089 .map(|x| x.first_row_index as usize)
1090 .unwrap_or(*total_rows);
1091
1092 Ok(Some(PageMetadata {
1093 num_rows: Some(next_rows - page.first_row_index as usize),
1094 num_levels: None,
1095 is_dict: false,
1096 }))
1097 } else {
1098 Ok(None)
1099 }
1100 }
1101 }
1102 }
1103
1104 fn skip_next_page(&mut self) -> Result<()> {
1105 match &mut self.state {
1106 SerializedPageReaderState::Values {
1107 offset,
1108 remaining_bytes,
1109 next_page_header,
1110 page_index,
1111 require_dictionary,
1112 } => {
1113 if let Some(buffered_header) = next_page_header.take() {
1114 verify_page_size(
1115 buffered_header.compressed_page_size,
1116 buffered_header.uncompressed_page_size,
1117 *remaining_bytes,
1118 )?;
1119 *offset += buffered_header.compressed_page_size as u64;
1121 *remaining_bytes -= buffered_header.compressed_page_size as u64;
1122 } else {
1123 let mut read = self.reader.get_read(*offset)?;
1124 let (header_len, header) = Self::read_page_header_len(
1125 &self.context,
1126 &mut read,
1127 *page_index,
1128 *require_dictionary,
1129 )?;
1130 verify_page_header_len(header_len, *remaining_bytes)?;
1131 verify_page_size(
1132 header.compressed_page_size,
1133 header.uncompressed_page_size,
1134 *remaining_bytes,
1135 )?;
1136 let data_page_size = header.compressed_page_size as u64;
1137 *offset += header_len as u64 + data_page_size;
1138 *remaining_bytes -= header_len as u64 + data_page_size;
1139 }
1140 if *require_dictionary {
1141 *require_dictionary = false;
1142 } else {
1143 *page_index += 1;
1144 }
1145 Ok(())
1146 }
1147 SerializedPageReaderState::Pages {
1148 page_locations,
1149 dictionary_page,
1150 page_index,
1151 ..
1152 } => {
1153 if dictionary_page.is_some() {
1154 dictionary_page.take();
1156 } else {
1157 if page_locations.pop_front().is_some() {
1159 *page_index += 1;
1160 }
1161 }
1162
1163 Ok(())
1164 }
1165 }
1166 }
1167
1168 fn at_record_boundary(&mut self) -> Result<bool> {
1169 match &mut self.state {
1170 SerializedPageReaderState::Values { .. } => match self.peek_next_page()? {
1171 None => Ok(true),
1172 Some(metadata) => Ok(metadata.num_rows.is_some()),
1175 },
1176 SerializedPageReaderState::Pages { .. } => Ok(true),
1177 }
1178 }
1179}
1180
1181#[cfg(test)]
1182mod tests {
1183 use std::collections::HashSet;
1184
1185 use bytes::Buf;
1186
1187 use crate::file::page_index::column_index::{
1188 ByteArrayColumnIndex, ColumnIndexMetaData, PrimitiveColumnIndex,
1189 };
1190 use crate::file::properties::{EnabledStatistics, WriterProperties};
1191
1192 use crate::basic::{self, BoundaryOrder, ColumnOrder, Encoding, SortOrder};
1193 use crate::column::reader::ColumnReader;
1194 use crate::data_type::private::ParquetValueType;
1195 use crate::data_type::{AsBytes, FixedLenByteArrayType, Int32Type};
1196 use crate::file::metadata::thrift::DataPageHeaderV2;
1197 use crate::file::writer::SerializedFileWriter;
1198 use crate::record::RowAccessor;
1199 use crate::schema::parser::parse_message_type;
1200 use crate::util::test_common::file_util::{get_test_file, get_test_path};
1201
1202 use super::*;
1203
1204 #[test]
1205 fn test_decode_page_invalid_offset() {
1206 let page_header = PageHeader {
1207 r#type: PageType::DATA_PAGE_V2,
1208 uncompressed_page_size: 10,
1209 compressed_page_size: 10,
1210 data_page_header: None,
1211 index_page_header: None,
1212 dictionary_page_header: None,
1213 crc: None,
1214 data_page_header_v2: Some(DataPageHeaderV2 {
1215 num_nulls: 0,
1216 num_rows: 0,
1217 num_values: 0,
1218 encoding: Encoding::PLAIN,
1219 definition_levels_byte_length: 11,
1220 repetition_levels_byte_length: 0,
1221 is_compressed: None,
1222 statistics: None,
1223 }),
1224 };
1225
1226 let buffer = Bytes::new();
1227 let err = decode_page(page_header, buffer, Type::INT32, None).unwrap_err();
1228 assert!(
1229 err.to_string()
1230 .contains("DataPage v2 header contains implausible values")
1231 );
1232 }
1233
1234 #[test]
1235 fn test_decode_unsupported_page() {
1236 let mut page_header = PageHeader {
1237 r#type: PageType::INDEX_PAGE,
1238 uncompressed_page_size: 10,
1239 compressed_page_size: 10,
1240 data_page_header: None,
1241 index_page_header: None,
1242 dictionary_page_header: None,
1243 crc: None,
1244 data_page_header_v2: None,
1245 };
1246 let buffer = Bytes::new();
1247 let err = decode_page(page_header.clone(), buffer.clone(), Type::INT32, None).unwrap_err();
1248 assert_eq!(
1249 err.to_string(),
1250 "Parquet error: Page type INDEX_PAGE is not supported"
1251 );
1252
1253 page_header.data_page_header_v2 = Some(DataPageHeaderV2 {
1254 num_nulls: 0,
1255 num_rows: 0,
1256 num_values: 0,
1257 encoding: Encoding::PLAIN,
1258 definition_levels_byte_length: 11,
1259 repetition_levels_byte_length: 0,
1260 is_compressed: None,
1261 statistics: None,
1262 });
1263 let err = decode_page(page_header, buffer, Type::INT32, None).unwrap_err();
1264 assert!(
1265 err.to_string()
1266 .contains("DataPage v2 header contains implausible values")
1267 );
1268 }
1269
1270 #[test]
1271 fn test_cursor_and_file_has_the_same_behaviour() {
1272 let mut buf: Vec<u8> = Vec::new();
1273 get_test_file("alltypes_plain.parquet")
1274 .read_to_end(&mut buf)
1275 .unwrap();
1276 let cursor = Bytes::from(buf);
1277 let read_from_cursor = SerializedFileReader::new(cursor).unwrap();
1278
1279 let test_file = get_test_file("alltypes_plain.parquet");
1280 let read_from_file = SerializedFileReader::new(test_file).unwrap();
1281
1282 let file_iter = read_from_file.get_row_iter(None).unwrap();
1283 let cursor_iter = read_from_cursor.get_row_iter(None).unwrap();
1284
1285 for (a, b) in file_iter.zip(cursor_iter) {
1286 assert_eq!(a.unwrap(), b.unwrap())
1287 }
1288 }
1289
1290 #[test]
1291 fn test_file_reader_try_from() {
1292 let test_file = get_test_file("alltypes_plain.parquet");
1294 let test_path_buf = get_test_path("alltypes_plain.parquet");
1295 let test_path = test_path_buf.as_path();
1296 let test_path_str = test_path.to_str().unwrap();
1297
1298 let reader = SerializedFileReader::try_from(test_file);
1299 assert!(reader.is_ok());
1300
1301 let reader = SerializedFileReader::try_from(test_path);
1302 assert!(reader.is_ok());
1303
1304 let reader = SerializedFileReader::try_from(test_path_str);
1305 assert!(reader.is_ok());
1306
1307 let reader = SerializedFileReader::try_from(test_path_str.to_string());
1308 assert!(reader.is_ok());
1309
1310 let test_path = Path::new("invalid.parquet");
1312 let test_path_str = test_path.to_str().unwrap();
1313
1314 let reader = SerializedFileReader::try_from(test_path);
1315 assert!(reader.is_err());
1316
1317 let reader = SerializedFileReader::try_from(test_path_str);
1318 assert!(reader.is_err());
1319
1320 let reader = SerializedFileReader::try_from(test_path_str.to_string());
1321 assert!(reader.is_err());
1322 }
1323
1324 #[test]
1325 fn test_file_reader_into_iter() {
1326 let path = get_test_path("alltypes_plain.parquet");
1327 let reader = SerializedFileReader::try_from(path.as_path()).unwrap();
1328 let iter = reader.into_iter();
1329 let values: Vec<_> = iter.flat_map(|x| x.unwrap().get_int(0)).collect();
1330
1331 assert_eq!(values, &[4, 5, 6, 7, 2, 3, 0, 1]);
1332 }
1333
1334 #[test]
1335 fn test_file_reader_into_iter_project() {
1336 let path = get_test_path("alltypes_plain.parquet");
1337 let reader = SerializedFileReader::try_from(path.as_path()).unwrap();
1338 let schema = "message schema { OPTIONAL INT32 id; }";
1339 let proj = parse_message_type(schema).ok();
1340 let iter = reader.into_iter().project(proj).unwrap();
1341 let values: Vec<_> = iter.flat_map(|x| x.unwrap().get_int(0)).collect();
1342
1343 assert_eq!(values, &[4, 5, 6, 7, 2, 3, 0, 1]);
1344 }
1345
1346 #[test]
1347 fn test_reuse_file_chunk() {
1348 let test_file = get_test_file("alltypes_plain.parquet");
1352 let reader = SerializedFileReader::new(test_file).unwrap();
1353 let row_group = reader.get_row_group(0).unwrap();
1354
1355 let mut page_readers = Vec::new();
1356 for i in 0..row_group.num_columns() {
1357 page_readers.push(row_group.get_column_page_reader(i).unwrap());
1358 }
1359
1360 for mut page_reader in page_readers {
1363 assert!(page_reader.get_next_page().is_ok());
1364 }
1365 }
1366
1367 #[test]
1368 fn test_file_reader() {
1369 let test_file = get_test_file("alltypes_plain.parquet");
1370 let reader_result = SerializedFileReader::new(test_file);
1371 assert!(reader_result.is_ok());
1372 let reader = reader_result.unwrap();
1373
1374 let metadata = reader.metadata();
1376 assert_eq!(metadata.num_row_groups(), 1);
1377
1378 let file_metadata = metadata.file_metadata();
1380 assert!(file_metadata.created_by().is_some());
1381 assert_eq!(
1382 file_metadata.created_by().unwrap(),
1383 "impala version 1.3.0-INTERNAL (build 8a48ddb1eff84592b3fc06bc6f51ec120e1fffc9)"
1384 );
1385 assert!(file_metadata.key_value_metadata().is_none());
1386 assert_eq!(file_metadata.num_rows(), 8);
1387 assert_eq!(file_metadata.version(), 1);
1388 assert_eq!(file_metadata.column_orders(), None);
1389
1390 let row_group_metadata = metadata.row_group(0);
1392 assert_eq!(row_group_metadata.num_columns(), 11);
1393 assert_eq!(row_group_metadata.num_rows(), 8);
1394 assert_eq!(row_group_metadata.total_byte_size(), 671);
1395 for i in 0..row_group_metadata.num_columns() {
1397 assert_eq!(file_metadata.column_order(i), ColumnOrder::UNDEFINED);
1398 }
1399
1400 let row_group_reader_result = reader.get_row_group(0);
1402 assert!(row_group_reader_result.is_ok());
1403 let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1404 assert_eq!(
1405 row_group_reader.num_columns(),
1406 row_group_metadata.num_columns()
1407 );
1408 assert_eq!(
1409 row_group_reader.metadata().total_byte_size(),
1410 row_group_metadata.total_byte_size()
1411 );
1412
1413 let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1416 assert!(page_reader_0_result.is_ok());
1417 let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1418 let mut page_count = 0;
1419 while let Some(page) = page_reader_0.get_next_page().unwrap() {
1420 let is_expected_page = match page {
1421 Page::DictionaryPage {
1422 buf,
1423 num_values,
1424 encoding,
1425 is_sorted,
1426 } => {
1427 assert_eq!(buf.len(), 32);
1428 assert_eq!(num_values, 8);
1429 assert_eq!(encoding, Encoding::PLAIN_DICTIONARY);
1430 assert!(!is_sorted);
1431 true
1432 }
1433 Page::DataPage {
1434 buf,
1435 num_values,
1436 encoding,
1437 def_level_encoding,
1438 rep_level_encoding,
1439 statistics,
1440 } => {
1441 assert_eq!(buf.len(), 11);
1442 assert_eq!(num_values, 8);
1443 assert_eq!(encoding, Encoding::PLAIN_DICTIONARY);
1444 assert_eq!(def_level_encoding, Encoding::RLE);
1445 #[expect(deprecated)]
1446 let expected_rep_level_encoding = Encoding::BIT_PACKED;
1447 assert_eq!(rep_level_encoding, expected_rep_level_encoding);
1448 assert!(statistics.is_none());
1449 true
1450 }
1451 Page::DataPageV2 { .. } => false,
1452 };
1453 assert!(is_expected_page);
1454 page_count += 1;
1455 }
1456 assert_eq!(page_count, 2);
1457 }
1458
1459 #[test]
1460 fn test_file_reader_datapage_v2() {
1461 let test_file = get_test_file("datapage_v2.snappy.parquet");
1462 let reader_result = SerializedFileReader::new(test_file);
1463 assert!(reader_result.is_ok());
1464 let reader = reader_result.unwrap();
1465
1466 let metadata = reader.metadata();
1468 assert_eq!(metadata.num_row_groups(), 1);
1469
1470 let file_metadata = metadata.file_metadata();
1472 assert!(file_metadata.created_by().is_some());
1473 assert_eq!(
1474 file_metadata.created_by().unwrap(),
1475 "parquet-mr version 1.8.1 (build 4aba4dae7bb0d4edbcf7923ae1339f28fd3f7fcf)"
1476 );
1477 assert!(file_metadata.key_value_metadata().is_some());
1478 assert_eq!(file_metadata.key_value_metadata().unwrap().len(), 1);
1479
1480 assert_eq!(file_metadata.num_rows(), 5);
1481 assert_eq!(file_metadata.version(), 1);
1482 assert_eq!(file_metadata.column_orders(), None);
1483
1484 let row_group_metadata = metadata.row_group(0);
1485
1486 for i in 0..row_group_metadata.num_columns() {
1488 assert_eq!(file_metadata.column_order(i), ColumnOrder::UNDEFINED);
1489 }
1490
1491 let row_group_reader_result = reader.get_row_group(0);
1493 assert!(row_group_reader_result.is_ok());
1494 let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1495 assert_eq!(
1496 row_group_reader.num_columns(),
1497 row_group_metadata.num_columns()
1498 );
1499 assert_eq!(
1500 row_group_reader.metadata().total_byte_size(),
1501 row_group_metadata.total_byte_size()
1502 );
1503
1504 let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1507 assert!(page_reader_0_result.is_ok());
1508 let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1509 let mut page_count = 0;
1510 while let Some(page) = page_reader_0.get_next_page().unwrap() {
1511 let is_expected_page = match page {
1512 Page::DictionaryPage {
1513 buf,
1514 num_values,
1515 encoding,
1516 is_sorted,
1517 } => {
1518 assert_eq!(buf.len(), 7);
1519 assert_eq!(num_values, 1);
1520 assert_eq!(encoding, Encoding::PLAIN);
1521 assert!(!is_sorted);
1522 true
1523 }
1524 Page::DataPageV2 {
1525 buf,
1526 num_values,
1527 encoding,
1528 num_nulls,
1529 num_rows,
1530 def_levels_byte_len,
1531 rep_levels_byte_len,
1532 is_compressed,
1533 statistics,
1534 } => {
1535 assert_eq!(buf.len(), 4);
1536 assert_eq!(num_values, 5);
1537 assert_eq!(encoding, Encoding::RLE_DICTIONARY);
1538 assert_eq!(num_nulls, 1);
1539 assert_eq!(num_rows, 5);
1540 assert_eq!(def_levels_byte_len, 2);
1541 assert_eq!(rep_levels_byte_len, 0);
1542 assert!(is_compressed);
1543 assert!(statistics.is_none()); true
1545 }
1546 Page::DataPage { .. } => false,
1547 };
1548 assert!(is_expected_page);
1549 page_count += 1;
1550 }
1551 assert_eq!(page_count, 2);
1552 }
1553
1554 #[cfg_attr(miri, ignore)] #[test]
1556 fn test_file_reader_empty_compressed_datapage_v2() {
1557 let test_file = get_test_file("page_v2_empty_compressed.parquet");
1559 let reader_result = SerializedFileReader::new(test_file);
1560 assert!(reader_result.is_ok());
1561 let reader = reader_result.unwrap();
1562
1563 let metadata = reader.metadata();
1565 assert_eq!(metadata.num_row_groups(), 1);
1566
1567 let file_metadata = metadata.file_metadata();
1569 assert!(file_metadata.created_by().is_some());
1570 assert_eq!(
1571 file_metadata.created_by().unwrap(),
1572 "parquet-cpp-arrow version 14.0.2"
1573 );
1574 assert!(file_metadata.key_value_metadata().is_some());
1575 assert_eq!(file_metadata.key_value_metadata().unwrap().len(), 1);
1576
1577 assert_eq!(file_metadata.num_rows(), 10);
1578 assert_eq!(file_metadata.version(), 2);
1579 let expected_order = ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED);
1580 assert_eq!(
1581 file_metadata.column_orders(),
1582 Some(vec![expected_order].as_ref())
1583 );
1584
1585 let row_group_metadata = metadata.row_group(0);
1586
1587 for i in 0..row_group_metadata.num_columns() {
1589 assert_eq!(file_metadata.column_order(i), expected_order);
1590 }
1591
1592 let row_group_reader_result = reader.get_row_group(0);
1594 assert!(row_group_reader_result.is_ok());
1595 let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1596 assert_eq!(
1597 row_group_reader.num_columns(),
1598 row_group_metadata.num_columns()
1599 );
1600 assert_eq!(
1601 row_group_reader.metadata().total_byte_size(),
1602 row_group_metadata.total_byte_size()
1603 );
1604
1605 let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1607 assert!(page_reader_0_result.is_ok());
1608 let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1609 let mut page_count = 0;
1610 while let Some(page) = page_reader_0.get_next_page().unwrap() {
1611 let is_expected_page = match page {
1612 Page::DictionaryPage {
1613 buf,
1614 num_values,
1615 encoding,
1616 is_sorted,
1617 } => {
1618 assert_eq!(buf.len(), 0);
1619 assert_eq!(num_values, 0);
1620 assert_eq!(encoding, Encoding::PLAIN);
1621 assert!(!is_sorted);
1622 true
1623 }
1624 Page::DataPageV2 {
1625 buf,
1626 num_values,
1627 encoding,
1628 num_nulls,
1629 num_rows,
1630 def_levels_byte_len,
1631 rep_levels_byte_len,
1632 is_compressed,
1633 statistics,
1634 } => {
1635 assert_eq!(buf.len(), 3);
1636 assert_eq!(num_values, 10);
1637 assert_eq!(encoding, Encoding::RLE_DICTIONARY);
1638 assert_eq!(num_nulls, 10);
1639 assert_eq!(num_rows, 10);
1640 assert_eq!(def_levels_byte_len, 2);
1641 assert_eq!(rep_levels_byte_len, 0);
1642 assert!(is_compressed);
1643 assert!(statistics.is_none()); true
1645 }
1646 Page::DataPage { .. } => false,
1647 };
1648 assert!(is_expected_page);
1649 page_count += 1;
1650 }
1651 assert_eq!(page_count, 2);
1652 }
1653
1654 #[test]
1655 fn test_file_reader_empty_datapage_v2() {
1656 let test_file = get_test_file("datapage_v2_empty_datapage.snappy.parquet");
1658 let reader_result = SerializedFileReader::new(test_file);
1659 assert!(reader_result.is_ok());
1660 let reader = reader_result.unwrap();
1661
1662 let metadata = reader.metadata();
1664 assert_eq!(metadata.num_row_groups(), 1);
1665
1666 let file_metadata = metadata.file_metadata();
1668 assert!(file_metadata.created_by().is_some());
1669 assert_eq!(
1670 file_metadata.created_by().unwrap(),
1671 "parquet-mr version 1.13.1 (build db4183109d5b734ec5930d870cdae161e408ddba)"
1672 );
1673 assert!(file_metadata.key_value_metadata().is_some());
1674 assert_eq!(file_metadata.key_value_metadata().unwrap().len(), 2);
1675
1676 assert_eq!(file_metadata.num_rows(), 1);
1677 assert_eq!(file_metadata.version(), 1);
1678 let expected_order = ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED);
1679 assert_eq!(
1680 file_metadata.column_orders(),
1681 Some(vec![expected_order].as_ref())
1682 );
1683
1684 let row_group_metadata = metadata.row_group(0);
1685
1686 for i in 0..row_group_metadata.num_columns() {
1688 assert_eq!(file_metadata.column_order(i), expected_order);
1689 }
1690
1691 let row_group_reader_result = reader.get_row_group(0);
1693 assert!(row_group_reader_result.is_ok());
1694 let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1695 assert_eq!(
1696 row_group_reader.num_columns(),
1697 row_group_metadata.num_columns()
1698 );
1699 assert_eq!(
1700 row_group_reader.metadata().total_byte_size(),
1701 row_group_metadata.total_byte_size()
1702 );
1703
1704 let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1706 assert!(page_reader_0_result.is_ok());
1707 let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1708 let mut page_count = 0;
1709 while let Some(page) = page_reader_0.get_next_page().unwrap() {
1710 let is_expected_page = match page {
1711 Page::DataPageV2 {
1712 buf,
1713 num_values,
1714 encoding,
1715 num_nulls,
1716 num_rows,
1717 def_levels_byte_len,
1718 rep_levels_byte_len,
1719 is_compressed,
1720 statistics,
1721 } => {
1722 assert_eq!(buf.len(), 2);
1723 assert_eq!(num_values, 1);
1724 assert_eq!(encoding, Encoding::PLAIN);
1725 assert_eq!(num_nulls, 1);
1726 assert_eq!(num_rows, 1);
1727 assert_eq!(def_levels_byte_len, 2);
1728 assert_eq!(rep_levels_byte_len, 0);
1729 assert!(is_compressed);
1730 assert!(statistics.is_none());
1731 true
1732 }
1733 _ => false,
1734 };
1735 assert!(is_expected_page);
1736 page_count += 1;
1737 }
1738 assert_eq!(page_count, 1);
1739 }
1740
1741 fn get_serialized_page_reader<R: ChunkReader>(
1742 file_reader: &SerializedFileReader<R>,
1743 row_group_idx: usize,
1744 column: usize,
1745 ) -> Result<SerializedPageReader<R>> {
1746 let row_group = {
1747 let row_group_metadata = file_reader.metadata.row_group(row_group_idx);
1748 let props = Arc::clone(&file_reader.props);
1749 let f = Arc::clone(&file_reader.chunk_reader);
1750 let page_index = file_reader.metadata.page_index_for_row_group(row_group_idx);
1751 SerializedRowGroupReader::new(f, row_group_metadata, page_index, props)?
1752 };
1753
1754 let col = row_group.metadata.column(column);
1755 let page_locations = file_reader
1756 .metadata
1757 .page_index()
1758 .map(|pi| pi.page_locations(row_group_idx, column))
1759 .unwrap_or(None)
1760 .cloned();
1761
1762 let props = Arc::clone(&row_group.props);
1763 SerializedPageReader::new_with_properties(
1764 Arc::clone(&row_group.chunk_reader),
1765 col,
1766 usize::try_from(row_group.metadata.num_rows())?,
1767 page_locations,
1768 props,
1769 )
1770 }
1771
1772 #[test]
1773 fn test_peek_next_page_offset_matches_actual() -> Result<()> {
1774 let test_file = get_test_file("alltypes_plain.parquet");
1775 let reader = SerializedFileReader::new(test_file)?;
1776
1777 let mut offset_set = HashSet::new();
1778 let num_row_groups = reader.metadata.num_row_groups();
1779 for row_group in 0..num_row_groups {
1780 let num_columns = reader.metadata.row_group(row_group).num_columns();
1781 for column in 0..num_columns {
1782 let mut page_reader = get_serialized_page_reader(&reader, row_group, column)?;
1783
1784 while let Ok(Some(page_offset)) = page_reader.peek_next_page_offset() {
1785 match &page_reader.state {
1786 SerializedPageReaderState::Pages {
1787 page_locations,
1788 dictionary_page,
1789 ..
1790 } => {
1791 if let Some(page) = dictionary_page {
1792 assert_eq!(page.offset as u64, page_offset);
1793 } else if let Some(page) = page_locations.front() {
1794 assert_eq!(page.offset as u64, page_offset);
1795 } else {
1796 unreachable!()
1797 }
1798 }
1799 SerializedPageReaderState::Values {
1800 offset,
1801 next_page_header,
1802 ..
1803 } => {
1804 assert!(next_page_header.is_some());
1805 assert_eq!(*offset, page_offset);
1806 }
1807 }
1808 let page = page_reader.get_next_page()?;
1809 assert!(page.is_some());
1810 let newly_inserted = offset_set.insert(page_offset);
1811 assert!(newly_inserted);
1812 }
1813 }
1814 }
1815
1816 Ok(())
1817 }
1818
1819 #[test]
1820 fn test_page_iterator() {
1821 let file = get_test_file("alltypes_plain.parquet");
1822 let file_reader = Arc::new(SerializedFileReader::new(file).unwrap());
1823
1824 let mut page_iterator = FilePageIterator::new(0, file_reader.clone()).unwrap();
1825
1826 let page = page_iterator.next();
1828 assert!(page.is_some());
1829 assert!(page.unwrap().is_ok());
1830
1831 let page = page_iterator.next();
1833 assert!(page.is_none());
1834
1835 let row_group_indices = Box::new(0..1);
1836 let mut page_iterator =
1837 FilePageIterator::with_row_groups(0, row_group_indices, file_reader).unwrap();
1838
1839 let page = page_iterator.next();
1841 assert!(page.is_some());
1842 assert!(page.unwrap().is_ok());
1843
1844 let page = page_iterator.next();
1846 assert!(page.is_none());
1847 }
1848
1849 #[test]
1850 fn test_file_reader_key_value_metadata() {
1851 let file = get_test_file("binary.parquet");
1852 let file_reader = Arc::new(SerializedFileReader::new(file).unwrap());
1853
1854 let metadata = file_reader
1855 .metadata
1856 .file_metadata()
1857 .key_value_metadata()
1858 .unwrap();
1859
1860 assert_eq!(metadata.len(), 3);
1861
1862 assert_eq!(metadata[0].key, "parquet.proto.descriptor");
1863
1864 assert_eq!(metadata[1].key, "writer.model.name");
1865 assert_eq!(metadata[1].value, Some("protobuf".to_owned()));
1866
1867 assert_eq!(metadata[2].key, "parquet.proto.class");
1868 assert_eq!(metadata[2].value, Some("foo.baz.Foobaz$Event".to_owned()));
1869 }
1870
1871 #[test]
1872 fn test_file_reader_optional_metadata() {
1873 let file = get_test_file("data_index_bloom_encoding_stats.parquet");
1875 let options = ReadOptionsBuilder::new()
1876 .with_encoding_stats_as_mask(false)
1877 .build();
1878 let file_reader = Arc::new(SerializedFileReader::new_with_options(file, options).unwrap());
1879
1880 let row_group_metadata = file_reader.metadata.row_group(0);
1881 let col0_metadata = row_group_metadata.column(0);
1882
1883 assert_eq!(col0_metadata.bloom_filter_offset().unwrap(), 192);
1885
1886 let page_encoding_stats = &col0_metadata.page_encoding_stats().unwrap()[0];
1888
1889 assert_eq!(page_encoding_stats.page_type, basic::PageType::DATA_PAGE);
1890 assert_eq!(page_encoding_stats.encoding, Encoding::PLAIN);
1891 assert_eq!(page_encoding_stats.count, 1);
1892
1893 assert_eq!(col0_metadata.column_index_offset().unwrap(), 156);
1895 assert_eq!(col0_metadata.column_index_length().unwrap(), 25);
1896
1897 assert_eq!(col0_metadata.offset_index_offset().unwrap(), 181);
1899 assert_eq!(col0_metadata.offset_index_length().unwrap(), 11);
1900 }
1901
1902 #[test]
1903 fn test_file_reader_page_stats_mask() {
1904 let file = get_test_file("alltypes_tiny_pages.parquet");
1905 let options = ReadOptionsBuilder::new()
1906 .with_encoding_stats_as_mask(true)
1907 .build();
1908 let file_reader = Arc::new(SerializedFileReader::new_with_options(file, options).unwrap());
1909
1910 let row_group_metadata = file_reader.metadata.row_group(0);
1911
1912 let page_encoding_stats = row_group_metadata
1914 .column(0)
1915 .page_encoding_stats_mask()
1916 .unwrap();
1917 assert!(page_encoding_stats.is_only(Encoding::PLAIN));
1918 let page_encoding_stats = row_group_metadata
1919 .column(2)
1920 .page_encoding_stats_mask()
1921 .unwrap();
1922 assert!(page_encoding_stats.is_only(Encoding::PLAIN_DICTIONARY));
1923 }
1924
1925 #[test]
1926 fn test_file_reader_page_stats_skipped() {
1927 let file = get_test_file("alltypes_tiny_pages.parquet");
1928
1929 let options = ReadOptionsBuilder::new()
1931 .with_encoding_stats_policy(ParquetStatisticsPolicy::SkipAll)
1932 .with_column_stats_policy(ParquetStatisticsPolicy::SkipAll)
1933 .build();
1934 let file_reader = Arc::new(
1935 SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1936 );
1937
1938 let row_group_metadata = file_reader.metadata.row_group(0);
1939 for column in row_group_metadata.columns() {
1940 assert!(column.page_encoding_stats().is_none());
1941 assert!(column.page_encoding_stats_mask().is_none());
1942 assert!(column.statistics().is_none());
1943 }
1944
1945 let options = ReadOptionsBuilder::new()
1947 .with_encoding_stats_as_mask(true)
1948 .with_encoding_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
1949 .with_column_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
1950 .build();
1951 let file_reader = Arc::new(
1952 SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1953 );
1954
1955 let row_group_metadata = file_reader.metadata.row_group(0);
1956 for (idx, column) in row_group_metadata.columns().iter().enumerate() {
1957 assert!(column.page_encoding_stats().is_none());
1958 assert_eq!(column.page_encoding_stats_mask().is_some(), idx == 0);
1959 assert_eq!(column.statistics().is_some(), idx == 0);
1960 }
1961 }
1962
1963 #[test]
1964 fn test_file_reader_size_stats_skipped() {
1965 let file = get_test_file("repeated_primitive_no_list.parquet");
1966
1967 let options = ReadOptionsBuilder::new()
1969 .with_size_stats_policy(ParquetStatisticsPolicy::SkipAll)
1970 .build();
1971 let file_reader = Arc::new(
1972 SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1973 );
1974
1975 let row_group_metadata = file_reader.metadata.row_group(0);
1976 for column in row_group_metadata.columns() {
1977 assert!(column.repetition_level_histogram().is_none());
1978 assert!(column.definition_level_histogram().is_none());
1979 assert!(column.unencoded_byte_array_data_bytes().is_none());
1980 }
1981
1982 let options = ReadOptionsBuilder::new()
1984 .with_encoding_stats_as_mask(true)
1985 .with_size_stats_policy(ParquetStatisticsPolicy::skip_except(&[1]))
1986 .build();
1987 let file_reader = Arc::new(
1988 SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1989 );
1990
1991 let row_group_metadata = file_reader.metadata.row_group(0);
1992 for (idx, column) in row_group_metadata.columns().iter().enumerate() {
1993 assert_eq!(column.repetition_level_histogram().is_some(), idx == 1);
1994 assert_eq!(column.definition_level_histogram().is_some(), idx == 1);
1995 assert_eq!(column.unencoded_byte_array_data_bytes().is_some(), idx == 1);
1996 }
1997 }
1998
1999 #[test]
2000 fn test_file_reader_with_no_filter() -> Result<()> {
2001 let test_file = get_test_file("alltypes_plain.parquet");
2002 let origin_reader = SerializedFileReader::new(test_file)?;
2003 let metadata = origin_reader.metadata();
2005 assert_eq!(metadata.num_row_groups(), 1);
2006 Ok(())
2007 }
2008
2009 #[test]
2010 fn test_file_reader_filter_row_groups_with_predicate() -> Result<()> {
2011 let test_file = get_test_file("alltypes_plain.parquet");
2012 let read_options = ReadOptionsBuilder::new()
2013 .with_predicate(Box::new(|_, _| false))
2014 .build();
2015 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2016 let metadata = reader.metadata();
2017 assert_eq!(metadata.num_row_groups(), 0);
2018 Ok(())
2019 }
2020
2021 #[test]
2022 fn test_file_reader_filter_row_groups_with_range() -> Result<()> {
2023 let test_file = get_test_file("alltypes_plain.parquet");
2024 let origin_reader = SerializedFileReader::new(test_file)?;
2025 let metadata = origin_reader.metadata();
2027 assert_eq!(metadata.num_row_groups(), 1);
2028 let mid = get_midpoint_offset(metadata.row_group(0));
2029
2030 let test_file = get_test_file("alltypes_plain.parquet");
2031 let read_options = ReadOptionsBuilder::new().with_range(0, mid + 1).build();
2032 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2033 let metadata = reader.metadata();
2034 assert_eq!(metadata.num_row_groups(), 1);
2035
2036 let test_file = get_test_file("alltypes_plain.parquet");
2037 let read_options = ReadOptionsBuilder::new().with_range(0, mid).build();
2038 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2039 let metadata = reader.metadata();
2040 assert_eq!(metadata.num_row_groups(), 0);
2041 Ok(())
2042 }
2043
2044 #[test]
2045 #[cfg_attr(miri, ignore)] fn test_file_reader_filter_row_groups_and_range() -> Result<()> {
2047 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2048 let origin_reader = SerializedFileReader::new(test_file)?;
2049 let metadata = origin_reader.metadata();
2050 let mid = get_midpoint_offset(metadata.row_group(0));
2051
2052 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2054 let read_options = ReadOptionsBuilder::new()
2055 .with_page_index()
2056 .with_predicate(Box::new(|_, _| true))
2057 .with_range(mid, mid + 1)
2058 .build();
2059 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2060 let metadata = reader.metadata();
2061 assert_eq!(metadata.num_row_groups(), 1);
2062 assert!(metadata.page_index().is_some());
2063
2064 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2066 let read_options = ReadOptionsBuilder::new()
2067 .with_page_index()
2068 .with_predicate(Box::new(|_, _| true))
2069 .with_range(0, mid)
2070 .build();
2071 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2072 let metadata = reader.metadata();
2073 assert_eq!(metadata.num_row_groups(), 0);
2074 assert!(metadata.page_index().is_none());
2075
2076 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2078 let read_options = ReadOptionsBuilder::new()
2079 .with_page_index()
2080 .with_predicate(Box::new(|_, _| false))
2081 .with_range(mid, mid + 1)
2082 .build();
2083 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2084 let metadata = reader.metadata();
2085 assert_eq!(metadata.num_row_groups(), 0);
2086 assert!(metadata.page_index().is_none());
2087
2088 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2090 let read_options = ReadOptionsBuilder::new()
2091 .with_page_index()
2092 .with_predicate(Box::new(|_, _| false))
2093 .with_range(0, mid)
2094 .build();
2095 let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2096 let metadata = reader.metadata();
2097 assert_eq!(metadata.num_row_groups(), 0);
2098 assert!(metadata.page_index().is_none());
2099 Ok(())
2100 }
2101
2102 #[test]
2103 fn test_file_reader_invalid_metadata() {
2104 let data = [
2105 255, 172, 1, 0, 50, 82, 65, 73, 1, 0, 0, 0, 169, 168, 168, 162, 87, 255, 16, 0, 0, 0,
2106 80, 65, 82, 49,
2107 ];
2108 let ret = SerializedFileReader::new(Bytes::copy_from_slice(&data));
2109 assert_eq!(
2110 ret.err().unwrap().to_string(),
2111 "Parquet error: Expected list element type of Struct but got List"
2112 );
2113 }
2114
2115 #[test]
2116 fn test_page_index_reader() {
2133 let test_file = get_test_file("data_index_bloom_encoding_stats.parquet");
2134 let builder = ReadOptionsBuilder::new();
2135 let options = builder.with_page_index().build();
2137 let reader_result = SerializedFileReader::new_with_options(test_file, options);
2138 let reader = reader_result.unwrap();
2139
2140 let metadata = reader.metadata();
2142 assert_eq!(metadata.num_row_groups(), 1);
2143
2144 let page_index = metadata.page_index().expect("page index should be present");
2145
2146 let Some(ColumnIndexMetaData::BYTE_ARRAY(index)) = page_index.column_index(0, 0) else {
2148 unreachable!()
2149 };
2150
2151 assert_eq!(index.boundary_order, BoundaryOrder::ASCENDING);
2152
2153 assert_eq!(index.num_pages(), 1);
2155
2156 let min = index.min_value(0).unwrap();
2157 let max = index.max_value(0).unwrap();
2158 assert_eq!(b"Hello", min.as_bytes());
2159 assert_eq!(b"today", max.as_bytes());
2160
2161 let offset_index = page_index
2163 .offset_index(0, 0)
2164 .expect("offset index should be present");
2165 let page_offset = offset_index
2166 .page_locations()
2167 .first()
2168 .expect("offset index too small");
2169
2170 assert_eq!(4, page_offset.offset);
2171 assert_eq!(152, page_offset.compressed_page_size);
2172 assert_eq!(0, page_offset.first_row_index);
2173 }
2174
2175 #[test]
2176 #[cfg_attr(miri, ignore)] fn test_page_index_reader_all_type() {
2178 let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2179 let builder = ReadOptionsBuilder::new();
2180 let options = builder.with_page_index().build();
2182 let reader_result = SerializedFileReader::new_with_options(test_file, options);
2183 let reader = reader_result.unwrap();
2184
2185 let metadata = reader.metadata();
2187 assert_eq!(metadata.num_row_groups(), 1);
2188
2189 let page_index = metadata.page_index().unwrap();
2190 let row_group_offset_indexes = metadata.page_index_for_row_group(0);
2191
2192 let row_group_metadata = metadata.row_group(0);
2194
2195 let ci = page_index.column_index(0, 0).unwrap();
2197 assert!(!ci.is_sorted());
2198 assert!(matches!(
2199 ci.get_boundary_order(),
2200 Some(BoundaryOrder::UNORDERED)
2201 ));
2202 if let ColumnIndexMetaData::INT32(index) = ci {
2203 check_native_page_index(
2204 index,
2205 325,
2206 get_row_group_min_max_bytes(row_group_metadata, 0),
2207 BoundaryOrder::UNORDERED,
2208 );
2209 assert_eq!(
2210 row_group_offset_indexes.page_locations(0).unwrap().len(),
2211 325
2212 );
2213 } else {
2214 unreachable!()
2215 }
2216 let ci = page_index.column_index(0, 1).unwrap();
2218 assert!(ci.is_sorted());
2219 if let ColumnIndexMetaData::BOOLEAN(index) = ci {
2220 assert_eq!(index.num_pages(), 82);
2221 assert_eq!(
2222 row_group_offset_indexes.page_locations(1).unwrap().len(),
2223 82
2224 );
2225 } else {
2226 unreachable!()
2227 }
2228 let ci = page_index.column_index(0, 2).unwrap();
2230 assert!(ci.is_sorted());
2231 if let ColumnIndexMetaData::INT32(index) = ci {
2232 check_native_page_index(
2233 index,
2234 325,
2235 get_row_group_min_max_bytes(row_group_metadata, 2),
2236 BoundaryOrder::ASCENDING,
2237 );
2238 assert_eq!(
2239 row_group_offset_indexes.page_locations(2).unwrap().len(),
2240 325
2241 );
2242 } else {
2243 unreachable!()
2244 }
2245 let ci = page_index.column_index(0, 3).unwrap();
2247 assert!(ci.is_sorted());
2248 if let ColumnIndexMetaData::INT32(index) = ci {
2249 check_native_page_index(
2250 index,
2251 325,
2252 get_row_group_min_max_bytes(row_group_metadata, 3),
2253 BoundaryOrder::ASCENDING,
2254 );
2255 assert_eq!(
2256 row_group_offset_indexes.page_locations(3).unwrap().len(),
2257 325
2258 );
2259 } else {
2260 unreachable!()
2261 }
2262 let ci = page_index.column_index(0, 4).unwrap();
2264 assert!(ci.is_sorted());
2265 if let ColumnIndexMetaData::INT32(index) = ci {
2266 check_native_page_index(
2267 index,
2268 325,
2269 get_row_group_min_max_bytes(row_group_metadata, 4),
2270 BoundaryOrder::ASCENDING,
2271 );
2272 assert_eq!(
2273 row_group_offset_indexes.page_locations(4).unwrap().len(),
2274 325
2275 );
2276 } else {
2277 unreachable!()
2278 }
2279 let ci = page_index.column_index(0, 5).unwrap();
2281 assert!(!ci.is_sorted());
2282 if let ColumnIndexMetaData::INT64(index) = ci {
2283 check_native_page_index(
2284 index,
2285 528,
2286 get_row_group_min_max_bytes(row_group_metadata, 5),
2287 BoundaryOrder::UNORDERED,
2288 );
2289 assert_eq!(
2290 row_group_offset_indexes.page_locations(5).unwrap().len(),
2291 528
2292 );
2293 } else {
2294 unreachable!()
2295 }
2296 let ci = page_index.column_index(0, 6).unwrap();
2298 assert!(ci.is_sorted());
2299 if let ColumnIndexMetaData::FLOAT(index) = ci {
2300 check_native_page_index(
2301 index,
2302 325,
2303 get_row_group_min_max_bytes(row_group_metadata, 6),
2304 BoundaryOrder::ASCENDING,
2305 );
2306 assert_eq!(
2307 row_group_offset_indexes.page_locations(6).unwrap().len(),
2308 325
2309 );
2310 } else {
2311 unreachable!()
2312 }
2313 let ci = page_index.column_index(0, 7).unwrap();
2315 assert!(!ci.is_sorted());
2316 if let ColumnIndexMetaData::DOUBLE(index) = ci {
2317 check_native_page_index(
2318 index,
2319 528,
2320 get_row_group_min_max_bytes(row_group_metadata, 7),
2321 BoundaryOrder::UNORDERED,
2322 );
2323 assert_eq!(
2324 row_group_offset_indexes.page_locations(7).unwrap().len(),
2325 528
2326 );
2327 } else {
2328 unreachable!()
2329 }
2330 let ci = page_index.column_index(0, 8).unwrap();
2332 assert!(!ci.is_sorted());
2333 if let ColumnIndexMetaData::BYTE_ARRAY(index) = ci {
2334 check_byte_array_page_index(
2335 index,
2336 974,
2337 get_row_group_min_max_bytes(row_group_metadata, 8),
2338 BoundaryOrder::UNORDERED,
2339 );
2340 assert_eq!(
2341 row_group_offset_indexes.page_locations(8).unwrap().len(),
2342 974
2343 );
2344 } else {
2345 unreachable!()
2346 }
2347 let ci = page_index.column_index(0, 9).unwrap();
2349 assert!(ci.is_sorted());
2350 if let ColumnIndexMetaData::BYTE_ARRAY(index) = ci {
2351 check_byte_array_page_index(
2352 index,
2353 352,
2354 get_row_group_min_max_bytes(row_group_metadata, 9),
2355 BoundaryOrder::ASCENDING,
2356 );
2357 assert_eq!(
2358 row_group_offset_indexes.page_locations(9).unwrap().len(),
2359 352
2360 );
2361 } else {
2362 unreachable!()
2363 }
2364 assert!(page_index.column_index(0, 10).is_none());
2367 let ci = page_index.column_index(0, 11).unwrap();
2369 assert!(ci.is_sorted());
2370 if let ColumnIndexMetaData::INT32(index) = ci {
2371 check_native_page_index(
2372 index,
2373 325,
2374 get_row_group_min_max_bytes(row_group_metadata, 11),
2375 BoundaryOrder::ASCENDING,
2376 );
2377 assert_eq!(
2378 row_group_offset_indexes.page_locations(11).unwrap().len(),
2379 325
2380 );
2381 } else {
2382 unreachable!()
2383 }
2384 let ci = page_index.column_index(0, 12).unwrap();
2386 assert!(!ci.is_sorted());
2387 if let ColumnIndexMetaData::INT32(index) = ci {
2388 check_native_page_index(
2389 index,
2390 325,
2391 get_row_group_min_max_bytes(row_group_metadata, 12),
2392 BoundaryOrder::UNORDERED,
2393 );
2394 assert_eq!(
2395 row_group_offset_indexes.page_locations(12).unwrap().len(),
2396 325
2397 );
2398 } else {
2399 unreachable!()
2400 }
2401 }
2402
2403 fn check_native_page_index<T: ParquetValueType>(
2404 row_group_index: &PrimitiveColumnIndex<T>,
2405 page_size: usize,
2406 min_max: (&[u8], &[u8]),
2407 boundary_order: BoundaryOrder,
2408 ) {
2409 assert_eq!(row_group_index.num_pages() as usize, page_size);
2410 assert_eq!(row_group_index.boundary_order, boundary_order);
2411 assert!(row_group_index.min_values().iter().all(|x| {
2412 x >= &T::try_from_le_slice(min_max.0).unwrap()
2413 && x <= &T::try_from_le_slice(min_max.1).unwrap()
2414 }));
2415 }
2416
2417 fn check_byte_array_page_index(
2418 row_group_index: &ByteArrayColumnIndex,
2419 page_size: usize,
2420 min_max: (&[u8], &[u8]),
2421 boundary_order: BoundaryOrder,
2422 ) {
2423 assert_eq!(row_group_index.num_pages() as usize, page_size);
2424 assert_eq!(row_group_index.boundary_order, boundary_order);
2425 for i in 0..row_group_index.num_pages() as usize {
2426 let x = row_group_index.min_value(i).unwrap();
2427 assert!(x >= min_max.0 && x <= min_max.1);
2428 }
2429 }
2430
2431 fn get_row_group_min_max_bytes(r: &RowGroupMetaData, col_num: usize) -> (&[u8], &[u8]) {
2432 let statistics = r.column(col_num).statistics().unwrap();
2433 (
2434 statistics.min_bytes_opt().unwrap_or_default(),
2435 statistics.max_bytes_opt().unwrap_or_default(),
2436 )
2437 }
2438
2439 #[test]
2440 #[cfg_attr(miri, ignore)] fn test_skip_next_page_with_dictionary_page() {
2442 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2443 let builder = ReadOptionsBuilder::new();
2444 let options = builder.with_page_index().build();
2446 let reader_result = SerializedFileReader::new_with_options(test_file, options);
2447 let reader = reader_result.unwrap();
2448
2449 let row_group_reader = reader.get_row_group(0).unwrap();
2450
2451 let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2453
2454 let mut vec = vec![];
2455
2456 let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2458 assert!(meta.is_dict);
2459
2460 column_page_reader.skip_next_page().unwrap();
2462
2463 let page = column_page_reader.get_next_page().unwrap().unwrap();
2465 assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2466
2467 for _i in 0..351 {
2469 let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2471 assert!(!meta.is_dict); vec.push(meta);
2473
2474 let page = column_page_reader.get_next_page().unwrap().unwrap();
2475 assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2476 }
2477
2478 assert!(column_page_reader.peek_next_page().unwrap().is_none());
2480 assert!(column_page_reader.get_next_page().unwrap().is_none());
2481
2482 assert_eq!(vec.len(), 351);
2484 }
2485
2486 #[test]
2487 #[cfg_attr(miri, ignore)] fn test_skip_page_with_offset_index() {
2489 let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2490 let builder = ReadOptionsBuilder::new();
2491 let options = builder.with_page_index().build();
2493 let reader_result = SerializedFileReader::new_with_options(test_file, options);
2494 let reader = reader_result.unwrap();
2495
2496 let row_group_reader = reader.get_row_group(0).unwrap();
2497
2498 let mut column_page_reader = row_group_reader.get_column_page_reader(4).unwrap();
2500
2501 let mut vec = vec![];
2502
2503 for i in 0..325 {
2504 if i % 2 == 0 {
2505 vec.push(column_page_reader.get_next_page().unwrap().unwrap());
2506 } else {
2507 column_page_reader.skip_next_page().unwrap();
2508 }
2509 }
2510 assert!(column_page_reader.peek_next_page().unwrap().is_none());
2512 assert!(column_page_reader.get_next_page().unwrap().is_none());
2513
2514 assert_eq!(vec.len(), 163);
2515 }
2516
2517 #[test]
2518 fn test_skip_page_without_offset_index() {
2519 let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2520
2521 let reader_result = SerializedFileReader::new(test_file);
2523 let reader = reader_result.unwrap();
2524
2525 let row_group_reader = reader.get_row_group(0).unwrap();
2526
2527 let mut column_page_reader = row_group_reader.get_column_page_reader(4).unwrap();
2529
2530 let mut vec = vec![];
2531
2532 for i in 0..325 {
2533 if i % 2 == 0 {
2534 vec.push(column_page_reader.get_next_page().unwrap().unwrap());
2535 } else {
2536 column_page_reader.peek_next_page().unwrap().unwrap();
2537 column_page_reader.skip_next_page().unwrap();
2538 }
2539 }
2540 assert!(column_page_reader.peek_next_page().unwrap().is_none());
2542 assert!(column_page_reader.get_next_page().unwrap().is_none());
2543
2544 assert_eq!(vec.len(), 163);
2545 }
2546
2547 #[test]
2548 #[cfg_attr(miri, ignore)] fn test_peek_page_with_dictionary_page() {
2550 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2551 let builder = ReadOptionsBuilder::new();
2552 let options = builder.with_page_index().build();
2554 let reader_result = SerializedFileReader::new_with_options(test_file, options);
2555 let reader = reader_result.unwrap();
2556 let row_group_reader = reader.get_row_group(0).unwrap();
2557
2558 let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2560
2561 let mut vec = vec![];
2562
2563 let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2564 assert!(meta.is_dict);
2565 let page = column_page_reader.get_next_page().unwrap().unwrap();
2566 assert!(matches!(page.page_type(), basic::PageType::DICTIONARY_PAGE));
2567
2568 for i in 0..352 {
2569 let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2570 if i != 351 {
2573 assert!((meta.num_rows == Some(21)) || (meta.num_rows == Some(20)));
2574 } else {
2575 assert_eq!(meta.num_rows, Some(10));
2578 }
2579 assert!(!meta.is_dict);
2580 vec.push(meta);
2581 let page = column_page_reader.get_next_page().unwrap().unwrap();
2582 assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2583 }
2584
2585 assert!(column_page_reader.peek_next_page().unwrap().is_none());
2587 assert!(column_page_reader.get_next_page().unwrap().is_none());
2588
2589 assert_eq!(vec.len(), 352);
2590 }
2591
2592 #[test]
2593 fn test_peek_page_with_dictionary_page_without_offset_index() {
2594 let test_file = get_test_file("alltypes_tiny_pages.parquet");
2595
2596 let reader_result = SerializedFileReader::new(test_file);
2597 let reader = reader_result.unwrap();
2598 let row_group_reader = reader.get_row_group(0).unwrap();
2599
2600 let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2602
2603 let mut vec = vec![];
2604
2605 let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2606 assert!(meta.is_dict);
2607 let page = column_page_reader.get_next_page().unwrap().unwrap();
2608 assert!(matches!(page.page_type(), basic::PageType::DICTIONARY_PAGE));
2609
2610 for i in 0..352 {
2611 let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2612 if i != 351 {
2615 assert!((meta.num_levels == Some(21)) || (meta.num_levels == Some(20)));
2616 } else {
2617 assert_eq!(meta.num_levels, Some(10));
2620 }
2621 assert!(!meta.is_dict);
2622 vec.push(meta);
2623 let page = column_page_reader.get_next_page().unwrap().unwrap();
2624 assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2625 }
2626
2627 assert!(column_page_reader.peek_next_page().unwrap().is_none());
2629 assert!(column_page_reader.get_next_page().unwrap().is_none());
2630
2631 assert_eq!(vec.len(), 352);
2632 }
2633
2634 #[test]
2635 fn test_fixed_length_index() {
2636 let message_type = "
2637 message test_schema {
2638 OPTIONAL FIXED_LEN_BYTE_ARRAY (11) value (DECIMAL(25,2));
2639 }
2640 ";
2641
2642 let schema = parse_message_type(message_type).unwrap();
2643 let mut out = Vec::with_capacity(1024);
2644 let mut writer =
2645 SerializedFileWriter::new(&mut out, Arc::new(schema), Default::default()).unwrap();
2646
2647 let mut r = writer.next_row_group().unwrap();
2648 let mut c = r.next_column().unwrap().unwrap();
2649 c.typed::<FixedLenByteArrayType>()
2650 .write_batch(
2651 &[vec![0; 11].into(), vec![5; 11].into(), vec![3; 11].into()],
2652 Some(&[1, 1, 0, 1]),
2653 None,
2654 )
2655 .unwrap();
2656 c.close().unwrap();
2657 r.close().unwrap();
2658 writer.close().unwrap();
2659
2660 let b = Bytes::from(out);
2661 let options = ReadOptionsBuilder::new().with_page_index().build();
2662 let reader = SerializedFileReader::new_with_options(b, options).unwrap();
2663 let page_index = reader.metadata().page_index().unwrap();
2664
2665 match page_index.column_index(0, 0) {
2666 Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v)) => {
2667 assert_eq!(v.num_pages(), 1);
2668 assert_eq!(v.null_count(0).unwrap(), 1);
2669 assert_eq!(v.min_value(0).unwrap(), &[0; 11]);
2670 assert_eq!(v.max_value(0).unwrap(), &[5; 11]);
2671 }
2672 _ => unreachable!(),
2673 }
2674 }
2675
2676 #[test]
2677 fn test_multi_gz() {
2678 let file = get_test_file("concatenated_gzip_members.parquet");
2679 let reader = SerializedFileReader::new(file).unwrap();
2680 let row_group_reader = reader.get_row_group(0).unwrap();
2681 match row_group_reader.get_column_reader(0).unwrap() {
2682 ColumnReader::Int64ColumnReader(mut reader) => {
2683 let mut buffer = Vec::with_capacity(1024);
2684 let mut def_levels = Vec::with_capacity(1024);
2685 let (num_records, num_values, num_levels) = reader
2686 .read_records(1024, Some(&mut def_levels), None, &mut buffer)
2687 .unwrap();
2688
2689 assert_eq!(num_records, 513);
2690 assert_eq!(num_values, 513);
2691 assert_eq!(num_levels, 513);
2692
2693 let expected: Vec<i64> = (1..514).collect();
2694 assert_eq!(&buffer, &expected);
2695 }
2696 _ => unreachable!(),
2697 }
2698 }
2699
2700 #[test]
2701 #[cfg_attr(miri, ignore)] fn test_byte_stream_split_extended() {
2703 let path = format!(
2704 "{}/byte_stream_split_extended.gzip.parquet",
2705 arrow::util::test_util::parquet_test_data(),
2706 );
2707 let file = File::open(path).unwrap();
2708 let reader = Box::new(SerializedFileReader::new(file).expect("Failed to create reader"));
2709
2710 let mut iter = reader
2712 .get_row_iter(None)
2713 .expect("Failed to create row iterator");
2714
2715 let mut start = 0;
2716 let end = reader.metadata().file_metadata().num_rows();
2717
2718 let check_row = |row: Result<Row, ParquetError>| {
2719 assert!(row.is_ok());
2720 let r = row.unwrap();
2721 assert_eq!(r.get_float16(0).unwrap(), r.get_float16(1).unwrap());
2722 assert_eq!(r.get_float(2).unwrap(), r.get_float(3).unwrap());
2723 assert_eq!(r.get_double(4).unwrap(), r.get_double(5).unwrap());
2724 assert_eq!(r.get_int(6).unwrap(), r.get_int(7).unwrap());
2725 assert_eq!(r.get_long(8).unwrap(), r.get_long(9).unwrap());
2726 assert_eq!(r.get_bytes(10).unwrap(), r.get_bytes(11).unwrap());
2727 assert_eq!(r.get_decimal(12).unwrap(), r.get_decimal(13).unwrap());
2728 };
2729
2730 while start < end {
2731 match iter.next() {
2732 Some(row) => check_row(row),
2733 None => break,
2734 }
2735 start += 1;
2736 }
2737 }
2738
2739 #[test]
2740 fn test_filtered_rowgroup_metadata() {
2741 let message_type = "
2742 message test_schema {
2743 REQUIRED INT32 a;
2744 }
2745 ";
2746 let schema = Arc::new(parse_message_type(message_type).unwrap());
2747 let props = Arc::new(
2748 WriterProperties::builder()
2749 .set_statistics_enabled(EnabledStatistics::Page)
2750 .build(),
2751 );
2752 let mut file: File = tempfile::tempfile().unwrap();
2753 let mut file_writer = SerializedFileWriter::new(&mut file, schema, props).unwrap();
2754 let data = [1, 2, 3, 4, 5];
2755
2756 for idx in 0..5 {
2758 let data_i: Vec<i32> = data.iter().map(|x| x * (idx + 1)).collect();
2759 let mut row_group_writer = file_writer.next_row_group().unwrap();
2760 if let Some(mut writer) = row_group_writer.next_column().unwrap() {
2761 writer
2762 .typed::<Int32Type>()
2763 .write_batch(data_i.as_slice(), None, None)
2764 .unwrap();
2765 writer.close().unwrap();
2766 }
2767 row_group_writer.close().unwrap();
2768 file_writer.flushed_row_groups();
2769 }
2770 let file_metadata = file_writer.close().unwrap();
2771
2772 assert_eq!(file_metadata.file_metadata().num_rows(), 25);
2773 assert_eq!(file_metadata.num_row_groups(), 5);
2774
2775 let read_options = ReadOptionsBuilder::new()
2777 .with_page_index()
2778 .with_predicate(Box::new(|rgmeta, _| rgmeta.ordinal().unwrap_or(0) == 2))
2779 .build();
2780 let reader =
2781 SerializedFileReader::new_with_options(file.try_clone().unwrap(), read_options)
2782 .unwrap();
2783 let metadata = reader.metadata();
2784
2785 assert_eq!(metadata.num_row_groups(), 1);
2787 assert_eq!(metadata.row_group(0).ordinal(), Some(2));
2788
2789 assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
2791 let page_index = metadata.page_index().unwrap();
2792
2793 let col_stats = metadata.row_group(0).column(0).statistics().unwrap();
2794 let pg_idx = page_index.column_index(0, 0);
2795 let off_idx_i = page_index.offset_index(0, 0);
2796
2797 match pg_idx {
2799 Some(ColumnIndexMetaData::INT32(int_idx)) => {
2800 let min = col_stats.min_bytes_opt().unwrap().get_i32_le();
2801 let max = col_stats.max_bytes_opt().unwrap().get_i32_le();
2802 assert_eq!(int_idx.min_value(0), Some(min).as_ref());
2803 assert_eq!(int_idx.max_value(0), Some(max).as_ref());
2804 }
2805 _ => panic!("wrong stats type"),
2806 }
2807
2808 assert_eq!(
2810 off_idx_i.as_ref().unwrap().page_locations[0].offset,
2811 metadata.row_group(0).column(0).data_page_offset()
2812 );
2813
2814 let read_options = ReadOptionsBuilder::new()
2816 .with_page_index()
2817 .with_predicate(Box::new(|rgmeta, _| rgmeta.ordinal().unwrap_or(0) % 2 == 1))
2818 .build();
2819 let reader =
2820 SerializedFileReader::new_with_options(file.try_clone().unwrap(), read_options)
2821 .unwrap();
2822 let metadata = reader.metadata();
2823
2824 assert_eq!(metadata.num_row_groups(), 2);
2826 assert_eq!(metadata.row_group(0).ordinal(), Some(1));
2827 assert_eq!(metadata.row_group(1).ordinal(), Some(3));
2828
2829 assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
2831
2832 let page_index = metadata.page_index().unwrap();
2833
2834 for rg_idx in 0..metadata.num_row_groups() {
2835 let col_stats = metadata.row_group(rg_idx).column(0).statistics().unwrap();
2836 let pg_idx = page_index.column_index(rg_idx, 0);
2837 let off_idx_i = page_index.offset_index(rg_idx, 0);
2838
2839 match pg_idx {
2841 Some(ColumnIndexMetaData::INT32(int_idx)) => {
2842 let min = col_stats.min_bytes_opt().unwrap().get_i32_le();
2843 let max = col_stats.max_bytes_opt().unwrap().get_i32_le();
2844 assert_eq!(int_idx.min_value(0), Some(min).as_ref());
2845 assert_eq!(int_idx.max_value(0), Some(max).as_ref());
2846 }
2847 _ => panic!("wrong stats type"),
2848 }
2849
2850 assert_eq!(
2852 off_idx_i.as_ref().unwrap().page_locations[0].offset,
2853 metadata.row_group(rg_idx).column(0).data_page_offset()
2854 );
2855 }
2856 }
2857
2858 #[test]
2859 fn test_reuse_schema() {
2860 let file = get_test_file("alltypes_plain.parquet");
2861 let file_reader = SerializedFileReader::new(file.try_clone().unwrap()).unwrap();
2862 let schema = file_reader.metadata().file_metadata().schema_descr_ptr();
2863 let expected = file_reader.metadata;
2864
2865 let options = ReadOptionsBuilder::new()
2866 .with_parquet_schema(schema)
2867 .build();
2868 let file_reader = SerializedFileReader::new_with_options(file, options).unwrap();
2869
2870 assert_eq!(expected.as_ref(), file_reader.metadata.as_ref());
2871 assert!(Arc::ptr_eq(
2873 &expected.file_metadata().schema_descr_ptr(),
2874 &file_reader.metadata.file_metadata().schema_descr_ptr()
2875 ));
2876 }
2877
2878 #[test]
2879 fn test_read_unknown_logical_type() {
2880 let file = get_test_file("unknown-logical-type.parquet");
2881 let reader = SerializedFileReader::new(file).expect("Error opening file");
2882
2883 let schema = reader.metadata().file_metadata().schema_descr();
2884 assert_eq!(
2885 schema.column(0).logical_type_ref(),
2886 Some(&basic::LogicalType::String)
2887 );
2888 assert_eq!(
2889 schema.column(1).logical_type_ref(),
2890 Some(&basic::LogicalType::_Unknown { field_id: 2555 })
2891 );
2892 assert_eq!(schema.column(1).physical_type(), Type::BYTE_ARRAY);
2893
2894 let mut iter = reader
2895 .get_row_iter(None)
2896 .expect("Failed to create row iterator");
2897
2898 let mut num_rows = 0;
2899 while iter.next().is_some() {
2900 num_rows += 1;
2901 }
2902 assert_eq!(num_rows, reader.metadata().file_metadata().num_rows());
2903 }
2904}