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