1#[cfg(feature = "encryption")]
19use crate::encryption::decrypt::FileDecryptionProperties;
20use crate::errors::{ParquetError, Result};
21use crate::file::FOOTER_SIZE;
22use crate::file::metadata::parser::decode_metadata;
23use crate::file::metadata::thrift::parquet_schema_from_bytes;
24use crate::file::metadata::{
25 FooterTail, ParquetMetaData, ParquetMetaDataOptions, ParquetMetaDataPushDecoder,
26};
27use crate::file::reader::ChunkReader;
28use crate::schema::types::SchemaDescriptor;
29use bytes::Bytes;
30use std::sync::Arc;
31use std::{io::Read, ops::Range};
32
33use crate::DecodeResult;
34#[cfg(all(feature = "async", feature = "arrow"))]
35use crate::arrow::async_reader::{MetadataFetch, MetadataSuffixFetch};
36
37#[derive(Default, Debug)]
70pub struct ParquetMetaDataReader {
71 metadata: Option<ParquetMetaData>,
72 column_index: PageIndexPolicy,
73 offset_index: PageIndexPolicy,
74 prefetch_hint: Option<usize>,
75 metadata_options: Option<Arc<ParquetMetaDataOptions>>,
76 metadata_size: Option<usize>,
79 #[cfg(feature = "encryption")]
80 file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
81}
82
83#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
85pub enum PageIndexPolicy {
86 #[default]
88 Skip,
89 Optional,
91 Required,
93}
94
95impl From<bool> for PageIndexPolicy {
96 fn from(value: bool) -> Self {
97 match value {
98 true => Self::Required,
99 false => Self::Skip,
100 }
101 }
102}
103
104impl ParquetMetaDataReader {
105 pub fn new() -> Self {
107 Default::default()
108 }
109
110 pub fn new_with_metadata(metadata: ParquetMetaData) -> Self {
113 Self {
114 metadata: Some(metadata),
115 ..Default::default()
116 }
117 }
118
119 pub fn with_page_index_policy(self, policy: PageIndexPolicy) -> Self {
121 self.with_column_index_policy(policy)
122 .with_offset_index_policy(policy)
123 }
124
125 pub fn with_column_index_policy(mut self, policy: PageIndexPolicy) -> Self {
127 self.column_index = policy;
128 self
129 }
130
131 pub fn with_offset_index_policy(mut self, policy: PageIndexPolicy) -> Self {
133 self.offset_index = policy;
134 self
135 }
136
137 pub fn with_metadata_options(mut self, options: Option<ParquetMetaDataOptions>) -> Self {
139 self.metadata_options = options.map(Arc::new);
140 self
141 }
142
143 pub fn with_prefetch_hint(mut self, prefetch: Option<usize>) -> Self {
155 self.prefetch_hint = prefetch;
156 self
157 }
158
159 #[cfg(feature = "encryption")]
163 pub fn with_decryption_properties(
164 mut self,
165 properties: Option<std::sync::Arc<FileDecryptionProperties>>,
166 ) -> Self {
167 self.file_decryption_properties = properties;
168 self
169 }
170
171 pub fn has_metadata(&self) -> bool {
173 self.metadata.is_some()
174 }
175
176 pub fn finish(&mut self) -> Result<ParquetMetaData> {
178 self.metadata
179 .take()
180 .ok_or_else(|| general_err!("could not parse parquet metadata"))
181 }
182
183 pub fn parse_and_finish<R: ChunkReader>(mut self, reader: &R) -> Result<ParquetMetaData> {
202 self.try_parse(reader)?;
203 self.finish()
204 }
205
206 pub fn try_parse<R: ChunkReader>(&mut self, reader: &R) -> Result<()> {
212 self.try_parse_sized(reader, reader.len())
213 }
214
215 pub fn try_parse_sized<R: ChunkReader>(&mut self, reader: &R, file_size: u64) -> Result<()> {
288 self.metadata = match self.parse_metadata(reader) {
289 Ok(metadata) => Some(metadata),
290 Err(ParquetError::NeedMoreData(needed)) => {
291 return if file_size == reader.len() || needed as u64 > file_size {
294 Err(eof_err!(
295 "Parquet file too small. Size is {} but need {}",
296 file_size,
297 needed
298 ))
299 } else {
300 Err(ParquetError::NeedMoreData(needed))
302 };
303 }
304 Err(e) => return Err(e),
305 };
306
307 if self.column_index == PageIndexPolicy::Skip && self.offset_index == PageIndexPolicy::Skip
309 {
310 return Ok(());
311 }
312
313 self.read_page_indexes_sized(reader, file_size)
314 }
315
316 pub fn read_page_indexes<R: ChunkReader>(&mut self, reader: &R) -> Result<()> {
319 self.read_page_indexes_sized(reader, reader.len())
320 }
321
322 pub fn read_page_indexes_sized<R: ChunkReader>(
328 &mut self,
329 reader: &R,
330 file_size: u64,
331 ) -> Result<()> {
332 let Some(metadata) = self.metadata.take() else {
333 return Err(general_err!(
334 "Tried to read page indexes without ParquetMetaData metadata"
335 ));
336 };
337
338 let push_decoder = ParquetMetaDataPushDecoder::try_new_with_metadata(file_size, metadata)?
339 .with_offset_index_policy(self.offset_index)
340 .with_column_index_policy(self.column_index)
341 .with_metadata_options(self.metadata_options.clone());
342 let mut push_decoder = self.prepare_push_decoder(push_decoder);
343
344 let range = match needs_index_data(&mut push_decoder)? {
346 NeedsIndexData::No(metadata) => {
347 self.metadata = Some(metadata);
348 return Ok(());
349 }
350 NeedsIndexData::Yes(range) => range,
351 };
352
353 let file_range = file_size.saturating_sub(reader.len())..file_size;
356 if !(file_range.contains(&range.start) && file_range.contains(&range.end)) {
357 return if range.end > file_size {
359 Err(eof_err!(
360 "Parquet file too small. Range {range:?} is beyond file bounds {file_size}",
361 ))
362 } else {
363 Err(ParquetError::NeedMoreData(
365 (file_size - range.start).try_into()?,
366 ))
367 };
368 }
369
370 if let Some(metadata_size) = self.metadata_size {
373 let metadata_range = file_size.saturating_sub(metadata_size as u64)..file_size;
374 if range.end > metadata_range.start {
375 return Err(eof_err!(
376 "Parquet file too small. Page index range {range:?} overlaps with file metadata {metadata_range:?}",
377 ));
378 }
379 }
380
381 let bytes_needed = usize::try_from(range.end - range.start)?;
383 let bytes = reader.get_bytes(range.start - file_range.start, bytes_needed)?;
384
385 push_decoder.push_range(range, bytes)?;
386 let metadata = parse_index_data(&mut push_decoder)?;
387 self.metadata = Some(metadata);
388
389 Ok(())
390 }
391
392 #[cfg(all(feature = "async", feature = "arrow"))]
399 pub async fn load_and_finish<F: MetadataFetch>(
400 mut self,
401 fetch: F,
402 file_size: u64,
403 ) -> Result<ParquetMetaData> {
404 self.try_load(fetch, file_size).await?;
405 self.finish()
406 }
407
408 #[cfg(all(feature = "async", feature = "arrow"))]
415 pub async fn load_via_suffix_and_finish<F: MetadataSuffixFetch>(
416 mut self,
417 fetch: F,
418 ) -> Result<ParquetMetaData> {
419 self.try_load_via_suffix(fetch).await?;
420 self.finish()
421 }
422 #[cfg(all(feature = "async", feature = "arrow"))]
428 pub async fn try_load<F: MetadataFetch>(&mut self, mut fetch: F, file_size: u64) -> Result<()> {
429 let (metadata, remainder) = self.load_metadata(&mut fetch, file_size).await?;
430
431 self.metadata = Some(metadata);
432
433 if self.column_index == PageIndexPolicy::Skip && self.offset_index == PageIndexPolicy::Skip
435 {
436 return Ok(());
437 }
438
439 self.load_page_index_with_remainder(fetch, remainder).await
440 }
441
442 #[cfg(all(feature = "async", feature = "arrow"))]
448 pub async fn try_load_via_suffix<F: MetadataSuffixFetch>(
449 &mut self,
450 mut fetch: F,
451 ) -> Result<()> {
452 let (metadata, remainder) = self.load_metadata_via_suffix(&mut fetch).await?;
453
454 self.metadata = Some(metadata);
455
456 if self.column_index == PageIndexPolicy::Skip && self.offset_index == PageIndexPolicy::Skip
458 {
459 return Ok(());
460 }
461
462 self.load_page_index_with_remainder(fetch, remainder).await
463 }
464
465 #[cfg(all(feature = "async", feature = "arrow"))]
468 pub async fn load_page_index<F: MetadataFetch>(&mut self, fetch: F) -> Result<()> {
469 self.load_page_index_with_remainder(fetch, None).await
470 }
471
472 #[cfg(all(feature = "async", feature = "arrow"))]
473 async fn load_page_index_with_remainder<F: MetadataFetch>(
474 &mut self,
475 mut fetch: F,
476 remainder: Option<(usize, Bytes)>,
477 ) -> Result<()> {
478 let Some(metadata) = self.metadata.take() else {
479 return Err(general_err!("Footer metadata is not present"));
480 };
481
482 let file_size = u64::MAX;
485 let push_decoder = ParquetMetaDataPushDecoder::try_new_with_metadata(file_size, metadata)?
486 .with_offset_index_policy(self.offset_index)
487 .with_column_index_policy(self.column_index)
488 .with_metadata_options(self.metadata_options.clone());
489 let mut push_decoder = self.prepare_push_decoder(push_decoder);
490
491 let range = match needs_index_data(&mut push_decoder)? {
493 NeedsIndexData::No(metadata) => {
494 self.metadata = Some(metadata);
495 return Ok(());
496 }
497 NeedsIndexData::Yes(range) => range,
498 };
499
500 let bytes = match &remainder {
501 Some((remainder_start, remainder)) if *remainder_start as u64 <= range.start => {
502 let remainder_start = *remainder_start as u64;
503 let offset = usize::try_from(range.start - remainder_start)?;
504 let end = usize::try_from(range.end - remainder_start)?;
505 if end > remainder.len() {
506 return Err(general_err!(
507 "Corrupted parquet file: index data range ({:?}) exceeds remainder length ({})",
508 range,
509 remainder.len()
510 ));
511 }
512 remainder.slice(offset..end)
513 }
514 _ => fetch.fetch(range.clone()).await?,
516 };
517
518 if bytes.len() as u64 != range.end - range.start {
520 return Err(general_err!(
521 "Corrupted parquet file: index data length mismatch, expected {}, got {}",
522 range.end - range.start,
523 bytes.len()
524 ));
525 }
526 push_decoder.push_range(range.clone(), bytes)?;
527 let metadata = parse_index_data(&mut push_decoder)?;
528 self.metadata = Some(metadata);
529 Ok(())
530 }
531
532 fn parse_metadata<R: ChunkReader>(&mut self, chunk_reader: &R) -> Result<ParquetMetaData> {
535 let file_size = chunk_reader.len();
537 if file_size < (FOOTER_SIZE as u64) {
538 return Err(ParquetError::NeedMoreData(FOOTER_SIZE));
539 }
540
541 let mut footer = [0_u8; FOOTER_SIZE];
542 chunk_reader
543 .get_read(file_size - FOOTER_SIZE as u64)?
544 .read_exact(&mut footer)?;
545
546 let footer = FooterTail::try_new(&footer)?;
547 let metadata_len = footer.metadata_length();
548 let footer_metadata_len = FOOTER_SIZE + metadata_len;
549 self.metadata_size = Some(footer_metadata_len);
550
551 if footer_metadata_len as u64 > file_size {
552 return Err(ParquetError::NeedMoreData(footer_metadata_len));
553 }
554
555 let start = file_size - footer_metadata_len as u64;
556 let bytes = chunk_reader.get_bytes(start, metadata_len)?;
557 self.decode_footer_metadata(bytes, file_size, footer)
558 }
559
560 pub fn metadata_size(&self) -> Option<usize> {
563 self.metadata_size
564 }
565
566 #[cfg(all(feature = "async", feature = "arrow"))]
570 fn get_prefetch_size(&self) -> usize {
571 if let Some(prefetch) = self.prefetch_hint
572 && prefetch > FOOTER_SIZE
573 {
574 return prefetch;
575 }
576 FOOTER_SIZE
577 }
578
579 #[cfg(all(feature = "async", feature = "arrow"))]
580 async fn load_metadata<F: MetadataFetch>(
581 &self,
582 fetch: &mut F,
583 file_size: u64,
584 ) -> Result<(ParquetMetaData, Option<(usize, Bytes)>)> {
585 let prefetch = self.get_prefetch_size() as u64;
586
587 if file_size < FOOTER_SIZE as u64 {
588 return Err(eof_err!("file size of {} is less than footer", file_size));
589 }
590
591 let footer_start = file_size.saturating_sub(prefetch);
595
596 let suffix = fetch.fetch(footer_start..file_size).await?;
597 let suffix_len = suffix.len();
598 let fetch_len = (file_size - footer_start)
599 .try_into()
600 .expect("footer size should never be larger than u32");
601 if suffix_len < fetch_len {
602 return Err(eof_err!(
603 "metadata requires {} bytes, but could only read {}",
604 fetch_len,
605 suffix_len
606 ));
607 }
608
609 let mut footer = [0; FOOTER_SIZE];
610 footer.copy_from_slice(&suffix[suffix_len - FOOTER_SIZE..suffix_len]);
611
612 let footer = FooterTail::try_new(&footer)?;
613 let length = footer.metadata_length();
614
615 if file_size < (length + FOOTER_SIZE) as u64 {
616 return Err(eof_err!(
617 "file size of {} is less than footer + metadata {}",
618 file_size,
619 length + FOOTER_SIZE
620 ));
621 }
622
623 if length > suffix_len - FOOTER_SIZE {
625 let metadata_start = file_size - (length + FOOTER_SIZE) as u64;
626 let meta = fetch
627 .fetch(metadata_start..(file_size - FOOTER_SIZE as u64))
628 .await?;
629 Ok((self.decode_footer_metadata(meta, file_size, footer)?, None))
630 } else {
631 let metadata_start = (file_size - (length + FOOTER_SIZE) as u64 - footer_start)
632 .try_into()
633 .expect("metadata length should never be larger than u32");
634 let slice = suffix.slice(metadata_start..suffix_len - FOOTER_SIZE);
635 Ok((
636 self.decode_footer_metadata(slice, file_size, footer)?,
637 Some((footer_start as usize, suffix.slice(..metadata_start))),
638 ))
639 }
640 }
641
642 #[cfg(all(feature = "async", feature = "arrow"))]
643 async fn load_metadata_via_suffix<F: MetadataSuffixFetch>(
644 &self,
645 fetch: &mut F,
646 ) -> Result<(ParquetMetaData, Option<(usize, Bytes)>)> {
647 let prefetch = self.get_prefetch_size();
648
649 let suffix = fetch.fetch_suffix(prefetch).await?;
650 let suffix_len = suffix.len();
651
652 if suffix_len < FOOTER_SIZE {
653 return Err(eof_err!(
654 "footer metadata requires {} bytes, but could only read {}",
655 FOOTER_SIZE,
656 suffix_len
657 ));
658 }
659
660 let mut footer = [0; FOOTER_SIZE];
661 footer.copy_from_slice(&suffix[suffix_len - FOOTER_SIZE..suffix_len]);
662
663 let footer = FooterTail::try_new(&footer)?;
664 let length = footer.metadata_length();
665 let file_size = (length + FOOTER_SIZE) as u64;
668
669 let metadata_offset = length + FOOTER_SIZE;
671 if length > suffix_len - FOOTER_SIZE {
672 let meta = fetch.fetch_suffix(metadata_offset).await?;
673
674 if meta.len() < metadata_offset {
675 return Err(eof_err!(
676 "metadata requires {} bytes, but could only read {}",
677 metadata_offset,
678 meta.len()
679 ));
680 }
681
682 let meta = meta.slice(0..length);
684 Ok((self.decode_footer_metadata(meta, file_size, footer)?, None))
685 } else {
686 let metadata_start = suffix_len - metadata_offset;
687 let slice = suffix.slice(metadata_start..suffix_len - FOOTER_SIZE);
688 Ok((
689 self.decode_footer_metadata(slice, file_size, footer)?,
690 Some((0, suffix.slice(..metadata_start))),
691 ))
692 }
693 }
694
695 pub(crate) fn decode_footer_metadata(
709 &self,
710 buf: Bytes,
711 file_size: u64,
712 footer_tail: FooterTail,
713 ) -> Result<ParquetMetaData> {
714 let ending_offset = file_size.checked_sub(FOOTER_SIZE as u64).ok_or_else(|| {
719 general_err!(
720 "file size {file_size} is smaller than footer size {}",
721 FOOTER_SIZE
722 )
723 })?;
724
725 let starting_offset = ending_offset.checked_sub(buf.len() as u64).ok_or_else(|| {
726 general_err!(
727 "file size {file_size} is smaller than buffer size {} + footer size {}",
728 buf.len(),
729 FOOTER_SIZE
730 )
731 })?;
732
733 let range = starting_offset..ending_offset;
734
735 let push_decoder =
736 ParquetMetaDataPushDecoder::try_new_with_footer_tail(file_size, footer_tail)?
737 .with_page_index_policy(PageIndexPolicy::Skip)
739 .with_metadata_options(self.metadata_options.clone());
740
741 let mut push_decoder = self.prepare_push_decoder(push_decoder);
742 push_decoder.push_range(range, buf)?;
743 match push_decoder.try_decode()? {
744 DecodeResult::Data(metadata) => Ok(metadata),
745 DecodeResult::Finished => Err(general_err!(
746 "could not parse parquet metadata -- previously finished"
747 )),
748 DecodeResult::NeedsData(ranges) => Err(general_err!(
749 "could not parse parquet metadata, needs ranges {:?}",
750 ranges
751 )),
752 }
753 }
754
755 #[cfg(feature = "encryption")]
757 fn prepare_push_decoder(
758 &self,
759 push_decoder: ParquetMetaDataPushDecoder,
760 ) -> ParquetMetaDataPushDecoder {
761 push_decoder.with_file_decryption_properties(
762 self.file_decryption_properties
763 .as_ref()
764 .map(std::sync::Arc::clone),
765 )
766 }
767 #[cfg(not(feature = "encryption"))]
768 fn prepare_push_decoder(
769 &self,
770 push_decoder: ParquetMetaDataPushDecoder,
771 ) -> ParquetMetaDataPushDecoder {
772 push_decoder
773 }
774
775 pub fn decode_metadata(buf: &[u8]) -> Result<ParquetMetaData> {
783 decode_metadata(buf, None)
784 }
785
786 pub fn decode_metadata_with_options(
791 buf: &[u8],
792 options: Option<&ParquetMetaDataOptions>,
793 ) -> Result<ParquetMetaData> {
794 decode_metadata(buf, options)
795 }
796
797 pub fn decode_schema(buf: &[u8]) -> Result<Arc<SchemaDescriptor>> {
800 Ok(Arc::new(parquet_schema_from_bytes(buf)?))
801 }
802}
803
804enum NeedsIndexData {
807 No(ParquetMetaData),
809 Yes(Range<u64>),
811}
812
813fn needs_index_data(push_decoder: &mut ParquetMetaDataPushDecoder) -> Result<NeedsIndexData> {
816 match push_decoder.try_decode()? {
817 DecodeResult::NeedsData(ranges) => {
818 let range = ranges
819 .into_iter()
820 .reduce(|a, b| a.start.min(b.start)..a.end.max(b.end))
821 .ok_or_else(|| general_err!("Internal error: no ranges provided"))?;
822 Ok(NeedsIndexData::Yes(range))
823 }
824 DecodeResult::Data(metadata) => Ok(NeedsIndexData::No(metadata)),
825 DecodeResult::Finished => Err(general_err!("Internal error: decoder was finished")),
826 }
827}
828
829fn parse_index_data(push_decoder: &mut ParquetMetaDataPushDecoder) -> Result<ParquetMetaData> {
832 match push_decoder.try_decode()? {
833 DecodeResult::NeedsData(_) => Err(general_err!(
834 "Internal error: decoder still needs data after reading required range"
835 )),
836 DecodeResult::Data(metadata) => Ok(metadata),
837 DecodeResult::Finished => Err(general_err!("Internal error: decoder was finished")),
838 }
839}
840
841#[cfg(test)]
842mod tests {
843 use super::*;
844 use crate::file::metadata::PageIndex;
845 use crate::file::reader::Length;
846 use crate::util::test_common::file_util::get_test_file;
847 use std::ops::Range;
848
849 #[test]
850 fn test_parse_metadata_size_smaller_than_footer() {
851 let test_file = tempfile::tempfile().unwrap();
852 let err = ParquetMetaDataReader::new()
853 .parse_metadata(&test_file)
854 .unwrap_err();
855 assert!(matches!(err, ParquetError::NeedMoreData(FOOTER_SIZE)));
856 }
857
858 #[test]
859 fn test_parse_metadata_corrupt_footer() {
860 let data = Bytes::from(vec![1, 2, 3, 4, 5, 6, 7, 8]);
861 let reader_result = ParquetMetaDataReader::new().parse_metadata(&data);
862 assert_eq!(
863 reader_result.unwrap_err().to_string(),
864 "Parquet error: Invalid Parquet file. Corrupt footer"
865 );
866 }
867
868 #[test]
869 fn test_parse_metadata_invalid_start() {
870 let test_file = Bytes::from(vec![255, 0, 0, 0, b'P', b'A', b'R', b'1']);
871 let err = ParquetMetaDataReader::new()
872 .parse_metadata(&test_file)
873 .unwrap_err();
874 assert!(matches!(err, ParquetError::NeedMoreData(263)));
875 }
876
877 #[test]
878 #[cfg_attr(miri, ignore)] fn test_try_parse() {
880 let file = get_test_file("alltypes_tiny_pages.parquet");
881 let len = file.len();
882
883 let mut reader =
884 ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
885
886 let bytes_for_range = |range: Range<u64>| {
887 file.get_bytes(range.start, (range.end - range.start).try_into().unwrap())
888 .unwrap()
889 };
890
891 let bytes = bytes_for_range(0..len);
893 reader.try_parse(&bytes).unwrap();
894 let metadata = reader.finish().unwrap();
895 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
896
897 let bytes = bytes_for_range(320000..len);
899 reader.try_parse_sized(&bytes, len).unwrap();
900 let metadata = reader.finish().unwrap();
901 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
902
903 let bytes = bytes_for_range(323583..len);
905 reader.try_parse_sized(&bytes, len).unwrap();
906 let metadata = reader.finish().unwrap();
907 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
908
909 let bytes = bytes_for_range(323584..len);
911 match reader.try_parse_sized(&bytes, len).unwrap_err() {
913 ParquetError::NeedMoreData(needed) => {
915 let bytes = bytes_for_range(len - needed as u64..len);
916 reader.try_parse_sized(&bytes, len).unwrap();
917 let metadata = reader.finish().unwrap();
918 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
919 }
920 _ => panic!("unexpected error"),
921 }
922
923 let mut reader =
925 ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
926 let mut bytes = bytes_for_range(452505..len);
927 loop {
928 match reader.try_parse_sized(&bytes, len) {
929 Ok(()) => break,
930 Err(ParquetError::NeedMoreData(needed)) => {
931 bytes = bytes_for_range(len - needed as u64..len);
932 if reader.has_metadata() {
933 reader.read_page_indexes_sized(&bytes, len).unwrap();
934 break;
935 }
936 }
937 _ => panic!("unexpected error"),
938 }
939 }
940 let metadata = reader.finish().unwrap();
941 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
942
943 let bytes = bytes_for_range(323584..len);
945 let reader_result = reader.try_parse_sized(&bytes, len - 323584).unwrap_err();
946 assert_eq!(
947 reader_result.to_string(),
948 "EOF: Parquet file too small. Range 323583..452504 is beyond file bounds 130649"
949 );
950
951 let mut reader = ParquetMetaDataReader::new();
953 let bytes = bytes_for_range(452505..len);
954 match reader.try_parse_sized(&bytes, len).unwrap_err() {
956 ParquetError::NeedMoreData(needed) => {
958 let bytes = bytes_for_range(len - needed as u64..len);
959 reader.try_parse_sized(&bytes, len).unwrap();
960 reader.finish().unwrap();
961 }
962 _ => panic!("unexpected error"),
963 }
964
965 let reader_result = reader.try_parse(&bytes).unwrap_err();
967 assert_eq!(
968 reader_result.to_string(),
969 "EOF: Parquet file too small. Size is 1728 but need 1729"
970 );
971
972 let bytes = bytes_for_range(0..1000);
974 let reader_result = reader.try_parse_sized(&bytes, len).unwrap_err();
975 assert_eq!(
976 reader_result.to_string(),
977 "Parquet error: Invalid Parquet file. Corrupt footer"
978 );
979
980 let bytes = bytes_for_range(452510..len);
982 let reader_result = reader.try_parse_sized(&bytes, len - 452505).unwrap_err();
983 assert_eq!(
984 reader_result.to_string(),
985 "EOF: Parquet file too small. Size is 1728 but need 1729"
986 );
987 }
988}
989
990#[cfg(all(feature = "async", feature = "arrow", test))]
991mod async_tests {
992 use super::*;
993
994 use arrow::{array::Int32Array, datatypes::DataType};
995 use arrow_array::RecordBatch;
996 use arrow_schema::{Field, Schema};
997 use bytes::Bytes;
998 use futures::FutureExt;
999 use futures::future::BoxFuture;
1000 use std::fs::File;
1001 use std::future::Future;
1002 use std::io::{Read, Seek, SeekFrom};
1003 use std::ops::Range;
1004 use std::sync::Arc;
1005 use std::sync::atomic::{AtomicUsize, Ordering};
1006 use tempfile::NamedTempFile;
1007
1008 use crate::arrow::ArrowWriter;
1009 use crate::file::metadata::PageIndex;
1010 use crate::file::properties::WriterProperties;
1011 use crate::file::reader::Length;
1012 use crate::util::test_common::file_util::get_test_file;
1013
1014 struct MetadataFetchFn<F>(F);
1015
1016 impl<F, Fut> MetadataFetch for MetadataFetchFn<F>
1017 where
1018 F: FnMut(Range<u64>) -> Fut + Send,
1019 Fut: Future<Output = Result<Bytes>> + Send,
1020 {
1021 fn fetch(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1022 async move { self.0(range).await }.boxed()
1023 }
1024 }
1025
1026 struct MetadataSuffixFetchFn<F1, F2>(F1, F2);
1027
1028 impl<F1, Fut, F2> MetadataFetch for MetadataSuffixFetchFn<F1, F2>
1029 where
1030 F1: FnMut(Range<u64>) -> Fut + Send,
1031 Fut: Future<Output = Result<Bytes>> + Send,
1032 F2: Send,
1033 {
1034 fn fetch(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1035 async move { self.0(range).await }.boxed()
1036 }
1037 }
1038
1039 impl<F1, Fut, F2> MetadataSuffixFetch for MetadataSuffixFetchFn<F1, F2>
1040 where
1041 F1: FnMut(Range<u64>) -> Fut + Send,
1042 F2: FnMut(usize) -> Fut + Send,
1043 Fut: Future<Output = Result<Bytes>> + Send,
1044 {
1045 fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
1046 async move { self.1(suffix).await }.boxed()
1047 }
1048 }
1049
1050 fn read_range(file: &mut File, range: Range<u64>) -> Result<Bytes> {
1051 file.seek(SeekFrom::Start(range.start))?;
1052 let len = range.end - range.start;
1053 let mut buf = Vec::with_capacity(len.try_into().unwrap());
1054 file.take(len).read_to_end(&mut buf)?;
1055 Ok(buf.into())
1056 }
1057
1058 fn read_suffix(file: &mut File, suffix: usize) -> Result<Bytes> {
1059 let file_len = file.len();
1060 file.seek(SeekFrom::End(0 - suffix.min(file_len as _) as i64))?;
1062 let mut buf = Vec::with_capacity(suffix);
1063 file.take(suffix as _).read_to_end(&mut buf)?;
1064 Ok(buf.into())
1065 }
1066
1067 #[tokio::test]
1068 async fn test_simple() {
1069 let mut file = get_test_file("nulls.snappy.parquet");
1070 let len = file.len();
1071
1072 let expected = ParquetMetaDataReader::new()
1073 .parse_and_finish(&file)
1074 .unwrap();
1075 let expected = expected.file_metadata().schema();
1076 let fetch_count = AtomicUsize::new(0);
1077
1078 let mut fetch = |range| {
1079 fetch_count.fetch_add(1, Ordering::SeqCst);
1080 futures::future::ready(read_range(&mut file, range))
1081 };
1082
1083 let input = MetadataFetchFn(&mut fetch);
1084 let actual = ParquetMetaDataReader::new()
1085 .load_and_finish(input, len)
1086 .await
1087 .unwrap();
1088 assert_eq!(actual.file_metadata().schema(), expected);
1089 assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1090
1091 fetch_count.store(0, Ordering::SeqCst);
1093 let input = MetadataFetchFn(&mut fetch);
1094 let actual = ParquetMetaDataReader::new()
1095 .with_prefetch_hint(Some(7))
1096 .load_and_finish(input, len)
1097 .await
1098 .unwrap();
1099 assert_eq!(actual.file_metadata().schema(), expected);
1100 assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1101
1102 fetch_count.store(0, Ordering::SeqCst);
1104 let input = MetadataFetchFn(&mut fetch);
1105 let actual = ParquetMetaDataReader::new()
1106 .with_prefetch_hint(Some(10))
1107 .load_and_finish(input, len)
1108 .await
1109 .unwrap();
1110 assert_eq!(actual.file_metadata().schema(), expected);
1111 assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1112
1113 fetch_count.store(0, Ordering::SeqCst);
1115 let input = MetadataFetchFn(&mut fetch);
1116 let actual = ParquetMetaDataReader::new()
1117 .with_prefetch_hint(Some(500))
1118 .load_and_finish(input, len)
1119 .await
1120 .unwrap();
1121 assert_eq!(actual.file_metadata().schema(), expected);
1122 assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1123
1124 fetch_count.store(0, Ordering::SeqCst);
1126 let input = MetadataFetchFn(&mut fetch);
1127 let actual = ParquetMetaDataReader::new()
1128 .with_prefetch_hint(Some(428))
1129 .load_and_finish(input, len)
1130 .await
1131 .unwrap();
1132 assert_eq!(actual.file_metadata().schema(), expected);
1133 assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1134
1135 let input = MetadataFetchFn(&mut fetch);
1136 let err = ParquetMetaDataReader::new()
1137 .load_and_finish(input, 4)
1138 .await
1139 .unwrap_err()
1140 .to_string();
1141 assert_eq!(err, "EOF: file size of 4 is less than footer");
1142
1143 let input = MetadataFetchFn(&mut fetch);
1144 let err = ParquetMetaDataReader::new()
1145 .load_and_finish(input, 20)
1146 .await
1147 .unwrap_err()
1148 .to_string();
1149 assert_eq!(err, "Parquet error: Invalid Parquet file. Corrupt footer");
1150 }
1151
1152 #[tokio::test]
1153 async fn test_suffix() {
1154 let mut file = get_test_file("nulls.snappy.parquet");
1155 let mut file2 = file.try_clone().unwrap();
1156
1157 let expected = ParquetMetaDataReader::new()
1158 .parse_and_finish(&file)
1159 .unwrap();
1160 let expected = expected.file_metadata().schema();
1161 let fetch_count = AtomicUsize::new(0);
1162 let suffix_fetch_count = AtomicUsize::new(0);
1163
1164 let mut fetch = |range| {
1165 fetch_count.fetch_add(1, Ordering::SeqCst);
1166 futures::future::ready(read_range(&mut file, range))
1167 };
1168 let mut suffix_fetch = |suffix| {
1169 suffix_fetch_count.fetch_add(1, Ordering::SeqCst);
1170 futures::future::ready(read_suffix(&mut file2, suffix))
1171 };
1172
1173 let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1174 let actual = ParquetMetaDataReader::new()
1175 .load_via_suffix_and_finish(input)
1176 .await
1177 .unwrap();
1178 assert_eq!(actual.file_metadata().schema(), expected);
1179 assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1180 assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 2);
1181
1182 fetch_count.store(0, Ordering::SeqCst);
1184 suffix_fetch_count.store(0, Ordering::SeqCst);
1185 let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1186 let actual = ParquetMetaDataReader::new()
1187 .with_prefetch_hint(Some(7))
1188 .load_via_suffix_and_finish(input)
1189 .await
1190 .unwrap();
1191 assert_eq!(actual.file_metadata().schema(), expected);
1192 assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1193 assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 2);
1194
1195 fetch_count.store(0, Ordering::SeqCst);
1197 suffix_fetch_count.store(0, Ordering::SeqCst);
1198 let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1199 let actual = ParquetMetaDataReader::new()
1200 .with_prefetch_hint(Some(10))
1201 .load_via_suffix_and_finish(input)
1202 .await
1203 .unwrap();
1204 assert_eq!(actual.file_metadata().schema(), expected);
1205 assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1206 assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 2);
1207
1208 fetch_count.store(0, Ordering::SeqCst);
1210 suffix_fetch_count.store(0, Ordering::SeqCst);
1211 let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1212 let actual = ParquetMetaDataReader::new()
1213 .with_prefetch_hint(Some(500))
1214 .load_via_suffix_and_finish(input)
1215 .await
1216 .unwrap();
1217 assert_eq!(actual.file_metadata().schema(), expected);
1218 assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1219 assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 1);
1220
1221 fetch_count.store(0, Ordering::SeqCst);
1223 suffix_fetch_count.store(0, Ordering::SeqCst);
1224 let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1225 let actual = ParquetMetaDataReader::new()
1226 .with_prefetch_hint(Some(428))
1227 .load_via_suffix_and_finish(input)
1228 .await
1229 .unwrap();
1230 assert_eq!(actual.file_metadata().schema(), expected);
1231 assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1232 assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 1);
1233 }
1234
1235 #[cfg(feature = "encryption")]
1236 #[tokio::test]
1237 async fn test_suffix_with_encryption() {
1238 let mut file = get_test_file("uniform_encryption.parquet.encrypted");
1239 let mut file2 = file.try_clone().unwrap();
1240
1241 let mut fetch = |range| futures::future::ready(read_range(&mut file, range));
1242 let mut suffix_fetch = |suffix| futures::future::ready(read_suffix(&mut file2, suffix));
1243
1244 let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1245
1246 let key_code: &[u8] = b"0123456789012345";
1247 let decryption_properties = FileDecryptionProperties::builder(key_code.to_vec())
1248 .build()
1249 .unwrap();
1250
1251 let expected = ParquetMetaDataReader::new()
1253 .with_decryption_properties(Some(decryption_properties))
1254 .load_via_suffix_and_finish(input)
1255 .await
1256 .unwrap();
1257 assert_eq!(expected.num_row_groups(), 1);
1258 }
1259
1260 #[tokio::test]
1261 async fn test_page_index() {
1262 let mut file = get_test_file("alltypes_tiny_pages.parquet");
1263 let len = file.len();
1264 let fetch_count = AtomicUsize::new(0);
1265 let mut fetch = |range| {
1266 fetch_count.fetch_add(1, Ordering::SeqCst);
1267 futures::future::ready(read_range(&mut file, range))
1268 };
1269
1270 let f = MetadataFetchFn(&mut fetch);
1271 let mut loader =
1272 ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
1273 loader.try_load(f, len).await.unwrap();
1274 assert_eq!(fetch_count.load(Ordering::SeqCst), 3);
1275 let metadata = loader.finish().unwrap();
1276 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1277
1278 fetch_count.store(0, Ordering::SeqCst);
1280 let f = MetadataFetchFn(&mut fetch);
1281 let mut loader = ParquetMetaDataReader::new()
1282 .with_page_index_policy(PageIndexPolicy::Required)
1283 .with_prefetch_hint(Some(1729));
1284 loader.try_load(f, len).await.unwrap();
1285 assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1286 let metadata = loader.finish().unwrap();
1287 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1288
1289 fetch_count.store(0, Ordering::SeqCst);
1291 let f = MetadataFetchFn(&mut fetch);
1292 let mut loader = ParquetMetaDataReader::new()
1293 .with_page_index_policy(PageIndexPolicy::Required)
1294 .with_prefetch_hint(Some(130649));
1295 loader.try_load(f, len).await.unwrap();
1296 assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1297 let metadata = loader.finish().unwrap();
1298 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1299
1300 fetch_count.store(0, Ordering::SeqCst);
1302 let f = MetadataFetchFn(&mut fetch);
1303 let metadata = ParquetMetaDataReader::new()
1304 .with_page_index_policy(PageIndexPolicy::Required)
1305 .with_prefetch_hint(Some(130650))
1306 .load_and_finish(f, len)
1307 .await
1308 .unwrap();
1309 assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1310 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1311
1312 fetch_count.store(0, Ordering::SeqCst);
1314 let f = MetadataFetchFn(&mut fetch);
1315 let metadata = ParquetMetaDataReader::new()
1316 .with_page_index_policy(PageIndexPolicy::Required)
1317 .with_prefetch_hint(Some((len - 1000) as usize)) .load_and_finish(f, len)
1319 .await
1320 .unwrap();
1321 assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1322 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1323
1324 fetch_count.store(0, Ordering::SeqCst);
1326 let f = MetadataFetchFn(&mut fetch);
1327 let metadata = ParquetMetaDataReader::new()
1328 .with_page_index_policy(PageIndexPolicy::Required)
1329 .with_prefetch_hint(Some(len as usize)) .load_and_finish(f, len)
1331 .await
1332 .unwrap();
1333 assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1334 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1335
1336 fetch_count.store(0, Ordering::SeqCst);
1338 let f = MetadataFetchFn(&mut fetch);
1339 let metadata = ParquetMetaDataReader::new()
1340 .with_page_index_policy(PageIndexPolicy::Required)
1341 .with_prefetch_hint(Some((len + 1000) as usize)) .load_and_finish(f, len)
1343 .await
1344 .unwrap();
1345 assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1346 assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
1347 }
1348
1349 fn write_parquet_file(offset_index_disabled: bool) -> Result<NamedTempFile> {
1350 let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
1351 let batch = RecordBatch::try_new(
1352 schema.clone(),
1353 vec![Arc::new(Int32Array::from(vec![1, 2, 3]))],
1354 )?;
1355
1356 let file = NamedTempFile::new().unwrap();
1357
1358 let props = WriterProperties::builder()
1360 .set_offset_index_disabled(offset_index_disabled)
1361 .build();
1362
1363 let mut writer = ArrowWriter::try_new(file.reopen()?, schema, Some(props))?;
1364 writer.write(&batch)?;
1365 writer.close()?;
1366
1367 Ok(file)
1368 }
1369
1370 fn read_and_check(file: &File, policy: PageIndexPolicy) -> Result<ParquetMetaData> {
1371 let mut reader = ParquetMetaDataReader::new().with_page_index_policy(policy);
1372 reader.try_parse(file)?;
1373 reader.finish()
1374 }
1375
1376 #[test]
1377 fn test_page_index_policy() {
1378 let f = write_parquet_file(false).unwrap();
1380 read_and_check(f.as_file(), PageIndexPolicy::Required).unwrap();
1381 read_and_check(f.as_file(), PageIndexPolicy::Optional).unwrap();
1382 read_and_check(f.as_file(), PageIndexPolicy::Skip).unwrap();
1383
1384 let f = write_parquet_file(true).unwrap();
1386 let res = read_and_check(f.as_file(), PageIndexPolicy::Required);
1387 assert!(matches!(
1388 res,
1389 Err(ParquetError::General(e)) if e == "missing offset index"
1390 ));
1391 read_and_check(f.as_file(), PageIndexPolicy::Optional).unwrap();
1392 read_and_check(f.as_file(), PageIndexPolicy::Skip).unwrap();
1393 }
1394}