1#![cfg_attr(feature = "encryption", doc = "```rust")]
136#![cfg_attr(not(feature = "encryption"), doc = "```ignore")]
137#[cfg(feature = "experimental")]
185#[doc(hidden)]
186pub mod array_reader;
187#[cfg(not(feature = "experimental"))]
188mod array_reader;
189pub(crate) use array_reader::ByteArrayDecoderPlain;
192pub mod arrow_reader;
193pub mod arrow_writer;
194mod buffer;
195pub(crate) use buffer::offset_buffer::OffsetBuffer;
196mod decoder;
197
198#[cfg(feature = "async")]
199pub mod async_reader;
200#[cfg(feature = "async")]
201pub mod async_writer;
202
203pub mod push_decoder;
204
205mod in_memory_row_group;
206mod record_reader;
207
208#[cfg(feature = "experimental")]
209#[doc(hidden)]
210pub mod schema;
211#[cfg(not(feature = "experimental"))]
212mod schema;
213
214use std::fmt::Debug;
215
216pub use self::arrow_writer::ArrowWriter;
217#[cfg(feature = "async")]
218pub use self::async_reader::ParquetRecordBatchStreamBuilder;
219#[cfg(feature = "async")]
220pub use self::async_writer::AsyncArrowWriter;
221use crate::schema::types::SchemaDescriptor;
222use arrow_schema::{FieldRef, Schema};
223
224pub use self::schema::{
225 ArrowSchemaConverter, FieldLevels, add_encoded_arrow_schema_to_metadata, encode_arrow_schema,
226 parquet_to_arrow_field_levels, parquet_to_arrow_field_levels_with_virtual,
227 parquet_to_arrow_schema, parquet_to_arrow_schema_by_columns, virtual_type::*,
228};
229
230pub const ARROW_SCHEMA_META_KEY: &str = "ARROW:schema";
235
236pub const PARQUET_FIELD_ID_META_KEY: &str = "PARQUET:field_id";
242
243#[derive(Debug, Clone, PartialEq, Eq)]
269pub struct ProjectionMask {
270 mask: Option<Vec<bool>>,
290}
291
292impl ProjectionMask {
293 pub fn all() -> Self {
295 Self { mask: None }
296 }
297
298 pub fn none(len: usize) -> Self {
300 Self {
301 mask: Some(vec![false; len]),
302 }
303 }
304
305 pub fn leaves(schema: &SchemaDescriptor, indices: impl IntoIterator<Item = usize>) -> Self {
311 let mut mask = vec![false; schema.num_columns()];
312 for leaf_idx in indices {
313 mask[leaf_idx] = true;
314 }
315 Self { mask: Some(mask) }
316 }
317
318 pub fn roots(schema: &SchemaDescriptor, indices: impl IntoIterator<Item = usize>) -> Self {
324 let num_root_columns = schema.root_schema().get_fields().len();
325 let mut root_mask = vec![false; num_root_columns];
326 for root_idx in indices {
327 root_mask[root_idx] = true;
328 }
329
330 let mask = (0..schema.num_columns())
331 .map(|leaf_idx| {
332 let root_idx = schema.get_column_root_idx(leaf_idx);
333 root_mask[root_idx]
334 })
335 .collect();
336
337 Self { mask: Some(mask) }
338 }
339
340 pub fn columns<'a>(
371 schema: &SchemaDescriptor,
372 names: impl IntoIterator<Item = &'a str>,
373 ) -> Self {
374 let mut mask = vec![false; schema.num_columns()];
375 for name in names {
376 let name_path: Vec<&str> = name.split('.').collect();
377 for (idx, col) in schema.columns().iter().enumerate() {
378 let path = col.path().parts();
379 if name_path.len() > path.len() {
381 continue;
382 }
383 if name_path.iter().zip(path.iter()).all(|(a, b)| a == b) {
385 mask[idx] = true;
386 }
387 }
388 }
389
390 Self { mask: Some(mask) }
391 }
392
393 pub fn leaf_included(&self, leaf_idx: usize) -> bool {
395 self.mask.as_ref().map(|m| m[leaf_idx]).unwrap_or(true)
396 }
397
398 pub fn union(&mut self, other: &Self) {
407 match (self.mask.as_ref(), other.mask.as_ref()) {
408 (None, _) | (_, None) => self.mask = None,
409 (Some(a), Some(b)) => {
410 debug_assert_eq!(a.len(), b.len());
411 let mask = a.iter().zip(b.iter()).map(|(&a, &b)| a || b).collect();
412 self.mask = Some(mask);
413 }
414 }
415 }
416
417 pub fn intersect(&mut self, other: &Self) {
426 match (self.mask.as_ref(), other.mask.as_ref()) {
427 (None, _) => self.mask.clone_from(&other.mask),
428 (_, None) => {}
429 (Some(a), Some(b)) => {
430 debug_assert_eq!(a.len(), b.len());
431 let mask = a.iter().zip(b.iter()).map(|(&a, &b)| a && b).collect();
432 self.mask = Some(mask);
433 }
434 }
435 }
436
437 pub(crate) fn without_nested_types(&self, schema: &SchemaDescriptor) -> Option<Self> {
442 let num_leaves = schema.num_columns();
443
444 let num_roots = schema.root_schema().get_fields().len();
446 let mut root_leaf_counts = vec![0usize; num_roots];
447 for leaf_idx in 0..num_leaves {
448 let root_idx = schema.get_column_root_idx(leaf_idx);
449 root_leaf_counts[root_idx] += 1;
450 }
451
452 let mut included_leaves = Vec::new();
455 for leaf_idx in 0..num_leaves {
456 if self.leaf_included(leaf_idx) {
457 let root = schema.get_column_root(leaf_idx);
458 let root_idx = schema.get_column_root_idx(leaf_idx);
459 if root_leaf_counts[root_idx] == 1 && root.is_primitive() {
460 included_leaves.push(leaf_idx);
461 }
462 }
463 }
464
465 if included_leaves.is_empty() {
466 None
467 } else {
468 Some(ProjectionMask::leaves(schema, included_leaves))
469 }
470 }
471}
472
473pub fn parquet_column<'a>(
477 parquet_schema: &SchemaDescriptor,
478 arrow_schema: &'a Schema,
479 name: &str,
480) -> Option<(usize, &'a FieldRef)> {
481 let (root_idx, field) = arrow_schema.fields.find(name)?;
482 if field.data_type().is_nested() {
483 return None;
490 }
491
492 let parquet_idx = (0..parquet_schema.columns().len())
494 .find(|x| parquet_schema.get_column_root_idx(*x) == root_idx)?;
495 Some((parquet_idx, field))
496}
497
498#[cfg(test)]
499mod test {
500 use crate::arrow::ArrowWriter;
501 use crate::file::metadata::{
502 PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions, ParquetMetaDataReader,
503 ParquetMetaDataWriter,
504 };
505 use crate::file::properties::{EnabledStatistics, WriterProperties};
506 use crate::schema::parser::parse_message_type;
507 use crate::schema::types::SchemaDescriptor;
508 use arrow_array::{ArrayRef, Int32Array, RecordBatch};
509 use bytes::Bytes;
510 use std::sync::Arc;
511
512 use super::ProjectionMask;
513
514 #[test]
515 fn test_metadata_read_write_partial_offset() {
517 let parquet_bytes = create_parquet_file();
518
519 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
521 let original_metadata = ParquetMetaDataReader::new()
522 .with_metadata_options(Some(options))
523 .parse_and_finish(&parquet_bytes)
524 .unwrap();
525
526 let metadata_bytes = metadata_to_bytes(&original_metadata);
528 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
529 let err = ParquetMetaDataReader::new()
530 .with_metadata_options(Some(options))
531 .with_page_index_policy(PageIndexPolicy::Required) .parse_and_finish(&metadata_bytes)
533 .err()
534 .unwrap();
535 assert_eq!(
536 err.to_string(),
537 "EOF: Parquet file too small. Page index range 82..115 overlaps with file metadata 0..357"
538 );
539 }
540
541 #[test]
542 fn test_metadata_read_write_roundtrip() {
543 let parquet_bytes = create_parquet_file();
544
545 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
547 let original_metadata = ParquetMetaDataReader::new()
548 .with_metadata_options(Some(options))
549 .parse_and_finish(&parquet_bytes)
550 .unwrap();
551
552 let metadata_bytes = metadata_to_bytes(&original_metadata);
554 assert_ne!(
555 metadata_bytes.len(),
556 parquet_bytes.len(),
557 "metadata is subset of parquet"
558 );
559
560 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
561 let roundtrip_metadata = ParquetMetaDataReader::new()
562 .with_metadata_options(Some(options))
563 .parse_and_finish(&metadata_bytes)
564 .unwrap();
565
566 assert_eq!(original_metadata, roundtrip_metadata);
567 }
568
569 #[test]
570 #[cfg_attr(miri, ignore)] fn test_metadata_read_write_roundtrip_page_index() {
572 let parquet_bytes = create_parquet_file();
573
574 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
577 let original_metadata = ParquetMetaDataReader::new()
578 .with_metadata_options(Some(options))
579 .with_page_index_policy(PageIndexPolicy::Required)
580 .parse_and_finish(&parquet_bytes)
581 .unwrap();
582
583 let metadata_bytes = metadata_to_bytes(&original_metadata);
585 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
586 let roundtrip_metadata = ParquetMetaDataReader::new()
587 .with_metadata_options(Some(options))
588 .with_page_index_policy(PageIndexPolicy::Required)
589 .parse_and_finish(&metadata_bytes)
590 .unwrap();
591
592 let original_metadata = normalize_locations(original_metadata);
594 let roundtrip_metadata = normalize_locations(roundtrip_metadata);
595 assert_eq!(
596 format!("{original_metadata:#?}"),
597 format!("{roundtrip_metadata:#?}")
598 );
599 assert_eq!(original_metadata, roundtrip_metadata);
600 }
601
602 #[test]
603 fn test_metadata_read_write_roundtrip_missing_page_index() {
604 let parquet_bytes = create_parquet_file();
605
606 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
608 let original_metadata = ParquetMetaDataReader::new()
609 .with_metadata_options(Some(options))
610 .with_page_index_policy(PageIndexPolicy::Skip)
611 .parse_and_finish(&parquet_bytes)
612 .unwrap();
613
614 let metadata_bytes = metadata_to_bytes_no_page_idx(&original_metadata);
617 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
618 let roundtrip_metadata = ParquetMetaDataReader::new()
619 .with_metadata_options(Some(options))
620 .with_page_index_policy(PageIndexPolicy::Optional)
621 .parse_and_finish(&metadata_bytes)
622 .expect("page index locations should have been cleared");
623
624 assert!(roundtrip_metadata.page_index().is_none());
625 }
626
627 #[test]
628 fn test_metadata_read_write_roundtrip_offset_index_only() {
629 let array: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3]));
631 let batch = RecordBatch::try_from_iter(vec![("id", array)]).unwrap();
632 let props = WriterProperties::builder()
633 .set_statistics_enabled(EnabledStatistics::Chunk)
634 .build();
635 let mut buf = vec![];
636 let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props)).unwrap();
637 writer.write(&batch).unwrap();
638 writer.close().unwrap();
639
640 let read = |bytes: &Bytes| {
641 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
642 ParquetMetaDataReader::new()
643 .with_metadata_options(Some(options))
644 .with_page_index_policy(PageIndexPolicy::Optional)
645 .parse_and_finish(bytes)
646 .unwrap()
647 };
648 let original = read(&Bytes::from(buf));
649 let page_index = original.page_index().unwrap();
650 assert!(!page_index.has_column_indexes() && page_index.has_offset_indexes());
651 let roundtrip = read(&metadata_to_bytes(&original));
652 assert_eq!(
653 normalize_locations(original),
654 normalize_locations(roundtrip)
655 );
656 }
657
658 #[test]
659 fn test_metadata_read_write_roundtrip_custom_page_index() {
660 use crate::file::metadata::page_index::PageIndexProvider;
661 use crate::file::page_index::{
662 column_index::ColumnIndexMetaData, offset_index::OffsetIndexMetaData,
663 };
664
665 #[derive(Debug)]
667 struct Forward(Arc<dyn PageIndexProvider>);
668 impl PageIndexProvider for Forward {
669 fn has_offset_indexes(&self) -> bool {
670 self.0.has_offset_indexes()
671 }
672 fn has_column_indexes(&self) -> bool {
673 self.0.has_column_indexes()
674 }
675 fn column_index(&self, rg: usize, col: usize) -> Option<&ColumnIndexMetaData> {
676 self.0.column_index(rg, col)
677 }
678 fn offset_index(&self, rg: usize, col: usize) -> Option<&OffsetIndexMetaData> {
679 self.0.offset_index(rg, col)
680 }
681 fn as_any(&self) -> &dyn std::any::Any {
682 self
683 }
684 }
685
686 let read = |bytes: &Bytes| {
687 let options = ParquetMetaDataOptions::new().with_encoding_stats_as_mask(false);
688 ParquetMetaDataReader::new()
689 .with_metadata_options(Some(options))
690 .with_page_index_policy(PageIndexPolicy::Required)
691 .parse_and_finish(bytes)
692 .unwrap()
693 };
694 let original = read(&create_parquet_file());
695 let provider = Forward(original.page_index().unwrap().clone());
696 let custom = original
697 .clone()
698 .into_builder()
699 .set_page_index(Some(Arc::new(provider)))
700 .build();
701 let roundtrip = read(&metadata_to_bytes(&custom));
702 assert_eq!(
703 normalize_locations(original),
704 normalize_locations(roundtrip)
705 );
706 }
707
708 fn normalize_locations(metadata: ParquetMetaData) -> ParquetMetaData {
713 let mut metadata_builder = metadata.into_builder();
714 for rg in metadata_builder.take_row_groups() {
715 let mut rg_builder = rg.into_builder();
716 for col in rg_builder.take_columns() {
717 rg_builder = rg_builder.add_column_metadata(
718 col.into_builder()
719 .set_offset_index_offset(None)
720 .set_index_page_offset(None)
721 .set_column_index_offset(None)
722 .build()
723 .unwrap(),
724 );
725 }
726 let rg = rg_builder.build().unwrap();
727 metadata_builder = metadata_builder.add_row_group(rg);
728 }
729 metadata_builder.build()
730 }
731
732 fn create_parquet_file() -> Bytes {
734 let mut buf = vec![];
735 let data = vec![100, 200, 201, 300, 102, 33];
736 let array: ArrayRef = Arc::new(Int32Array::from(data));
737 let batch = RecordBatch::try_from_iter(vec![("id", array)]).unwrap();
738 let props = WriterProperties::builder()
739 .set_statistics_enabled(EnabledStatistics::Page)
740 .set_write_page_header_statistics(true)
741 .build();
742
743 let mut writer = ArrowWriter::try_new(&mut buf, batch.schema(), Some(props)).unwrap();
744 writer.write(&batch).unwrap();
745 writer.finish().unwrap();
746 drop(writer);
747
748 Bytes::from(buf)
749 }
750
751 fn metadata_to_bytes(metadata: &ParquetMetaData) -> Bytes {
753 let mut buf = vec![];
754 ParquetMetaDataWriter::new(&mut buf, metadata)
755 .finish()
756 .unwrap();
757 Bytes::from(buf)
758 }
759
760 fn metadata_to_bytes_no_page_idx(metadata: &ParquetMetaData) -> Bytes {
762 let mut buf = vec![];
763 ParquetMetaDataWriter::new(&mut buf, metadata)
764 .with_preserve_page_index_locations(false)
765 .finish()
766 .unwrap();
767 Bytes::from(buf)
768 }
769
770 #[test]
771 fn test_mask_from_column_names() {
772 let schema = parse_schema(
773 "
774 message test_schema {
775 OPTIONAL group a (MAP) {
776 REPEATED group key_value {
777 REQUIRED BYTE_ARRAY key (UTF8);
778 OPTIONAL group value (MAP) {
779 REPEATED group key_value {
780 REQUIRED INT32 key;
781 REQUIRED BOOLEAN value;
782 }
783 }
784 }
785 }
786 REQUIRED INT32 b;
787 REQUIRED DOUBLE c;
788 }
789 ",
790 );
791
792 let mask = ProjectionMask::columns(&schema, ["foo", "bar"]);
793 assert_eq!(mask.mask.unwrap(), vec![false; 5]);
794
795 let mask = ProjectionMask::columns(&schema, []);
796 assert_eq!(mask.mask.unwrap(), vec![false; 5]);
797
798 let mask = ProjectionMask::columns(&schema, ["a", "c"]);
799 assert_eq!(mask.mask.unwrap(), [true, true, true, false, true]);
800
801 let mask = ProjectionMask::columns(&schema, ["a.key_value.key", "c"]);
802 assert_eq!(mask.mask.unwrap(), [true, false, false, false, true]);
803
804 let mask = ProjectionMask::columns(&schema, ["a.key_value.value", "b"]);
805 assert_eq!(mask.mask.unwrap(), [false, true, true, true, false]);
806
807 let schema = parse_schema(
808 "
809 message test_schema {
810 OPTIONAL group a (LIST) {
811 REPEATED group list {
812 OPTIONAL group element (LIST) {
813 REPEATED group list {
814 OPTIONAL group element (LIST) {
815 REPEATED group list {
816 OPTIONAL BYTE_ARRAY element (UTF8);
817 }
818 }
819 }
820 }
821 }
822 }
823 REQUIRED INT32 b;
824 }
825 ",
826 );
827
828 let mask = ProjectionMask::columns(&schema, ["a", "b"]);
829 assert_eq!(mask.mask.unwrap(), [true, true]);
830
831 let mask = ProjectionMask::columns(&schema, ["a.list.element", "b"]);
832 assert_eq!(mask.mask.unwrap(), [true, true]);
833
834 let mask =
835 ProjectionMask::columns(&schema, ["a.list.element.list.element.list.element", "b"]);
836 assert_eq!(mask.mask.unwrap(), [true, true]);
837
838 let mask = ProjectionMask::columns(&schema, ["b"]);
839 assert_eq!(mask.mask.unwrap(), [false, true]);
840
841 let schema = parse_schema(
842 "
843 message test_schema {
844 OPTIONAL INT32 a;
845 OPTIONAL INT32 b;
846 OPTIONAL INT32 c;
847 OPTIONAL INT32 d;
848 OPTIONAL INT32 e;
849 }
850 ",
851 );
852
853 let mask = ProjectionMask::columns(&schema, ["a", "b"]);
854 assert_eq!(mask.mask.unwrap(), [true, true, false, false, false]);
855
856 let mask = ProjectionMask::columns(&schema, ["d", "b", "d"]);
857 assert_eq!(mask.mask.unwrap(), [false, true, false, true, false]);
858
859 let schema = parse_schema(
860 "
861 message test_schema {
862 OPTIONAL INT32 a;
863 OPTIONAL INT32 b;
864 OPTIONAL INT32 a;
865 OPTIONAL INT32 d;
866 OPTIONAL INT32 e;
867 }
868 ",
869 );
870
871 let mask = ProjectionMask::columns(&schema, ["a", "e"]);
872 assert_eq!(mask.mask.unwrap(), [true, false, true, false, true]);
873
874 let schema = parse_schema(
875 "
876 message test_schema {
877 OPTIONAL INT32 a;
878 OPTIONAL INT32 aa;
879 }
880 ",
881 );
882
883 let mask = ProjectionMask::columns(&schema, ["a"]);
884 assert_eq!(mask.mask.unwrap(), [true, false]);
885 }
886
887 #[test]
888 fn test_projection_mask_union() {
889 let mut mask1 = ProjectionMask {
890 mask: Some(vec![true, false, true]),
891 };
892 let mask2 = ProjectionMask {
893 mask: Some(vec![false, true, true]),
894 };
895 mask1.union(&mask2);
896 assert_eq!(mask1.mask, Some(vec![true, true, true]));
897
898 let mut mask1 = ProjectionMask { mask: None };
899 let mask2 = ProjectionMask {
900 mask: Some(vec![false, true, true]),
901 };
902 mask1.union(&mask2);
903 assert_eq!(mask1.mask, None);
904
905 let mut mask1 = ProjectionMask {
906 mask: Some(vec![true, false, true]),
907 };
908 let mask2 = ProjectionMask { mask: None };
909 mask1.union(&mask2);
910 assert_eq!(mask1.mask, None);
911
912 let mut mask1 = ProjectionMask { mask: None };
913 let mask2 = ProjectionMask { mask: None };
914 mask1.union(&mask2);
915 assert_eq!(mask1.mask, None);
916 }
917
918 #[test]
919 fn test_projection_mask_intersect() {
920 let mut mask1 = ProjectionMask {
921 mask: Some(vec![true, false, true]),
922 };
923 let mask2 = ProjectionMask {
924 mask: Some(vec![false, true, true]),
925 };
926 mask1.intersect(&mask2);
927 assert_eq!(mask1.mask, Some(vec![false, false, true]));
928
929 let mut mask1 = ProjectionMask { mask: None };
930 let mask2 = ProjectionMask {
931 mask: Some(vec![false, true, true]),
932 };
933 mask1.intersect(&mask2);
934 assert_eq!(mask1.mask, Some(vec![false, true, true]));
935
936 let mut mask1 = ProjectionMask {
937 mask: Some(vec![true, false, true]),
938 };
939 let mask2 = ProjectionMask { mask: None };
940 mask1.intersect(&mask2);
941 assert_eq!(mask1.mask, Some(vec![true, false, true]));
942
943 let mut mask1 = ProjectionMask { mask: None };
944 let mask2 = ProjectionMask { mask: None };
945 mask1.intersect(&mask2);
946 assert_eq!(mask1.mask, None);
947 }
948
949 #[test]
950 fn test_projection_mask_without_nested_no_nested() {
951 let schema = parse_schema(
953 "
954 message test_schema {
955 OPTIONAL INT32 a;
956 OPTIONAL INT32 b;
957 REQUIRED DOUBLE d;
958 }
959 ",
960 );
961
962 let mask = ProjectionMask::all();
963 assert_eq!(
965 Some(ProjectionMask::leaves(&schema, [0, 1, 2])),
966 mask.without_nested_types(&schema)
967 );
968
969 let mask = ProjectionMask::leaves(&schema, [1, 2]);
971 assert_eq!(Some(mask.clone()), mask.without_nested_types(&schema));
972 }
973
974 #[test]
975 fn test_projection_mask_without_nested_nested() {
976 let schema = parse_schema(
978 "
979 message test_schema {
980 OPTIONAL INT32 a;
981 OPTIONAL group b {
982 REQUIRED INT32 b1;
983 OPTIONAL INT64 b2;
984 }
985 OPTIONAL group c (LIST) {
986 REPEATED group list {
987 OPTIONAL INT32 element;
988 }
989 }
990 REQUIRED DOUBLE d;
991 }
992 ",
993 );
994
995 let mask = ProjectionMask::all();
997 assert_eq!(
998 Some(ProjectionMask::leaves(&schema, [0, 4])),
999 mask.without_nested_types(&schema)
1000 );
1001
1002 let mask = ProjectionMask::leaves(&schema, [1]);
1004 assert_eq!(None, mask.without_nested_types(&schema));
1005
1006 let mask = ProjectionMask::leaves(&schema, [1, 4]);
1008 assert_eq!(
1009 Some(ProjectionMask::leaves(&schema, [4])),
1010 mask.without_nested_types(&schema)
1011 );
1012
1013 let mask = ProjectionMask::leaves(&schema, [3]);
1015 assert_eq!(None, mask.without_nested_types(&schema));
1016 }
1017
1018 #[test]
1019 fn test_projection_mask_without_nested_map_only() {
1020 let schema = parse_schema(
1022 "
1023 message test_schema {
1024 required group my_map (MAP) {
1025 repeated group key_value {
1026 required binary key (STRING);
1027 optional int32 value;
1028 }
1029 }
1030 }
1031 ",
1032 );
1033
1034 let mask = ProjectionMask::all();
1035 assert_eq!(None, mask.without_nested_types(&schema));
1036
1037 let mask = ProjectionMask::leaves(&schema, [0]);
1039 assert_eq!(None, mask.without_nested_types(&schema));
1040
1041 let mask = ProjectionMask::leaves(&schema, [1]);
1043 assert_eq!(None, mask.without_nested_types(&schema));
1044 }
1045
1046 #[test]
1047 fn test_projection_mask_without_nested_map_with_non_nested() {
1048 let schema = parse_schema(
1051 "
1052 message test_schema {
1053 REQUIRED INT32 a;
1054 required group my_map (MAP) {
1055 repeated group key_value {
1056 required binary key (STRING);
1057 optional int32 value;
1058 }
1059 }
1060 REQUIRED INT32 b;
1061 }
1062 ",
1063 );
1064
1065 let mask = ProjectionMask::all();
1067 assert_eq!(
1068 Some(ProjectionMask::leaves(&schema, [0, 3])),
1069 mask.without_nested_types(&schema)
1070 );
1071
1072 let mask = ProjectionMask::leaves(&schema, [1, 2, 3]);
1074 assert_eq!(
1075 Some(ProjectionMask::leaves(&schema, [3])),
1076 mask.without_nested_types(&schema)
1077 );
1078
1079 let mask = ProjectionMask::leaves(&schema, [1, 2]);
1081 assert_eq!(None, mask.without_nested_types(&schema));
1082 }
1083
1084 #[test]
1085 fn test_projection_mask_without_nested_deeply_nested() {
1086 let schema = parse_schema(
1088 "
1089 message test_schema {
1090 OPTIONAL group a (MAP) {
1091 REPEATED group key_value {
1092 REQUIRED BYTE_ARRAY key (UTF8);
1093 OPTIONAL group value (MAP) {
1094 REPEATED group key_value {
1095 REQUIRED INT32 key;
1096 REQUIRED BOOLEAN value;
1097 }
1098 }
1099 }
1100 }
1101 REQUIRED INT32 b;
1102 REQUIRED DOUBLE c;
1103 ",
1104 );
1105
1106 let mask = ProjectionMask::all();
1107 assert_eq!(
1108 Some(ProjectionMask::leaves(&schema, [3, 4])),
1109 mask.without_nested_types(&schema)
1110 );
1111
1112 let mask = ProjectionMask::leaves(&schema, [0, 4]);
1114 assert_eq!(
1115 Some(ProjectionMask::leaves(&schema, [4])),
1116 mask.without_nested_types(&schema)
1117 );
1118
1119 let mask = ProjectionMask::leaves(&schema, [1, 2, 3]);
1121 assert_eq!(
1122 Some(ProjectionMask::leaves(&schema, [3])),
1123 mask.without_nested_types(&schema)
1124 );
1125
1126 let mask = ProjectionMask::leaves(&schema, [0]);
1128 assert_eq!(None, mask.without_nested_types(&schema));
1129 }
1130
1131 #[test]
1132 fn test_projection_mask_without_nested_list() {
1133 let schema = parse_schema(
1135 "
1136 message test_schema {
1137 required group my_list (LIST) {
1138 repeated group list {
1139 optional binary element (STRING);
1140 }
1141 }
1142 REQUIRED INT32 b;
1143 }
1144 ",
1145 );
1146
1147 let mask = ProjectionMask::all();
1148 assert_eq!(
1149 Some(ProjectionMask::leaves(&schema, [1])),
1150 mask.without_nested_types(&schema),
1151 );
1152
1153 let mask = ProjectionMask::leaves(&schema, [0]);
1155 assert_eq!(None, mask.without_nested_types(&schema));
1156
1157 let mask = ProjectionMask::leaves(&schema, [0, 1]);
1159 assert_eq!(
1160 Some(ProjectionMask::leaves(&schema, [1])),
1161 mask.without_nested_types(&schema),
1162 );
1163 }
1164
1165 #[test]
1166 fn test_projection_mask_without_nested_single_leaf_struct() {
1167 let schema = parse_schema(
1169 "
1170 message test_schema {
1171 OPTIONAL group address {
1172 REQUIRED BYTE_ARRAY street (UTF8);
1173 }
1174 REQUIRED INT32 id;
1175 }
1176 ",
1177 );
1178
1179 let mask = ProjectionMask::leaves(&schema, [0]);
1181 assert_eq!(None, mask.without_nested_types(&schema));
1182
1183 let mask = ProjectionMask::leaves(&schema, [0, 1]);
1185 assert_eq!(
1186 Some(ProjectionMask::leaves(&schema, [1])),
1187 mask.without_nested_types(&schema)
1188 );
1189
1190 let mask = ProjectionMask::all();
1192 assert_eq!(
1193 Some(ProjectionMask::leaves(&schema, [1])),
1194 mask.without_nested_types(&schema)
1195 );
1196 }
1197
1198 fn parse_schema(schema: &str) -> SchemaDescriptor {
1200 let parquet_group_type = parse_message_type(schema).unwrap();
1201 SchemaDescriptor::new(Arc::new(parquet_group_type))
1202 }
1203}