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