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