1use std::fmt::Formatter;
25use std::io::SeekFrom;
26use std::ops::Range;
27use std::pin::Pin;
28use std::sync::Arc;
29use std::task::{Context, Poll};
30
31use bytes::Bytes;
32use futures::future::{BoxFuture, FutureExt};
33use futures::stream::Stream;
34use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeek, AsyncSeekExt};
35
36use arrow_array::RecordBatch;
37use arrow_schema::{Schema, SchemaRef};
38
39use crate::arrow::arrow_reader::{
40 ArrowReaderBuilder, ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReader,
41};
42
43use crate::basic::{BloomFilterAlgorithm, BloomFilterCompression, BloomFilterHash};
44use crate::bloom_filter::{
45 SBBF_HEADER_SIZE_ESTIMATE, Sbbf, chunk_read_bloom_filter_header_and_offset,
46};
47use crate::errors::{ParquetError, Result};
48use crate::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
49
50mod metadata;
51pub use metadata::*;
52
53mod spawn;
54pub use spawn::SpawnedReader;
55
56#[cfg(feature = "object_store")]
57mod store;
58
59use crate::DecodeResult;
60use crate::arrow::push_decoder::{ParquetPushDecoder, ParquetPushDecoderBuilder, PushDecoderInput};
61#[allow(deprecated)]
62#[cfg(feature = "object_store")]
63pub use store::*;
64
65pub trait AsyncFileReader: Send {
163 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>>;
165
166 fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_, Result<Vec<Bytes>>> {
168 async move {
169 let mut result = Vec::with_capacity(ranges.len());
170
171 for range in ranges.into_iter() {
172 let data = self.get_bytes(range).await?;
173 result.push(data);
174 }
175
176 Ok(result)
177 }
178 .boxed()
179 }
180
181 fn get_metadata<'a>(
198 &'a mut self,
199 options: Option<&'a ArrowReaderOptions>,
200 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>>;
201}
202
203impl AsyncFileReader for Box<dyn AsyncFileReader + '_> {
205 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
206 self.as_mut().get_bytes(range)
207 }
208
209 fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_, Result<Vec<Bytes>>> {
210 self.as_mut().get_byte_ranges(ranges)
211 }
212
213 fn get_metadata<'a>(
214 &'a mut self,
215 options: Option<&'a ArrowReaderOptions>,
216 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
217 self.as_mut().get_metadata(options)
218 }
219}
220
221impl<T: AsyncFileReader + MetadataFetch + AsyncRead + AsyncSeek + Unpin> MetadataSuffixFetch for T {
222 fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
223 async move {
224 self.seek(SeekFrom::End(-(suffix as i64))).await?;
225 let mut buf = Vec::with_capacity(suffix);
226 self.take(suffix as _).read_to_end(&mut buf).await?;
227 Ok(buf.into())
228 }
229 .boxed()
230 }
231}
232
233impl<T: AsyncRead + AsyncSeek + Unpin + Send> AsyncFileReader for T {
234 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
235 async move {
236 self.seek(SeekFrom::Start(range.start)).await?;
237
238 let to_read = range.end - range.start;
239 let mut buffer = Vec::with_capacity(to_read.try_into()?);
240 let read = self.take(to_read).read_to_end(&mut buffer).await?;
241 if read as u64 != to_read {
242 return Err(eof_err!("expected to read {} bytes, got {}", to_read, read));
243 }
244
245 Ok(buffer.into())
246 }
247 .boxed()
248 }
249
250 fn get_metadata<'a>(
251 &'a mut self,
252 options: Option<&'a ArrowReaderOptions>,
253 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
254 async move {
255 let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
256 let parquet_metadata = metadata_reader.load_via_suffix_and_finish(self).await?;
257 Ok(Arc::new(parquet_metadata))
258 }
259 .boxed()
260 }
261}
262
263impl ArrowReaderMetadata {
264 pub async fn load_async<T: AsyncFileReader>(
268 input: &mut T,
269 options: ArrowReaderOptions,
270 ) -> Result<Self> {
271 let metadata = input.get_metadata(Some(&options)).await?;
272 Self::try_new(metadata, options)
273 }
274}
275
276#[doc(hidden)]
277pub struct AsyncReader<T>(T);
282
283pub type ParquetRecordBatchStreamBuilder<T> = ArrowReaderBuilder<AsyncReader<T>>;
303
304impl<T: AsyncFileReader + Send + 'static> ParquetRecordBatchStreamBuilder<T> {
305 pub async fn new(input: T) -> Result<Self> {
521 Self::new_with_options(input, Default::default()).await
522 }
523
524 pub async fn new_with_options(mut input: T, options: ArrowReaderOptions) -> Result<Self> {
527 let metadata = ArrowReaderMetadata::load_async(&mut input, options).await?;
528 Ok(Self::new_with_metadata(input, metadata))
529 }
530
531 pub fn new_with_metadata(input: T, metadata: ArrowReaderMetadata) -> Self {
577 Self::new_builder(AsyncReader(input), metadata)
578 }
579
580 pub async fn get_row_group_column_bloom_filter(
586 &mut self,
587 row_group_idx: usize,
588 column_idx: usize,
589 ) -> Result<Option<Sbbf>> {
590 let metadata = self.metadata.row_group(row_group_idx);
591 let column_metadata = metadata.column(column_idx);
592
593 let offset: u64 = if let Some(offset) = column_metadata.bloom_filter_offset() {
594 offset
595 .try_into()
596 .map_err(|_| ParquetError::General("Bloom filter offset is invalid".to_string()))?
597 } else {
598 return Ok(None);
599 };
600
601 let buffer = match column_metadata.bloom_filter_length() {
602 Some(length) => self.input.0.get_bytes(offset..offset + length as u64),
603 None => self
604 .input
605 .0
606 .get_bytes(offset..offset + SBBF_HEADER_SIZE_ESTIMATE as u64),
607 }
608 .await?;
609
610 let (header, bitset_offset) =
611 chunk_read_bloom_filter_header_and_offset(offset, buffer.clone())?;
612
613 match header.algorithm {
614 BloomFilterAlgorithm::BLOCK => {
615 }
617 }
618 match header.compression {
619 BloomFilterCompression::UNCOMPRESSED => {
620 }
622 }
623 match header.hash {
624 BloomFilterHash::XXHASH => {
625 }
627 }
628
629 let bitset = match column_metadata.bloom_filter_length() {
630 Some(_) => buffer.slice(
631 (TryInto::<usize>::try_into(bitset_offset).unwrap()
632 - TryInto::<usize>::try_into(offset).unwrap())..,
633 ),
634 None => {
635 let bitset_length: u64 = header.num_bytes.try_into().map_err(|_| {
636 ParquetError::General("Bloom filter length is invalid".to_string())
637 })?;
638 self.input
639 .0
640 .get_bytes(bitset_offset..bitset_offset + bitset_length)
641 .await?
642 }
643 };
644 Ok(Some(Sbbf::new(&bitset)))
645 }
646
647 pub fn build(self) -> Result<ParquetRecordBatchStream<T>> {
651 let Self {
652 input,
653 metadata,
654 schema,
655 fields,
656 batch_size,
657 row_groups,
658 projection,
659 filter,
660 selection,
661 row_selection_policy: selection_strategy,
662 limit,
663 offset,
664 metrics,
665 max_predicate_cache_size,
666 } = self;
667
668 let projection_len = projection.mask.as_ref().map_or(usize::MAX, |m| m.len());
671 let projected_fields = schema
672 .fields
673 .filter_leaves(|idx, _| idx < projection_len && projection.leaf_included(idx));
674 let projected_schema = Arc::new(Schema::new(projected_fields));
675
676 let decoder = ParquetPushDecoderBuilder {
677 input: PushDecoderInput::default(),
678 metadata,
679 schema,
680 fields,
681 projection,
682 filter,
683 selection,
684 row_selection_policy: selection_strategy,
685 batch_size,
686 row_groups,
687 limit,
688 offset,
689 metrics,
690 max_predicate_cache_size,
691 }
692 .build()?;
693
694 let request_state = RequestState::None { input: input.0 };
695
696 Ok(ParquetRecordBatchStream {
697 schema: projected_schema,
698 decoder,
699 request_state,
700 })
701 }
702}
703
704enum RequestState<T> {
708 None {
710 input: T,
711 },
712 Outstanding {
714 ranges: Vec<Range<u64>>,
716 future: BoxFuture<'static, Result<(T, Vec<Bytes>)>>,
721 },
722 Done,
723}
724
725impl<T> RequestState<T>
726where
727 T: AsyncFileReader + Unpin + Send + 'static,
728{
729 fn begin_request(mut input: T, ranges: Vec<Range<u64>>) -> Self {
731 let ranges_captured = ranges.clone();
732
733 let future = async move {
738 let data = input.get_byte_ranges(ranges_captured).await?;
739 Ok((input, data))
740 }
741 .boxed();
742 RequestState::Outstanding { ranges, future }
743 }
744}
745
746impl<T> std::fmt::Debug for RequestState<T> {
747 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
748 match self {
749 RequestState::None { input: _ } => f
750 .debug_struct("RequestState::None")
751 .field("input", &"...")
752 .finish(),
753 RequestState::Outstanding { ranges, .. } => f
754 .debug_struct("RequestState::Outstanding")
755 .field("ranges", &ranges)
756 .finish(),
757 RequestState::Done => {
758 write!(f, "RequestState::Done")
759 }
760 }
761 }
762}
763
764pub struct ParquetRecordBatchStream<T> {
784 schema: SchemaRef,
786 request_state: RequestState<T>,
788 decoder: ParquetPushDecoder,
790}
791
792impl<T> std::fmt::Debug for ParquetRecordBatchStream<T> {
793 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
794 f.debug_struct("ParquetRecordBatchStream")
795 .field("request_state", &self.request_state)
796 .finish()
797 }
798}
799
800impl<T> ParquetRecordBatchStream<T> {
801 pub fn schema(&self) -> &SchemaRef {
806 &self.schema
807 }
808}
809
810impl<T> ParquetRecordBatchStream<T>
811where
812 T: AsyncFileReader + Unpin + Send + 'static,
813{
814 pub async fn next_row_group(&mut self) -> Result<Option<ParquetRecordBatchReader>> {
828 loop {
829 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
832 match request_state {
833 RequestState::None { input } => {
835 match self.decoder.try_next_reader()? {
836 DecodeResult::NeedsData(ranges) => {
837 self.request_state = RequestState::begin_request(input, ranges);
838 continue; }
840 DecodeResult::Data(reader) => {
841 self.request_state = RequestState::None { input };
842 return Ok(Some(reader));
843 }
844 DecodeResult::Finished => return Ok(None),
845 }
846 }
847 RequestState::Outstanding { ranges, future } => {
848 let (input, data) = future.await?;
849 self.decoder.push_ranges(ranges, data)?;
851 self.request_state = RequestState::None { input };
852 continue; }
854 RequestState::Done => {
855 self.request_state = RequestState::Done;
856 return Ok(None);
857 }
858 }
859 }
860 }
861}
862
863impl<T> Stream for ParquetRecordBatchStream<T>
864where
865 T: AsyncFileReader + Unpin + Send + 'static,
866{
867 type Item = Result<RecordBatch>;
868 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
869 match self.poll_next_inner(cx) {
870 Ok(res) => {
871 res.map(|res| Ok(res).transpose())
874 }
875 Err(e) => {
876 self.request_state = RequestState::Done;
877 Poll::Ready(Some(Err(e)))
878 }
879 }
880 }
881}
882
883impl<T> ParquetRecordBatchStream<T>
884where
885 T: AsyncFileReader + Unpin + Send + 'static,
886{
887 fn poll_next_inner(&mut self, cx: &mut Context<'_>) -> Result<Poll<Option<RecordBatch>>> {
892 loop {
893 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
894 match request_state {
895 RequestState::None { input } => {
896 match self.decoder.try_decode()? {
898 DecodeResult::NeedsData(ranges) => {
899 self.request_state = RequestState::begin_request(input, ranges);
900 continue; }
902 DecodeResult::Data(batch) => {
903 self.request_state = RequestState::None { input };
904 return Ok(Poll::Ready(Some(batch)));
905 }
906 DecodeResult::Finished => {
907 self.request_state = RequestState::Done;
908 return Ok(Poll::Ready(None));
909 }
910 }
911 }
912 RequestState::Outstanding { ranges, mut future } => match future.poll_unpin(cx) {
913 Poll::Ready(result) => {
915 let (input, data) = result?;
916 self.decoder.push_ranges(ranges, data)?;
918 self.request_state = RequestState::None { input };
919 continue; }
921 Poll::Pending => {
922 self.request_state = RequestState::Outstanding { ranges, future };
923 return Ok(Poll::Pending);
924 }
925 },
926 RequestState::Done => {
927 self.request_state = RequestState::Done;
929 return Ok(Poll::Ready(None));
930 }
931 }
932 }
933 }
934}
935
936#[cfg(test)]
937mod tests {
938 use super::*;
939 use crate::arrow::arrow_reader::tests::test_row_numbers_with_multiple_row_groups_helper;
940 use crate::arrow::arrow_reader::{
941 ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowFilter, RowSelection, RowSelector,
942 };
943 use crate::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
944 use crate::arrow::schema::virtual_type::RowNumber;
945 use crate::arrow::{ArrowWriter, AsyncArrowWriter, ProjectionMask};
946 use crate::file::metadata::PageIndexPolicy;
947 use crate::file::metadata::ParquetMetaDataReader;
948 use crate::file::properties::WriterProperties;
949 use arrow::compute::kernels::cmp::eq;
950 use arrow::error::Result as ArrowResult;
951 use arrow_array::builder::{Float32Builder, ListBuilder, StringBuilder};
952 use arrow_array::cast::AsArray;
953 use arrow_array::types::Int32Type;
954 use arrow_array::{
955 Array, ArrayRef, BooleanArray, Int32Array, RecordBatchReader, Scalar, StringArray,
956 StructArray, UInt64Array,
957 };
958 use arrow_schema::{DataType, Field, Schema};
959 use futures::{StreamExt, TryStreamExt};
960 use rand::{Rng, rng};
961 use std::collections::HashMap;
962 use std::sync::{Arc, Mutex};
963 use tempfile::tempfile;
964
965 #[derive(Clone)]
966 struct TestReader {
967 data: Bytes,
968 metadata: Option<Arc<ParquetMetaData>>,
969 requests: Arc<Mutex<Vec<Range<usize>>>>,
970 }
971
972 impl TestReader {
973 fn new(data: Bytes) -> Self {
974 Self {
975 data,
976 metadata: Default::default(),
977 requests: Default::default(),
978 }
979 }
980 }
981
982 impl AsyncFileReader for TestReader {
983 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
984 let range = range.clone();
985 self.requests
986 .lock()
987 .unwrap()
988 .push(range.start as usize..range.end as usize);
989 futures::future::ready(Ok(self
990 .data
991 .slice(range.start as usize..range.end as usize)))
992 .boxed()
993 }
994
995 fn get_metadata<'a>(
996 &'a mut self,
997 options: Option<&'a ArrowReaderOptions>,
998 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
999 let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
1000 self.metadata = Some(Arc::new(
1001 metadata_reader.parse_and_finish(&self.data).unwrap(),
1002 ));
1003 futures::future::ready(Ok(self.metadata.clone().unwrap().clone())).boxed()
1004 }
1005 }
1006
1007 #[tokio::test]
1008 async fn test_async_reader() {
1009 let testdata = arrow::util::test_util::parquet_test_data();
1010 let path = format!("{testdata}/alltypes_plain.parquet");
1011 let data = Bytes::from(std::fs::read(path).unwrap());
1012
1013 let async_reader = TestReader::new(data.clone());
1014
1015 let requests = async_reader.requests.clone();
1016 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1017 .await
1018 .unwrap();
1019
1020 let metadata = builder.metadata().clone();
1021 assert_eq!(metadata.num_row_groups(), 1);
1022
1023 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1024 let stream = builder
1025 .with_projection(mask.clone())
1026 .with_batch_size(1024)
1027 .build()
1028 .unwrap();
1029
1030 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1031
1032 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1033 .unwrap()
1034 .with_projection(mask)
1035 .with_batch_size(104)
1036 .build()
1037 .unwrap()
1038 .collect::<ArrowResult<Vec<_>>>()
1039 .unwrap();
1040
1041 assert_eq!(async_batches, sync_batches);
1042
1043 let requests = requests.lock().unwrap();
1044 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1045 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1046
1047 assert_eq!(
1048 &requests[..],
1049 &[
1050 offset_1 as usize..(offset_1 + length_1) as usize,
1051 offset_2 as usize..(offset_2 + length_2) as usize
1052 ]
1053 );
1054 }
1055
1056 #[tokio::test]
1057 async fn test_async_reader_with_next_row_group() {
1058 let testdata = arrow::util::test_util::parquet_test_data();
1059 let path = format!("{testdata}/alltypes_plain.parquet");
1060 let data = Bytes::from(std::fs::read(path).unwrap());
1061
1062 let async_reader = TestReader::new(data.clone());
1063
1064 let requests = async_reader.requests.clone();
1065 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1066 .await
1067 .unwrap();
1068
1069 let metadata = builder.metadata().clone();
1070 assert_eq!(metadata.num_row_groups(), 1);
1071
1072 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1073 let mut stream = builder
1074 .with_projection(mask.clone())
1075 .with_batch_size(1024)
1076 .build()
1077 .unwrap();
1078
1079 let mut readers = vec![];
1080 while let Some(reader) = stream.next_row_group().await.unwrap() {
1081 readers.push(reader);
1082 }
1083
1084 let async_batches: Vec<_> = readers
1085 .into_iter()
1086 .flat_map(|r| r.map(|v| v.unwrap()).collect::<Vec<_>>())
1087 .collect();
1088
1089 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1090 .unwrap()
1091 .with_projection(mask)
1092 .with_batch_size(104)
1093 .build()
1094 .unwrap()
1095 .collect::<ArrowResult<Vec<_>>>()
1096 .unwrap();
1097
1098 assert_eq!(async_batches, sync_batches);
1099
1100 let requests = requests.lock().unwrap();
1101 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1102 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1103
1104 assert_eq!(
1105 &requests[..],
1106 &[
1107 offset_1 as usize..(offset_1 + length_1) as usize,
1108 offset_2 as usize..(offset_2 + length_2) as usize
1109 ]
1110 );
1111 }
1112
1113 #[tokio::test]
1114 async fn test_async_reader_with_index() {
1115 let testdata = arrow::util::test_util::parquet_test_data();
1116 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1117 let data = Bytes::from(std::fs::read(path).unwrap());
1118
1119 let async_reader = TestReader::new(data.clone());
1120
1121 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1122 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1123 .await
1124 .unwrap();
1125
1126 let metadata_with_index = builder.metadata();
1128 assert_eq!(metadata_with_index.num_row_groups(), 1);
1129
1130 let offset_index = metadata_with_index.offset_index().unwrap();
1132 let column_index = metadata_with_index.column_index().unwrap();
1133
1134 assert_eq!(offset_index.len(), metadata_with_index.num_row_groups());
1135 assert_eq!(column_index.len(), metadata_with_index.num_row_groups());
1136
1137 let num_columns = metadata_with_index
1138 .file_metadata()
1139 .schema_descr()
1140 .num_columns();
1141
1142 offset_index
1144 .iter()
1145 .for_each(|x| assert_eq!(x.len(), num_columns));
1146 column_index
1147 .iter()
1148 .for_each(|x| assert_eq!(x.len(), num_columns));
1149
1150 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1151 let stream = builder
1152 .with_projection(mask.clone())
1153 .with_batch_size(1024)
1154 .build()
1155 .unwrap();
1156
1157 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1158
1159 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1160 .unwrap()
1161 .with_projection(mask)
1162 .with_batch_size(1024)
1163 .build()
1164 .unwrap()
1165 .collect::<ArrowResult<Vec<_>>>()
1166 .unwrap();
1167
1168 assert_eq!(async_batches, sync_batches);
1169 }
1170
1171 #[tokio::test]
1172 async fn test_async_reader_with_limit() {
1173 let testdata = arrow::util::test_util::parquet_test_data();
1174 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1175 let data = Bytes::from(std::fs::read(path).unwrap());
1176
1177 let metadata = ParquetMetaDataReader::new()
1178 .parse_and_finish(&data)
1179 .unwrap();
1180 let metadata = Arc::new(metadata);
1181
1182 assert_eq!(metadata.num_row_groups(), 1);
1183
1184 let async_reader = TestReader::new(data.clone());
1185
1186 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1187 .await
1188 .unwrap();
1189
1190 assert_eq!(builder.metadata().num_row_groups(), 1);
1191
1192 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1193 let stream = builder
1194 .with_projection(mask.clone())
1195 .with_batch_size(1024)
1196 .with_limit(1)
1197 .build()
1198 .unwrap();
1199
1200 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1201
1202 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1203 .unwrap()
1204 .with_projection(mask)
1205 .with_batch_size(1024)
1206 .with_limit(1)
1207 .build()
1208 .unwrap()
1209 .collect::<ArrowResult<Vec<_>>>()
1210 .unwrap();
1211
1212 assert_eq!(async_batches, sync_batches);
1213 }
1214
1215 #[tokio::test]
1216 async fn test_async_reader_skip_pages() {
1217 let testdata = arrow::util::test_util::parquet_test_data();
1218 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1219 let data = Bytes::from(std::fs::read(path).unwrap());
1220
1221 let async_reader = TestReader::new(data.clone());
1222
1223 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1224 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1225 .await
1226 .unwrap();
1227
1228 assert_eq!(builder.metadata().num_row_groups(), 1);
1229
1230 let selection = RowSelection::from(vec![
1231 RowSelector::skip(21), RowSelector::select(21), RowSelector::skip(41), RowSelector::select(41), RowSelector::skip(25), RowSelector::select(25), RowSelector::skip(7116), RowSelector::select(10), ]);
1240
1241 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![9]);
1242
1243 let stream = builder
1244 .with_projection(mask.clone())
1245 .with_row_selection(selection.clone())
1246 .build()
1247 .expect("building stream");
1248
1249 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1250
1251 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1252 .unwrap()
1253 .with_projection(mask)
1254 .with_batch_size(1024)
1255 .with_row_selection(selection)
1256 .build()
1257 .unwrap()
1258 .collect::<ArrowResult<Vec<_>>>()
1259 .unwrap();
1260
1261 assert_eq!(async_batches, sync_batches);
1262 }
1263
1264 #[tokio::test]
1265 async fn test_fuzz_async_reader_selection() {
1266 let testdata = arrow::util::test_util::parquet_test_data();
1267 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1268 let data = Bytes::from(std::fs::read(path).unwrap());
1269
1270 let mut rand = rng();
1271
1272 for _ in 0..100 {
1273 let mut expected_rows = 0;
1274 let mut total_rows = 0;
1275 let mut skip = false;
1276 let mut selectors = vec![];
1277
1278 while total_rows < 7300 {
1279 let row_count: usize = rand.random_range(1..100);
1280
1281 let row_count = row_count.min(7300 - total_rows);
1282
1283 selectors.push(RowSelector { row_count, skip });
1284
1285 total_rows += row_count;
1286 if !skip {
1287 expected_rows += row_count;
1288 }
1289
1290 skip = !skip;
1291 }
1292
1293 let selection = RowSelection::from(selectors);
1294
1295 let async_reader = TestReader::new(data.clone());
1296
1297 let options =
1298 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1299 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1300 .await
1301 .unwrap();
1302
1303 assert_eq!(builder.metadata().num_row_groups(), 1);
1304
1305 let col_idx: usize = rand.random_range(0..13);
1306 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1307
1308 let stream = builder
1309 .with_projection(mask.clone())
1310 .with_row_selection(selection.clone())
1311 .build()
1312 .expect("building stream");
1313
1314 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1315
1316 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1317
1318 assert_eq!(actual_rows, expected_rows);
1319 }
1320 }
1321
1322 #[tokio::test]
1323 async fn test_async_reader_zero_row_selector() {
1324 let testdata = arrow::util::test_util::parquet_test_data();
1326 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1327 let data = Bytes::from(std::fs::read(path).unwrap());
1328
1329 let mut rand = rng();
1330
1331 let mut expected_rows = 0;
1332 let mut total_rows = 0;
1333 let mut skip = false;
1334 let mut selectors = vec![];
1335
1336 selectors.push(RowSelector {
1337 row_count: 0,
1338 skip: false,
1339 });
1340
1341 while total_rows < 7300 {
1342 let row_count: usize = rand.random_range(1..100);
1343
1344 let row_count = row_count.min(7300 - total_rows);
1345
1346 selectors.push(RowSelector { row_count, skip });
1347
1348 total_rows += row_count;
1349 if !skip {
1350 expected_rows += row_count;
1351 }
1352
1353 skip = !skip;
1354 }
1355
1356 let selection = RowSelection::from(selectors);
1357
1358 let async_reader = TestReader::new(data.clone());
1359
1360 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1361 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1362 .await
1363 .unwrap();
1364
1365 assert_eq!(builder.metadata().num_row_groups(), 1);
1366
1367 let col_idx: usize = rand.random_range(0..13);
1368 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1369
1370 let stream = builder
1371 .with_projection(mask.clone())
1372 .with_row_selection(selection.clone())
1373 .build()
1374 .expect("building stream");
1375
1376 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1377
1378 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1379
1380 assert_eq!(actual_rows, expected_rows);
1381 }
1382
1383 #[tokio::test]
1384 async fn test_limit_multiple_row_groups() {
1385 let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]);
1386 let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]);
1387 let c = Int32Array::from_iter(0..6);
1388 let data = RecordBatch::try_from_iter([
1389 ("a", Arc::new(a) as ArrayRef),
1390 ("b", Arc::new(b) as ArrayRef),
1391 ("c", Arc::new(c) as ArrayRef),
1392 ])
1393 .unwrap();
1394
1395 let mut buf = Vec::with_capacity(1024);
1396 let props = WriterProperties::builder()
1397 .set_max_row_group_row_count(Some(3))
1398 .build();
1399 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1400 writer.write(&data).unwrap();
1401 writer.close().unwrap();
1402
1403 let data: Bytes = buf.into();
1404 let metadata = ParquetMetaDataReader::new()
1405 .parse_and_finish(&data)
1406 .unwrap();
1407
1408 assert_eq!(metadata.num_row_groups(), 2);
1409
1410 let test = TestReader::new(data);
1411
1412 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1413 .await
1414 .unwrap()
1415 .with_batch_size(1024)
1416 .with_limit(4)
1417 .build()
1418 .unwrap();
1419
1420 let batches: Vec<_> = stream.try_collect().await.unwrap();
1421 assert_eq!(batches.len(), 2);
1423
1424 let batch = &batches[0];
1425 assert_eq!(batch.num_rows(), 3);
1427 assert_eq!(batch.num_columns(), 3);
1428 let col2 = batch.column(2).as_primitive::<Int32Type>();
1429 assert_eq!(col2.values(), &[0, 1, 2]);
1430
1431 let batch = &batches[1];
1432 assert_eq!(batch.num_rows(), 1);
1434 assert_eq!(batch.num_columns(), 3);
1435 let col2 = batch.column(2).as_primitive::<Int32Type>();
1436 assert_eq!(col2.values(), &[3]);
1437
1438 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1439 .await
1440 .unwrap()
1441 .with_offset(2)
1442 .with_limit(3)
1443 .build()
1444 .unwrap();
1445
1446 let batches: Vec<_> = stream.try_collect().await.unwrap();
1447 assert_eq!(batches.len(), 2);
1449
1450 let batch = &batches[0];
1451 assert_eq!(batch.num_rows(), 1);
1453 assert_eq!(batch.num_columns(), 3);
1454 let col2 = batch.column(2).as_primitive::<Int32Type>();
1455 assert_eq!(col2.values(), &[2]);
1456
1457 let batch = &batches[1];
1458 assert_eq!(batch.num_rows(), 2);
1460 assert_eq!(batch.num_columns(), 3);
1461 let col2 = batch.column(2).as_primitive::<Int32Type>();
1462 assert_eq!(col2.values(), &[3, 4]);
1463
1464 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1465 .await
1466 .unwrap()
1467 .with_offset(4)
1468 .with_limit(20)
1469 .build()
1470 .unwrap();
1471
1472 let batches: Vec<_> = stream.try_collect().await.unwrap();
1473 assert_eq!(batches.len(), 1);
1475
1476 let batch = &batches[0];
1477 assert_eq!(batch.num_rows(), 2);
1479 assert_eq!(batch.num_columns(), 3);
1480 let col2 = batch.column(2).as_primitive::<Int32Type>();
1481 assert_eq!(col2.values(), &[4, 5]);
1482 }
1483
1484 #[tokio::test]
1485 async fn test_batch_size_overallocate() {
1486 let testdata = arrow::util::test_util::parquet_test_data();
1487 let path = format!("{testdata}/alltypes_plain.parquet");
1489 let data = Bytes::from(std::fs::read(path).unwrap());
1490
1491 let async_reader = TestReader::new(data.clone());
1492
1493 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1494 .await
1495 .unwrap();
1496
1497 let file_rows = builder.metadata().file_metadata().num_rows() as usize;
1498
1499 let builder = builder
1500 .with_projection(ProjectionMask::all())
1501 .with_batch_size(1024);
1502
1503 assert_ne!(1024, file_rows);
1506 assert_eq!(builder.batch_size, file_rows);
1507
1508 let _stream = builder.build().unwrap();
1509 }
1510
1511 #[tokio::test]
1512 async fn test_parquet_record_batch_stream_schema() {
1513 fn get_all_field_names(schema: &Schema) -> Vec<&String> {
1514 schema.flattened_fields().iter().map(|f| f.name()).collect()
1515 }
1516
1517 let mut metadata = HashMap::with_capacity(1);
1526 metadata.insert("key".to_string(), "value".to_string());
1527
1528 let nested_struct_array = StructArray::from(vec![
1529 (
1530 Arc::new(Field::new("d", DataType::Utf8, true)),
1531 Arc::new(StringArray::from(vec!["a", "b"])) as ArrayRef,
1532 ),
1533 (
1534 Arc::new(Field::new("e", DataType::Utf8, true)),
1535 Arc::new(StringArray::from(vec!["c", "d"])) as ArrayRef,
1536 ),
1537 ]);
1538 let struct_array = StructArray::from(vec![
1539 (
1540 Arc::new(Field::new("a", DataType::Int32, true)),
1541 Arc::new(Int32Array::from(vec![-1, 1])) as ArrayRef,
1542 ),
1543 (
1544 Arc::new(Field::new("b", DataType::UInt64, true)),
1545 Arc::new(UInt64Array::from(vec![1, 2])) as ArrayRef,
1546 ),
1547 (
1548 Arc::new(Field::new(
1549 "c",
1550 nested_struct_array.data_type().clone(),
1551 true,
1552 )),
1553 Arc::new(nested_struct_array) as ArrayRef,
1554 ),
1555 ]);
1556
1557 let schema =
1558 Arc::new(Schema::new(struct_array.fields().clone()).with_metadata(metadata.clone()));
1559 let record_batch = RecordBatch::from(struct_array)
1560 .with_schema(schema.clone())
1561 .unwrap();
1562
1563 let mut file = tempfile().unwrap();
1565 let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1566 writer.write(&record_batch).unwrap();
1567 writer.close().unwrap();
1568
1569 let all_fields = ["a", "b", "c", "d", "e"];
1570 let projections = [
1572 (vec![], vec![]),
1573 (vec![0], vec!["a"]),
1574 (vec![0, 1], vec!["a", "b"]),
1575 (vec![0, 1, 2], vec!["a", "b", "c", "d"]),
1576 (vec![0, 1, 2, 3], vec!["a", "b", "c", "d", "e"]),
1577 ];
1578
1579 for (indices, expected_projected_names) in projections {
1581 let assert_schemas = |builder: SchemaRef, reader: SchemaRef, batch: SchemaRef| {
1582 assert_eq!(get_all_field_names(&builder), all_fields);
1584 assert_eq!(builder.metadata, metadata);
1585 assert_eq!(get_all_field_names(&reader), expected_projected_names);
1587 assert_eq!(reader.metadata, HashMap::default());
1588 assert_eq!(get_all_field_names(&batch), expected_projected_names);
1589 assert_eq!(batch.metadata, HashMap::default());
1590 };
1591
1592 let builder =
1593 ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1594 let sync_builder_schema = builder.schema().clone();
1595 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices.clone());
1596 let mut reader = builder.with_projection(mask).build().unwrap();
1597 let sync_reader_schema = reader.schema();
1598 let batch = reader.next().unwrap().unwrap();
1599 let sync_batch_schema = batch.schema();
1600 assert_schemas(sync_builder_schema, sync_reader_schema, sync_batch_schema);
1601
1602 let file = tokio::fs::File::from(file.try_clone().unwrap());
1604 let builder = ParquetRecordBatchStreamBuilder::new(file).await.unwrap();
1605 let async_builder_schema = builder.schema().clone();
1606 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices);
1607 let mut reader = builder.with_projection(mask).build().unwrap();
1608 let async_reader_schema = reader.schema().clone();
1609 let batch = reader.next().await.unwrap().unwrap();
1610 let async_batch_schema = batch.schema();
1611 assert_schemas(
1612 async_builder_schema,
1613 async_reader_schema,
1614 async_batch_schema,
1615 );
1616 }
1617 }
1618
1619 #[tokio::test]
1620 async fn test_nested_skip() {
1621 let schema = Arc::new(Schema::new(vec![
1622 Field::new("col_1", DataType::UInt64, false),
1623 Field::new_list("col_2", Field::new_list_field(DataType::Utf8, true), true),
1624 ]));
1625
1626 let props = WriterProperties::builder()
1628 .set_data_page_row_count_limit(256)
1629 .set_write_batch_size(256)
1630 .set_max_row_group_row_count(Some(1024));
1631
1632 let mut file = tempfile().unwrap();
1634 let mut writer =
1635 ArrowWriter::try_new(&mut file, schema.clone(), Some(props.build())).unwrap();
1636
1637 let mut builder = ListBuilder::new(StringBuilder::new());
1638 for id in 0..1024 {
1639 match id % 3 {
1640 0 => builder.append_value([Some("val_1".to_string()), Some(format!("id_{id}"))]),
1641 1 => builder.append_value([Some(format!("id_{id}"))]),
1642 _ => builder.append_null(),
1643 }
1644 }
1645 let refs = vec![
1646 Arc::new(UInt64Array::from_iter_values(0..1024)) as ArrayRef,
1647 Arc::new(builder.finish()) as ArrayRef,
1648 ];
1649
1650 let batch = RecordBatch::try_new(schema.clone(), refs).unwrap();
1651 writer.write(&batch).unwrap();
1652 writer.close().unwrap();
1653
1654 let selections = [
1655 RowSelection::from(vec![
1656 RowSelector::skip(313),
1657 RowSelector::select(1),
1658 RowSelector::skip(709),
1659 RowSelector::select(1),
1660 ]),
1661 RowSelection::from(vec![
1662 RowSelector::skip(255),
1663 RowSelector::select(1),
1664 RowSelector::skip(767),
1665 RowSelector::select(1),
1666 ]),
1667 RowSelection::from(vec![
1668 RowSelector::select(255),
1669 RowSelector::skip(1),
1670 RowSelector::select(767),
1671 RowSelector::skip(1),
1672 ]),
1673 RowSelection::from(vec![
1674 RowSelector::skip(254),
1675 RowSelector::select(1),
1676 RowSelector::select(1),
1677 RowSelector::skip(767),
1678 RowSelector::select(1),
1679 ]),
1680 ];
1681
1682 for selection in selections {
1683 let expected = selection.row_count();
1684 let mut reader = ParquetRecordBatchStreamBuilder::new_with_options(
1686 tokio::fs::File::from_std(file.try_clone().unwrap()),
1687 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
1688 )
1689 .await
1690 .unwrap();
1691
1692 reader = reader.with_row_selection(selection);
1693
1694 let mut stream = reader.build().unwrap();
1695
1696 let mut total_rows = 0;
1697 while let Some(rb) = stream.next().await {
1698 let rb = rb.unwrap();
1699 total_rows += rb.num_rows();
1700 }
1701 assert_eq!(total_rows, expected);
1702 }
1703 }
1704
1705 #[tokio::test]
1706 async fn empty_offset_index_doesnt_panic_in_read_row_group() {
1707 use tokio::fs::File;
1708 let testdata = arrow::util::test_util::parquet_test_data();
1709 let path = format!("{testdata}/alltypes_plain.parquet");
1710 let mut file = File::open(&path).await.unwrap();
1711 let file_size = file.metadata().await.unwrap().len();
1712 let mut metadata = ParquetMetaDataReader::new()
1713 .with_page_index_policy(PageIndexPolicy::Required)
1714 .load_and_finish(&mut file, file_size)
1715 .await
1716 .unwrap();
1717
1718 metadata.set_offset_index(Some(vec![]));
1719 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1720 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1721 let reader =
1722 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1723 .build()
1724 .unwrap();
1725
1726 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1727 assert_eq!(result.len(), 1);
1728 }
1729
1730 #[tokio::test]
1731 async fn non_empty_offset_index_doesnt_panic_in_read_row_group() {
1732 use tokio::fs::File;
1733 let testdata = arrow::util::test_util::parquet_test_data();
1734 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1735 let mut file = File::open(&path).await.unwrap();
1736 let file_size = file.metadata().await.unwrap().len();
1737 let metadata = ParquetMetaDataReader::new()
1738 .with_page_index_policy(PageIndexPolicy::Required)
1739 .load_and_finish(&mut file, file_size)
1740 .await
1741 .unwrap();
1742
1743 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1744 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1745 let reader =
1746 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1747 .build()
1748 .unwrap();
1749
1750 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1751 assert_eq!(result.len(), 8);
1752 }
1753
1754 #[tokio::test]
1755 async fn empty_offset_index_doesnt_panic_in_column_chunks() {
1756 use tempfile::TempDir;
1757 use tokio::fs::File;
1758 fn write_metadata_to_local_file(
1759 metadata: ParquetMetaData,
1760 file: impl AsRef<std::path::Path>,
1761 ) {
1762 use crate::file::metadata::ParquetMetaDataWriter;
1763 use std::fs::File;
1764 let file = File::create(file).unwrap();
1765 ParquetMetaDataWriter::new(file, &metadata)
1766 .finish()
1767 .unwrap()
1768 }
1769
1770 fn read_metadata_from_local_file(file: impl AsRef<std::path::Path>) -> ParquetMetaData {
1771 use std::fs::File;
1772 let file = File::open(file).unwrap();
1773 ParquetMetaDataReader::new()
1774 .with_page_index_policy(PageIndexPolicy::Required)
1775 .parse_and_finish(&file)
1776 .unwrap()
1777 }
1778
1779 let testdata = arrow::util::test_util::parquet_test_data();
1780 let path = format!("{testdata}/alltypes_plain.parquet");
1781 let mut file = File::open(&path).await.unwrap();
1782 let file_size = file.metadata().await.unwrap().len();
1783 let metadata = ParquetMetaDataReader::new()
1784 .with_page_index_policy(PageIndexPolicy::Required)
1785 .load_and_finish(&mut file, file_size)
1786 .await
1787 .unwrap();
1788
1789 let tempdir = TempDir::new().unwrap();
1790 let metadata_path = tempdir.path().join("thrift_metadata.dat");
1791 write_metadata_to_local_file(metadata, &metadata_path);
1792 let metadata = read_metadata_from_local_file(&metadata_path);
1793
1794 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1795 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1796 let reader =
1797 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1798 .build()
1799 .unwrap();
1800
1801 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1803 assert_eq!(result.len(), 1);
1804 }
1805
1806 #[tokio::test]
1807 async fn test_cached_array_reader_sparse_offset_error() {
1808 use futures::TryStreamExt;
1809
1810 use crate::arrow::arrow_reader::{ArrowPredicateFn, RowFilter, RowSelection, RowSelector};
1811 use arrow_array::{BooleanArray, RecordBatch};
1812
1813 let testdata = arrow::util::test_util::parquet_test_data();
1814 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1815 let data = Bytes::from(std::fs::read(path).unwrap());
1816
1817 let async_reader = TestReader::new(data);
1818
1819 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1821 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1822 .await
1823 .unwrap();
1824
1825 let selection = RowSelection::from(vec![RowSelector::skip(22), RowSelector::select(3)]);
1829
1830 let parquet_schema = builder.parquet_schema();
1834 let proj = ProjectionMask::leaves(parquet_schema, vec![0]);
1835 let always_true = ArrowPredicateFn::new(proj.clone(), |batch: RecordBatch| {
1836 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
1837 });
1838 let filter = RowFilter::new(vec![Box::new(always_true)]);
1839
1840 let stream = builder
1843 .with_batch_size(8)
1844 .with_projection(proj)
1845 .with_row_selection(selection)
1846 .with_row_filter(filter)
1847 .build()
1848 .unwrap();
1849
1850 let _result: Vec<_> = stream.try_collect().await.unwrap();
1853 }
1854
1855 #[tokio::test]
1856 async fn test_predicate_cache_disabled() {
1857 let k = Int32Array::from_iter_values(0..10);
1858 let data = RecordBatch::try_from_iter([("k", Arc::new(k) as ArrayRef)]).unwrap();
1859
1860 let mut buf = Vec::new();
1861 let props = WriterProperties::builder()
1863 .set_data_page_row_count_limit(1)
1864 .set_write_batch_size(1)
1865 .set_max_row_group_row_count(Some(10))
1866 .set_write_page_header_statistics(true)
1867 .build();
1868 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1869 writer.write(&data).unwrap();
1870 writer.close().unwrap();
1871
1872 let data = Bytes::from(buf);
1873 let metadata = ParquetMetaDataReader::new()
1874 .with_page_index_policy(PageIndexPolicy::Required)
1875 .parse_and_finish(&data)
1876 .unwrap();
1877 let parquet_schema = metadata.file_metadata().schema_descr_ptr();
1878
1879 let build_filter = || {
1881 let scalar = Int32Array::from_iter_values([5]);
1882 let predicate = ArrowPredicateFn::new(
1883 ProjectionMask::leaves(&parquet_schema, vec![0]),
1884 move |batch| eq(batch.column(0), &Scalar::new(&scalar)),
1885 );
1886 RowFilter::new(vec![Box::new(predicate)])
1887 };
1888
1889 let selection = RowSelection::from(vec![RowSelector::skip(5), RowSelector::select(1)]);
1891
1892 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1893 let reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1894
1895 let reader_with_cache = TestReader::new(data.clone());
1897 let requests_with_cache = reader_with_cache.requests.clone();
1898 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1899 reader_with_cache,
1900 reader_metadata.clone(),
1901 )
1902 .with_batch_size(1000)
1903 .with_row_selection(selection.clone())
1904 .with_row_filter(build_filter())
1905 .build()
1906 .unwrap();
1907 let batches_with_cache: Vec<_> = stream.try_collect().await.unwrap();
1908
1909 let reader_without_cache = TestReader::new(data);
1911 let requests_without_cache = reader_without_cache.requests.clone();
1912 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1913 reader_without_cache,
1914 reader_metadata,
1915 )
1916 .with_batch_size(1000)
1917 .with_row_selection(selection)
1918 .with_row_filter(build_filter())
1919 .with_max_predicate_cache_size(0) .build()
1921 .unwrap();
1922 let batches_without_cache: Vec<_> = stream.try_collect().await.unwrap();
1923
1924 assert_eq!(batches_with_cache, batches_without_cache);
1925
1926 let requests_with_cache = requests_with_cache.lock().unwrap();
1927 let requests_without_cache = requests_without_cache.lock().unwrap();
1928
1929 assert_eq!(requests_with_cache.len(), 11);
1931 assert_eq!(requests_without_cache.len(), 2);
1932
1933 assert_eq!(
1935 requests_with_cache.iter().map(|r| r.len()).sum::<usize>(),
1936 433
1937 );
1938 assert_eq!(
1939 requests_without_cache
1940 .iter()
1941 .map(|r| r.len())
1942 .sum::<usize>(),
1943 92
1944 );
1945 }
1946
1947 #[test]
1948 fn test_row_numbers_with_multiple_row_groups() {
1949 test_row_numbers_with_multiple_row_groups_helper(
1950 false,
1951 |path, selection, _row_filter, batch_size| {
1952 let runtime = tokio::runtime::Builder::new_current_thread()
1953 .enable_all()
1954 .build()
1955 .expect("Could not create runtime");
1956 runtime.block_on(async move {
1957 let file = tokio::fs::File::open(path).await.unwrap();
1958 let row_number_field = Arc::new(
1959 Field::new("row_number", DataType::Int64, false)
1960 .with_extension_type(RowNumber),
1961 );
1962 let options = ArrowReaderOptions::new()
1963 .with_virtual_columns(vec![row_number_field])
1964 .unwrap();
1965 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
1966 .await
1967 .unwrap()
1968 .with_row_selection(selection)
1969 .with_batch_size(batch_size)
1970 .build()
1971 .expect("Could not create reader");
1972 reader.try_collect::<Vec<_>>().await.unwrap()
1973 })
1974 },
1975 );
1976 }
1977
1978 #[test]
1979 fn test_row_numbers_with_multiple_row_groups_and_filter() {
1980 test_row_numbers_with_multiple_row_groups_helper(
1981 true,
1982 |path, selection, row_filter, batch_size| {
1983 let runtime = tokio::runtime::Builder::new_current_thread()
1984 .enable_all()
1985 .build()
1986 .expect("Could not create runtime");
1987 runtime.block_on(async move {
1988 let file = tokio::fs::File::open(path).await.unwrap();
1989 let row_number_field = Arc::new(
1990 Field::new("row_number", DataType::Int64, false)
1991 .with_extension_type(RowNumber),
1992 );
1993 let options = ArrowReaderOptions::new()
1994 .with_virtual_columns(vec![row_number_field])
1995 .unwrap();
1996 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
1997 .await
1998 .unwrap()
1999 .with_row_selection(selection)
2000 .with_row_filter(row_filter.expect("No row filter"))
2001 .with_batch_size(batch_size)
2002 .build()
2003 .expect("Could not create reader");
2004 reader.try_collect::<Vec<_>>().await.unwrap()
2005 })
2006 },
2007 );
2008 }
2009
2010 #[tokio::test]
2011 async fn test_nested_lists() -> Result<()> {
2012 let list_inner_field = Arc::new(Field::new("item", DataType::Float32, true));
2014 let table_schema = Arc::new(Schema::new(vec![
2015 Field::new("id", DataType::Int32, false),
2016 Field::new("vector", DataType::List(list_inner_field.clone()), true),
2017 ]));
2018
2019 let mut list_builder =
2020 ListBuilder::new(Float32Builder::new()).with_field(list_inner_field.clone());
2021 list_builder.values().append_slice(&[10.0, 10.0, 10.0]);
2022 list_builder.append(true);
2023 list_builder.values().append_slice(&[20.0, 20.0, 20.0]);
2024 list_builder.append(true);
2025 list_builder.values().append_slice(&[30.0, 30.0, 30.0]);
2026 list_builder.append(true);
2027 list_builder.values().append_slice(&[40.0, 40.0, 40.0]);
2028 list_builder.append(true);
2029 let list_array = list_builder.finish();
2030
2031 let data = vec![RecordBatch::try_new(
2032 table_schema.clone(),
2033 vec![
2034 Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
2035 Arc::new(list_array),
2036 ],
2037 )?];
2038
2039 let mut buffer = Vec::new();
2040 let mut writer = AsyncArrowWriter::try_new(&mut buffer, table_schema, None)?;
2041
2042 for batch in data {
2043 writer.write(&batch).await?;
2044 }
2045
2046 writer.close().await?;
2047
2048 let reader = TestReader::new(Bytes::from(buffer));
2049 let builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
2050
2051 let predicate = ArrowPredicateFn::new(ProjectionMask::all(), |batch| {
2052 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
2053 });
2054
2055 let projection_mask = ProjectionMask::all();
2056
2057 let mut stream = builder
2058 .with_row_filter(RowFilter::new(vec![Box::new(predicate)]))
2059 .with_projection(projection_mask)
2060 .build()?;
2061
2062 while let Some(batch) = stream.next().await {
2063 let _ = batch.unwrap(); }
2065
2066 Ok(())
2067 }
2068}