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
56pub use crate::arrow::arrow_reader::RowGroupSelection;
59
60#[cfg(feature = "object_store")]
61mod store;
62
63use crate::DecodeResult;
64use crate::arrow::push_decoder::{ParquetPushDecoder, ParquetPushDecoderBuilder, PushDecoderInput};
65#[cfg(feature = "object_store")]
66pub use store::*;
67
68pub trait AsyncFileReader: Send {
166 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>>;
168
169 fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_, Result<Vec<Bytes>>> {
171 async move {
172 let mut result = Vec::with_capacity(ranges.len());
173
174 for range in ranges {
175 let data = self.get_bytes(range).await?;
176 result.push(data);
177 }
178
179 Ok(result)
180 }
181 .boxed()
182 }
183
184 fn get_metadata<'a>(
201 &'a mut self,
202 options: Option<&'a ArrowReaderOptions>,
203 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>>;
204}
205
206impl AsyncFileReader for Box<dyn AsyncFileReader + '_> {
208 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
209 self.as_mut().get_bytes(range)
210 }
211
212 fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_, Result<Vec<Bytes>>> {
213 self.as_mut().get_byte_ranges(ranges)
214 }
215
216 fn get_metadata<'a>(
217 &'a mut self,
218 options: Option<&'a ArrowReaderOptions>,
219 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
220 self.as_mut().get_metadata(options)
221 }
222}
223
224impl<T: AsyncFileReader + MetadataFetch + AsyncRead + AsyncSeek + Unpin> MetadataSuffixFetch for T {
225 fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
226 async move {
227 self.seek(SeekFrom::End(-(suffix as i64))).await?;
228 let mut buf = Vec::with_capacity(suffix);
229 self.take(suffix as _).read_to_end(&mut buf).await?;
230 Ok(buf.into())
231 }
232 .boxed()
233 }
234}
235
236impl<T: AsyncRead + AsyncSeek + Unpin + Send> AsyncFileReader for T {
237 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
238 async move {
239 self.seek(SeekFrom::Start(range.start)).await?;
240
241 let to_read = range.end - range.start;
242 let mut buffer = Vec::with_capacity(to_read.try_into()?);
243 let read = self.take(to_read).read_to_end(&mut buffer).await?;
244 if read as u64 != to_read {
245 return Err(eof_err!("expected to read {} bytes, got {}", to_read, read));
246 }
247
248 Ok(buffer.into())
249 }
250 .boxed()
251 }
252
253 fn get_metadata<'a>(
254 &'a mut self,
255 options: Option<&'a ArrowReaderOptions>,
256 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
257 async move {
258 let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
259 let parquet_metadata = metadata_reader.load_via_suffix_and_finish(self).await?;
260 Ok(Arc::new(parquet_metadata))
261 }
262 .boxed()
263 }
264}
265
266impl ArrowReaderMetadata {
267 pub async fn load_async<T: AsyncFileReader>(
271 input: &mut T,
272 options: ArrowReaderOptions,
273 ) -> Result<Self> {
274 let metadata = input.get_metadata(Some(&options)).await?;
275 Self::try_new(metadata, options)
276 }
277}
278
279#[doc(hidden)]
280pub struct AsyncReader<T>(T);
285
286pub type ParquetRecordBatchStreamBuilder<T> = ArrowReaderBuilder<AsyncReader<T>>;
306
307impl<T: AsyncFileReader + Send + 'static> ParquetRecordBatchStreamBuilder<T> {
308 pub async fn new(input: T) -> Result<Self> {
524 Self::new_with_options(input, Default::default()).await
525 }
526
527 pub async fn new_with_options(mut input: T, options: ArrowReaderOptions) -> Result<Self> {
530 let metadata = ArrowReaderMetadata::load_async(&mut input, options).await?;
531 Ok(Self::new_with_metadata(input, metadata))
532 }
533
534 pub fn new_with_metadata(input: T, metadata: ArrowReaderMetadata) -> Self {
580 Self::new_builder(AsyncReader(input), metadata)
581 }
582
583 pub async fn get_row_group_column_bloom_filter(
589 &mut self,
590 row_group_idx: usize,
591 column_idx: usize,
592 ) -> Result<Option<Sbbf>> {
593 let metadata = self.metadata.row_group(row_group_idx);
594 let column_metadata = metadata.column(column_idx);
595
596 let offset: u64 = if let Some(offset) = column_metadata.bloom_filter_offset() {
597 offset
598 .try_into()
599 .map_err(|_| ParquetError::General("Bloom filter offset is invalid".to_string()))?
600 } else {
601 return Ok(None);
602 };
603
604 let buffer = match column_metadata.bloom_filter_length() {
605 Some(length) => self.input.0.get_bytes(offset..offset + length as u64),
606 None => self
607 .input
608 .0
609 .get_bytes(offset..offset + SBBF_HEADER_SIZE_ESTIMATE as u64),
610 }
611 .await?;
612
613 let (header, bitset_offset) =
614 chunk_read_bloom_filter_header_and_offset(offset, buffer.clone())?;
615
616 match header.algorithm {
617 BloomFilterAlgorithm::BLOCK => {
618 }
620 }
621 match header.compression {
622 BloomFilterCompression::UNCOMPRESSED => {
623 }
625 }
626 match header.hash {
627 BloomFilterHash::XXHASH => {
628 }
630 }
631
632 let bitset = match column_metadata.bloom_filter_length() {
633 Some(_) => buffer.slice(
634 (TryInto::<usize>::try_into(bitset_offset).unwrap()
635 - TryInto::<usize>::try_into(offset).unwrap())..,
636 ),
637 None => {
638 let bitset_length: u64 = header.num_bytes.try_into().map_err(|_| {
639 ParquetError::General("Bloom filter length is invalid".to_string())
640 })?;
641 self.input
642 .0
643 .get_bytes(bitset_offset..bitset_offset + bitset_length)
644 .await?
645 }
646 };
647 Ok(Some(Sbbf::new(&bitset)))
648 }
649
650 pub fn with_row_group_selections(
665 mut self,
666 row_group_selections: Vec<RowGroupSelection>,
667 ) -> Self {
668 self.row_group_plan
669 .set_row_group_selections(row_group_selections);
670 self
671 }
672
673 pub fn build(self) -> Result<ParquetRecordBatchStream<T>> {
677 let Self {
678 input,
679 metadata,
680 schema,
681 fields,
682 batch_size,
683 row_group_plan,
684 projection,
685 filter,
686 row_selection_policy: selection_strategy,
687 limit,
688 offset,
689 metrics,
690 max_predicate_cache_size,
691 } = self;
692
693 let projection_len = projection.mask.as_ref().map_or(usize::MAX, |m| m.len());
696 let projected_fields = schema
697 .fields
698 .filter_leaves(|idx, _| idx < projection_len && projection.leaf_included(idx));
699 let projected_schema = Arc::new(Schema::new(projected_fields));
700
701 let decoder = ParquetPushDecoderBuilder {
702 input: PushDecoderInput::default(),
703 metadata,
704 schema,
705 fields,
706 projection,
707 filter,
708 row_group_plan,
709 row_selection_policy: selection_strategy,
710 batch_size,
711 limit,
712 offset,
713 metrics,
714 max_predicate_cache_size,
715 }
716 .build()?;
717
718 let request_state = RequestState::None { input: input.0 };
719
720 Ok(ParquetRecordBatchStream {
721 schema: projected_schema,
722 decoder,
723 request_state,
724 })
725 }
726}
727
728enum RequestState<T> {
732 None {
734 input: T,
735 },
736 Outstanding {
738 ranges: Vec<Range<u64>>,
740 future: BoxFuture<'static, Result<(T, Vec<Bytes>)>>,
745 },
746 Done,
747}
748
749impl<T> RequestState<T>
750where
751 T: AsyncFileReader + Unpin + Send + 'static,
752{
753 fn begin_request(mut input: T, ranges: Vec<Range<u64>>) -> Self {
755 let ranges_captured = ranges.clone();
756
757 let future = async move {
762 let data = input.get_byte_ranges(ranges_captured).await?;
763 Ok((input, data))
764 }
765 .boxed();
766 RequestState::Outstanding { ranges, future }
767 }
768}
769
770impl<T> std::fmt::Debug for RequestState<T> {
771 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
772 match self {
773 RequestState::None { input: _ } => f
774 .debug_struct("RequestState::None")
775 .field("input", &"...")
776 .finish(),
777 RequestState::Outstanding { ranges, .. } => f
778 .debug_struct("RequestState::Outstanding")
779 .field("ranges", &ranges)
780 .finish(),
781 RequestState::Done => {
782 write!(f, "RequestState::Done")
783 }
784 }
785 }
786}
787
788pub struct ParquetRecordBatchStream<T> {
808 schema: SchemaRef,
810 request_state: RequestState<T>,
812 decoder: ParquetPushDecoder,
814}
815
816impl<T> std::fmt::Debug for ParquetRecordBatchStream<T> {
817 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
818 f.debug_struct("ParquetRecordBatchStream")
819 .field("request_state", &self.request_state)
820 .finish()
821 }
822}
823
824impl<T> ParquetRecordBatchStream<T> {
825 pub fn schema(&self) -> &SchemaRef {
830 &self.schema
831 }
832}
833
834impl<T> ParquetRecordBatchStream<T>
835where
836 T: AsyncFileReader + Unpin + Send + 'static,
837{
838 pub async fn next_row_group(&mut self) -> Result<Option<ParquetRecordBatchReader>> {
852 loop {
853 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
856 match request_state {
857 RequestState::None { input } => {
859 match self.decoder.try_next_reader()? {
860 DecodeResult::NeedsData(ranges) => {
861 self.request_state = RequestState::begin_request(input, ranges);
862 }
864 DecodeResult::Data(reader) => {
865 self.request_state = RequestState::None { input };
866 return Ok(Some(reader));
867 }
868 DecodeResult::Finished => return Ok(None),
869 }
870 }
871 RequestState::Outstanding { ranges, future } => {
872 let (input, data) = future.await?;
873 self.decoder.push_ranges(ranges, data)?;
875 self.request_state = RequestState::None { input };
876 }
878 RequestState::Done => {
879 self.request_state = RequestState::Done;
880 return Ok(None);
881 }
882 }
883 }
884 }
885}
886
887impl<T> Stream for ParquetRecordBatchStream<T>
888where
889 T: AsyncFileReader + Unpin + Send + 'static,
890{
891 type Item = Result<RecordBatch>;
892 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
893 match self.poll_next_inner(cx) {
894 Ok(res) => {
895 res.map(|res| Ok(res).transpose())
898 }
899 Err(e) => {
900 self.request_state = RequestState::Done;
901 Poll::Ready(Some(Err(e)))
902 }
903 }
904 }
905}
906
907impl<T> ParquetRecordBatchStream<T>
908where
909 T: AsyncFileReader + Unpin + Send + 'static,
910{
911 fn poll_next_inner(&mut self, cx: &mut Context<'_>) -> Result<Poll<Option<RecordBatch>>> {
916 loop {
917 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
918 match request_state {
919 RequestState::None { input } => {
920 match self.decoder.try_decode()? {
922 DecodeResult::NeedsData(ranges) => {
923 self.request_state = RequestState::begin_request(input, ranges);
924 }
926 DecodeResult::Data(batch) => {
927 self.request_state = RequestState::None { input };
928 return Ok(Poll::Ready(Some(batch)));
929 }
930 DecodeResult::Finished => {
931 self.request_state = RequestState::Done;
932 return Ok(Poll::Ready(None));
933 }
934 }
935 }
936 RequestState::Outstanding { ranges, mut future } => match future.poll_unpin(cx) {
937 Poll::Ready(result) => {
939 let (input, data) = result?;
940 self.decoder.push_ranges(ranges, data)?;
942 self.request_state = RequestState::None { input };
943 }
945 Poll::Pending => {
946 self.request_state = RequestState::Outstanding { ranges, future };
947 return Ok(Poll::Pending);
948 }
949 },
950 RequestState::Done => {
951 self.request_state = RequestState::Done;
953 return Ok(Poll::Ready(None));
954 }
955 }
956 }
957 }
958}
959
960#[cfg(test)]
961mod tests {
962 use super::*;
963 use crate::arrow::arrow_reader::tests::test_row_numbers_with_multiple_row_groups_helper;
964 use crate::arrow::arrow_reader::{
965 ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowFilter, RowSelection, RowSelector,
966 };
967 use crate::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
968 use crate::arrow::schema::virtual_type::RowNumber;
969 use crate::arrow::{ArrowWriter, AsyncArrowWriter, ProjectionMask};
970 use crate::file::metadata::ParquetMetaDataReader;
971 use crate::file::metadata::{PageIndex, PageIndexPolicy};
972 use crate::file::properties::WriterProperties;
973 use arrow::compute::kernels::cmp::eq;
974 use arrow::error::Result as ArrowResult;
975 use arrow_array::builder::{Float32Builder, ListBuilder, StringBuilder};
976 use arrow_array::cast::AsArray;
977 use arrow_array::types::Int32Type;
978 use arrow_array::{
979 Array, ArrayRef, BooleanArray, Int32Array, RecordBatchReader, Scalar, StringArray,
980 StructArray, UInt64Array,
981 };
982 use arrow_schema::{DataType, Field, Schema};
983 use futures::{StreamExt, TryStreamExt};
984 use rand::{RngExt, rng};
985 use std::collections::HashMap;
986 use std::sync::{Arc, Mutex};
987 use tempfile::tempfile;
988
989 #[derive(Clone)]
990 struct TestReader {
991 data: Bytes,
992 metadata: Option<Arc<ParquetMetaData>>,
993 requests: Arc<Mutex<Vec<Range<usize>>>>,
994 }
995
996 impl TestReader {
997 fn new(data: Bytes) -> Self {
998 Self {
999 data,
1000 metadata: Default::default(),
1001 requests: Default::default(),
1002 }
1003 }
1004 }
1005
1006 impl AsyncFileReader for TestReader {
1007 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1008 let range = range.clone();
1009 self.requests
1010 .lock()
1011 .unwrap()
1012 .push(range.start as usize..range.end as usize);
1013 futures::future::ready(Ok(self
1014 .data
1015 .slice(range.start as usize..range.end as usize)))
1016 .boxed()
1017 }
1018
1019 fn get_metadata<'a>(
1020 &'a mut self,
1021 options: Option<&'a ArrowReaderOptions>,
1022 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
1023 let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
1024 self.metadata = Some(Arc::new(
1025 metadata_reader.parse_and_finish(&self.data).unwrap(),
1026 ));
1027 futures::future::ready(Ok(self.metadata.clone().unwrap().clone())).boxed()
1028 }
1029 }
1030
1031 #[tokio::test]
1032 async fn test_async_reader() {
1033 let testdata = arrow::util::test_util::parquet_test_data();
1034 let path = format!("{testdata}/alltypes_plain.parquet");
1035 let data = Bytes::from(std::fs::read(path).unwrap());
1036
1037 let async_reader = TestReader::new(data.clone());
1038
1039 let requests = async_reader.requests.clone();
1040 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1041 .await
1042 .unwrap();
1043
1044 let metadata = builder.metadata().clone();
1045 assert_eq!(metadata.num_row_groups(), 1);
1046
1047 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1048 let stream = builder
1049 .with_projection(mask.clone())
1050 .with_batch_size(1024)
1051 .build()
1052 .unwrap();
1053
1054 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1055
1056 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1057 .unwrap()
1058 .with_projection(mask)
1059 .with_batch_size(104)
1060 .build()
1061 .unwrap()
1062 .collect::<ArrowResult<Vec<_>>>()
1063 .unwrap();
1064
1065 assert_eq!(async_batches, sync_batches);
1066
1067 let requests = requests.lock().unwrap();
1068 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1069 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1070
1071 assert_eq!(
1072 &requests[..],
1073 &[
1074 offset_1 as usize..(offset_1 + length_1) as usize,
1075 offset_2 as usize..(offset_2 + length_2) as usize
1076 ]
1077 );
1078 }
1079
1080 #[tokio::test]
1081 async fn test_async_reader_row_group_local_selections() {
1082 let batch = RecordBatch::try_from_iter([(
1083 "a",
1084 Arc::new(Int32Array::from_iter_values(0..6)) as ArrayRef,
1085 )])
1086 .unwrap();
1087 let mut data = Vec::new();
1088 let properties = WriterProperties::builder()
1089 .set_max_row_group_row_count(Some(3))
1090 .build();
1091 let mut writer = ArrowWriter::try_new(&mut data, batch.schema(), Some(properties)).unwrap();
1092 writer.write(&batch).unwrap();
1093 writer.close().unwrap();
1094
1095 let stream = ParquetRecordBatchStreamBuilder::new(TestReader::new(data.into()))
1096 .await
1097 .unwrap()
1098 .with_row_group_selections(vec![
1099 RowGroupSelection::new(1, Some(RowSelection::from(vec![RowSelector::select(1)]))),
1100 RowGroupSelection::new(
1101 0,
1102 Some(RowSelection::from(vec![
1103 RowSelector::skip(1),
1104 RowSelector::select(2),
1105 ])),
1106 ),
1107 ])
1108 .build()
1109 .unwrap();
1110
1111 let batches: Vec<_> = stream.try_collect().await.unwrap();
1112 assert_eq!(batches.len(), 2);
1113 assert_eq!(
1114 batches[0].column(0).as_primitive::<Int32Type>().values(),
1115 &[3]
1116 );
1117 assert_eq!(
1118 batches[1].column(0).as_primitive::<Int32Type>().values(),
1119 &[1, 2]
1120 );
1121 }
1122
1123 #[tokio::test]
1124 async fn test_async_reader_with_next_row_group() {
1125 let testdata = arrow::util::test_util::parquet_test_data();
1126 let path = format!("{testdata}/alltypes_plain.parquet");
1127 let data = Bytes::from(std::fs::read(path).unwrap());
1128
1129 let async_reader = TestReader::new(data.clone());
1130
1131 let requests = async_reader.requests.clone();
1132 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1133 .await
1134 .unwrap();
1135
1136 let metadata = builder.metadata().clone();
1137 assert_eq!(metadata.num_row_groups(), 1);
1138
1139 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1140 let mut stream = builder
1141 .with_projection(mask.clone())
1142 .with_batch_size(1024)
1143 .build()
1144 .unwrap();
1145
1146 let mut readers = vec![];
1147 while let Some(reader) = stream.next_row_group().await.unwrap() {
1148 readers.push(reader);
1149 }
1150
1151 let async_batches: Vec<_> = readers
1152 .into_iter()
1153 .flat_map(|r| r.map(|v| v.unwrap()).collect::<Vec<_>>())
1154 .collect();
1155
1156 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1157 .unwrap()
1158 .with_projection(mask)
1159 .with_batch_size(104)
1160 .build()
1161 .unwrap()
1162 .collect::<ArrowResult<Vec<_>>>()
1163 .unwrap();
1164
1165 assert_eq!(async_batches, sync_batches);
1166
1167 let requests = requests.lock().unwrap();
1168 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1169 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1170
1171 assert_eq!(
1172 &requests[..],
1173 &[
1174 offset_1 as usize..(offset_1 + length_1) as usize,
1175 offset_2 as usize..(offset_2 + length_2) as usize
1176 ]
1177 );
1178 }
1179
1180 #[tokio::test]
1181 async fn test_async_reader_with_index() {
1182 let testdata = arrow::util::test_util::parquet_test_data();
1183 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1184 let data = Bytes::from(std::fs::read(path).unwrap());
1185
1186 let async_reader = TestReader::new(data.clone());
1187
1188 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1189 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1190 .await
1191 .unwrap();
1192
1193 let metadata_with_index = builder.metadata();
1195 assert_eq!(metadata_with_index.num_row_groups(), 1);
1196
1197 let page_index = metadata_with_index.page_index().unwrap();
1199 let num_rowgroups = metadata_with_index.num_row_groups();
1200 let num_columns = metadata_with_index
1201 .file_metadata()
1202 .schema_descr()
1203 .num_columns();
1204 for rgidx in 0..num_rowgroups {
1205 let column_index = page_index.column_indexes_for_rowgroup(rgidx);
1206 let offset_index = page_index.offset_indexes_for_rowgroup(rgidx);
1207 assert!(column_index.is_some_and(|ci| ci.len() == num_columns));
1208 assert!(offset_index.is_some_and(|oi| oi.len() == num_columns));
1209 for colidx in 0..num_columns {
1211 assert!(page_index.offset_index(rgidx, colidx).is_some());
1212 }
1213 }
1214
1215 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1216 let stream = builder
1217 .with_projection(mask.clone())
1218 .with_batch_size(1024)
1219 .build()
1220 .unwrap();
1221
1222 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1223
1224 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1225 .unwrap()
1226 .with_projection(mask)
1227 .with_batch_size(1024)
1228 .build()
1229 .unwrap()
1230 .collect::<ArrowResult<Vec<_>>>()
1231 .unwrap();
1232
1233 assert_eq!(async_batches, sync_batches);
1234 }
1235
1236 #[tokio::test]
1237 async fn test_async_reader_with_limit() {
1238 let testdata = arrow::util::test_util::parquet_test_data();
1239 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1240 let data = Bytes::from(std::fs::read(path).unwrap());
1241
1242 let metadata = ParquetMetaDataReader::new()
1243 .parse_and_finish(&data)
1244 .unwrap();
1245 let metadata = Arc::new(metadata);
1246
1247 assert_eq!(metadata.num_row_groups(), 1);
1248
1249 let async_reader = TestReader::new(data.clone());
1250
1251 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1252 .await
1253 .unwrap();
1254
1255 assert_eq!(builder.metadata().num_row_groups(), 1);
1256
1257 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1258 let stream = builder
1259 .with_projection(mask.clone())
1260 .with_batch_size(1024)
1261 .with_limit(1)
1262 .build()
1263 .unwrap();
1264
1265 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1266
1267 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1268 .unwrap()
1269 .with_projection(mask)
1270 .with_batch_size(1024)
1271 .with_limit(1)
1272 .build()
1273 .unwrap()
1274 .collect::<ArrowResult<Vec<_>>>()
1275 .unwrap();
1276
1277 assert_eq!(async_batches, sync_batches);
1278 }
1279
1280 #[tokio::test]
1281 async fn test_async_reader_skip_pages() {
1282 let testdata = arrow::util::test_util::parquet_test_data();
1283 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1284 let data = Bytes::from(std::fs::read(path).unwrap());
1285
1286 let async_reader = TestReader::new(data.clone());
1287
1288 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1289 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1290 .await
1291 .unwrap();
1292
1293 assert_eq!(builder.metadata().num_row_groups(), 1);
1294
1295 let selection = RowSelection::from(vec![
1296 RowSelector::skip(21), RowSelector::select(21), RowSelector::skip(41), RowSelector::select(41), RowSelector::skip(25), RowSelector::select(25), RowSelector::skip(7116), RowSelector::select(10), ]);
1305
1306 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![9]);
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 sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1317 .unwrap()
1318 .with_projection(mask)
1319 .with_batch_size(1024)
1320 .with_row_selection(selection)
1321 .build()
1322 .unwrap()
1323 .collect::<ArrowResult<Vec<_>>>()
1324 .unwrap();
1325
1326 assert_eq!(async_batches, sync_batches);
1327 }
1328
1329 #[tokio::test]
1330 async fn test_fuzz_async_reader_selection() {
1331 let testdata = arrow::util::test_util::parquet_test_data();
1332 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1333 let data = Bytes::from(std::fs::read(path).unwrap());
1334
1335 let mut rand = rng();
1336
1337 for _ in 0..100 {
1338 let mut expected_rows = 0;
1339 let mut total_rows = 0;
1340 let mut skip = false;
1341 let mut selectors = vec![];
1342
1343 while total_rows < 7300 {
1344 let row_count: usize = rand.random_range(1..100);
1345
1346 let row_count = row_count.min(7300 - total_rows);
1347
1348 selectors.push(RowSelector { row_count, skip });
1349
1350 total_rows += row_count;
1351 if !skip {
1352 expected_rows += row_count;
1353 }
1354
1355 skip = !skip;
1356 }
1357
1358 let selection = RowSelection::from(selectors);
1359
1360 let async_reader = TestReader::new(data.clone());
1361
1362 let options =
1363 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1364 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1365 .await
1366 .unwrap();
1367
1368 assert_eq!(builder.metadata().num_row_groups(), 1);
1369
1370 let col_idx: usize = rand.random_range(0..13);
1371 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1372
1373 let stream = builder
1374 .with_projection(mask.clone())
1375 .with_row_selection(selection.clone())
1376 .build()
1377 .expect("building stream");
1378
1379 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1380
1381 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1382
1383 assert_eq!(actual_rows, expected_rows);
1384 }
1385 }
1386
1387 #[tokio::test]
1388 async fn test_async_reader_zero_row_selector() {
1389 let testdata = arrow::util::test_util::parquet_test_data();
1391 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1392 let data = Bytes::from(std::fs::read(path).unwrap());
1393
1394 let mut rand = rng();
1395
1396 let mut expected_rows = 0;
1397 let mut total_rows = 0;
1398 let mut skip = false;
1399 let mut selectors = vec![];
1400
1401 selectors.push(RowSelector {
1402 row_count: 0,
1403 skip: false,
1404 });
1405
1406 while total_rows < 7300 {
1407 let row_count: usize = rand.random_range(1..100);
1408
1409 let row_count = row_count.min(7300 - total_rows);
1410
1411 selectors.push(RowSelector { row_count, skip });
1412
1413 total_rows += row_count;
1414 if !skip {
1415 expected_rows += row_count;
1416 }
1417
1418 skip = !skip;
1419 }
1420
1421 let selection = RowSelection::from(selectors);
1422
1423 let async_reader = TestReader::new(data.clone());
1424
1425 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1426 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1427 .await
1428 .unwrap();
1429
1430 assert_eq!(builder.metadata().num_row_groups(), 1);
1431
1432 let col_idx: usize = rand.random_range(0..13);
1433 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1434
1435 let stream = builder
1436 .with_projection(mask.clone())
1437 .with_row_selection(selection.clone())
1438 .build()
1439 .expect("building stream");
1440
1441 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1442
1443 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1444
1445 assert_eq!(actual_rows, expected_rows);
1446 }
1447
1448 #[tokio::test]
1449 async fn test_limit_multiple_row_groups() {
1450 let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]);
1451 let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]);
1452 let c = Int32Array::from_iter(0..6);
1453 let data = RecordBatch::try_from_iter([
1454 ("a", Arc::new(a) as ArrayRef),
1455 ("b", Arc::new(b) as ArrayRef),
1456 ("c", Arc::new(c) as ArrayRef),
1457 ])
1458 .unwrap();
1459
1460 let mut buf = Vec::with_capacity(1024);
1461 let props = WriterProperties::builder()
1462 .set_max_row_group_row_count(Some(3))
1463 .build();
1464 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1465 writer.write(&data).unwrap();
1466 writer.close().unwrap();
1467
1468 let data: Bytes = buf.into();
1469 let metadata = ParquetMetaDataReader::new()
1470 .parse_and_finish(&data)
1471 .unwrap();
1472
1473 assert_eq!(metadata.num_row_groups(), 2);
1474
1475 let test = TestReader::new(data);
1476
1477 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1478 .await
1479 .unwrap()
1480 .with_batch_size(1024)
1481 .with_limit(4)
1482 .build()
1483 .unwrap();
1484
1485 let batches: Vec<_> = stream.try_collect().await.unwrap();
1486 assert_eq!(batches.len(), 2);
1488
1489 let batch = &batches[0];
1490 assert_eq!(batch.num_rows(), 3);
1492 assert_eq!(batch.num_columns(), 3);
1493 let col2 = batch.column(2).as_primitive::<Int32Type>();
1494 assert_eq!(col2.values(), &[0, 1, 2]);
1495
1496 let batch = &batches[1];
1497 assert_eq!(batch.num_rows(), 1);
1499 assert_eq!(batch.num_columns(), 3);
1500 let col2 = batch.column(2).as_primitive::<Int32Type>();
1501 assert_eq!(col2.values(), &[3]);
1502
1503 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1504 .await
1505 .unwrap()
1506 .with_offset(2)
1507 .with_limit(3)
1508 .build()
1509 .unwrap();
1510
1511 let batches: Vec<_> = stream.try_collect().await.unwrap();
1512 assert_eq!(batches.len(), 2);
1514
1515 let batch = &batches[0];
1516 assert_eq!(batch.num_rows(), 1);
1518 assert_eq!(batch.num_columns(), 3);
1519 let col2 = batch.column(2).as_primitive::<Int32Type>();
1520 assert_eq!(col2.values(), &[2]);
1521
1522 let batch = &batches[1];
1523 assert_eq!(batch.num_rows(), 2);
1525 assert_eq!(batch.num_columns(), 3);
1526 let col2 = batch.column(2).as_primitive::<Int32Type>();
1527 assert_eq!(col2.values(), &[3, 4]);
1528
1529 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1530 .await
1531 .unwrap()
1532 .with_offset(4)
1533 .with_limit(20)
1534 .build()
1535 .unwrap();
1536
1537 let batches: Vec<_> = stream.try_collect().await.unwrap();
1538 assert_eq!(batches.len(), 1);
1540
1541 let batch = &batches[0];
1542 assert_eq!(batch.num_rows(), 2);
1544 assert_eq!(batch.num_columns(), 3);
1545 let col2 = batch.column(2).as_primitive::<Int32Type>();
1546 assert_eq!(col2.values(), &[4, 5]);
1547 }
1548
1549 #[tokio::test]
1550 async fn test_batch_size_overallocate() {
1551 let testdata = arrow::util::test_util::parquet_test_data();
1552 let path = format!("{testdata}/alltypes_plain.parquet");
1554 let data = Bytes::from(std::fs::read(path).unwrap());
1555
1556 let async_reader = TestReader::new(data.clone());
1557
1558 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1559 .await
1560 .unwrap();
1561
1562 let file_rows = builder.metadata().file_metadata().num_rows() as usize;
1563
1564 let builder = builder
1565 .with_projection(ProjectionMask::all())
1566 .with_batch_size(1024);
1567
1568 assert_ne!(1024, file_rows);
1571 assert_eq!(builder.batch_size, file_rows);
1572
1573 let _stream = builder.build().unwrap();
1574 }
1575
1576 #[tokio::test]
1577 async fn test_parquet_record_batch_stream_schema() {
1578 fn get_all_field_names(schema: &Schema) -> Vec<&String> {
1579 schema.flattened_fields().iter().map(|f| f.name()).collect()
1580 }
1581
1582 let mut metadata = HashMap::with_capacity(1);
1591 metadata.insert("key".to_string(), "value".to_string());
1592
1593 let nested_struct_array = StructArray::from(vec![
1594 (
1595 Arc::new(Field::new("d", DataType::Utf8, true)),
1596 Arc::new(StringArray::from(vec!["a", "b"])) as ArrayRef,
1597 ),
1598 (
1599 Arc::new(Field::new("e", DataType::Utf8, true)),
1600 Arc::new(StringArray::from(vec!["c", "d"])) as ArrayRef,
1601 ),
1602 ]);
1603 let struct_array = StructArray::from(vec![
1604 (
1605 Arc::new(Field::new("a", DataType::Int32, true)),
1606 Arc::new(Int32Array::from(vec![-1, 1])) as ArrayRef,
1607 ),
1608 (
1609 Arc::new(Field::new("b", DataType::UInt64, true)),
1610 Arc::new(UInt64Array::from(vec![1, 2])) as ArrayRef,
1611 ),
1612 (
1613 Arc::new(Field::new(
1614 "c",
1615 nested_struct_array.data_type().clone(),
1616 true,
1617 )),
1618 Arc::new(nested_struct_array) as ArrayRef,
1619 ),
1620 ]);
1621
1622 let schema =
1623 Arc::new(Schema::new(struct_array.fields().clone()).with_metadata(metadata.clone()));
1624 let record_batch = RecordBatch::from(struct_array)
1625 .with_schema(schema.clone())
1626 .unwrap();
1627
1628 let mut file = tempfile().unwrap();
1630 let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1631 writer.write(&record_batch).unwrap();
1632 writer.close().unwrap();
1633
1634 let all_fields = ["a", "b", "c", "d", "e"];
1635 let projections = [
1637 (vec![], vec![]),
1638 (vec![0], vec!["a"]),
1639 (vec![0, 1], vec!["a", "b"]),
1640 (vec![0, 1, 2], vec!["a", "b", "c", "d"]),
1641 (vec![0, 1, 2, 3], vec!["a", "b", "c", "d", "e"]),
1642 ];
1643
1644 for (indices, expected_projected_names) in projections {
1646 let assert_schemas = |builder: SchemaRef, reader: SchemaRef, batch: SchemaRef| {
1647 assert_eq!(get_all_field_names(&builder), all_fields);
1649 assert_eq!(builder.metadata, metadata);
1650 assert_eq!(get_all_field_names(&reader), expected_projected_names);
1652 assert_eq!(reader.metadata, HashMap::default());
1653 assert_eq!(get_all_field_names(&batch), expected_projected_names);
1654 assert_eq!(batch.metadata, HashMap::default());
1655 };
1656
1657 let builder =
1658 ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1659 let sync_builder_schema = builder.schema().clone();
1660 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices.clone());
1661 let mut reader = builder.with_projection(mask).build().unwrap();
1662 let sync_reader_schema = reader.schema();
1663 let batch = reader.next().unwrap().unwrap();
1664 let sync_batch_schema = batch.schema();
1665 assert_schemas(sync_builder_schema, sync_reader_schema, sync_batch_schema);
1666
1667 let file = tokio::fs::File::from(file.try_clone().unwrap());
1669 let builder = ParquetRecordBatchStreamBuilder::new(file).await.unwrap();
1670 let async_builder_schema = builder.schema().clone();
1671 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices);
1672 let mut reader = builder.with_projection(mask).build().unwrap();
1673 let async_reader_schema = reader.schema().clone();
1674 let batch = reader.next().await.unwrap().unwrap();
1675 let async_batch_schema = batch.schema();
1676 assert_schemas(
1677 async_builder_schema,
1678 async_reader_schema,
1679 async_batch_schema,
1680 );
1681 }
1682 }
1683
1684 #[tokio::test]
1685 async fn test_nested_skip() {
1686 let schema = Arc::new(Schema::new(vec![
1687 Field::new("col_1", DataType::UInt64, false),
1688 Field::new_list("col_2", Field::new_list_field(DataType::Utf8, true), true),
1689 ]));
1690
1691 let props = WriterProperties::builder()
1693 .set_data_page_row_count_limit(256)
1694 .set_write_batch_size(256)
1695 .set_max_row_group_row_count(Some(1024));
1696
1697 let mut file = tempfile().unwrap();
1699 let mut writer =
1700 ArrowWriter::try_new(&mut file, schema.clone(), Some(props.build())).unwrap();
1701
1702 let mut builder = ListBuilder::new(StringBuilder::new());
1703 for id in 0..1024 {
1704 match id % 3 {
1705 0 => builder.append_value([Some("val_1".to_string()), Some(format!("id_{id}"))]),
1706 1 => builder.append_value([Some(format!("id_{id}"))]),
1707 _ => builder.append_null(),
1708 }
1709 }
1710 let refs = vec![
1711 Arc::new(UInt64Array::from_iter_values(0..1024)) as ArrayRef,
1712 Arc::new(builder.finish()) as ArrayRef,
1713 ];
1714
1715 let batch = RecordBatch::try_new(schema.clone(), refs).unwrap();
1716 writer.write(&batch).unwrap();
1717 writer.close().unwrap();
1718
1719 let selections = [
1720 RowSelection::from(vec![
1721 RowSelector::skip(313),
1722 RowSelector::select(1),
1723 RowSelector::skip(709),
1724 RowSelector::select(1),
1725 ]),
1726 RowSelection::from(vec![
1727 RowSelector::skip(255),
1728 RowSelector::select(1),
1729 RowSelector::skip(767),
1730 RowSelector::select(1),
1731 ]),
1732 RowSelection::from(vec![
1733 RowSelector::select(255),
1734 RowSelector::skip(1),
1735 RowSelector::select(767),
1736 RowSelector::skip(1),
1737 ]),
1738 RowSelection::from(vec![
1739 RowSelector::skip(254),
1740 RowSelector::select(1),
1741 RowSelector::select(1),
1742 RowSelector::skip(767),
1743 RowSelector::select(1),
1744 ]),
1745 ];
1746
1747 for selection in selections {
1748 let expected = selection.row_count();
1749 let mut reader = ParquetRecordBatchStreamBuilder::new_with_options(
1751 tokio::fs::File::from_std(file.try_clone().unwrap()),
1752 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
1753 )
1754 .await
1755 .unwrap();
1756
1757 reader = reader.with_row_selection(selection);
1758
1759 let mut stream = reader.build().unwrap();
1760
1761 let mut total_rows = 0;
1762 while let Some(rb) = stream.next().await {
1763 let rb = rb.unwrap();
1764 total_rows += rb.num_rows();
1765 }
1766 assert_eq!(total_rows, expected);
1767 }
1768 }
1769
1770 #[tokio::test]
1771 async fn empty_offset_index_doesnt_panic_in_read_row_group() {
1772 use tokio::fs::File;
1773 let testdata = arrow::util::test_util::parquet_test_data();
1774 let path = format!("{testdata}/alltypes_plain.parquet");
1775 let mut file = File::open(&path).await.unwrap();
1776 let file_size = file.metadata().await.unwrap().len();
1777 let mut metadata = ParquetMetaDataReader::new()
1778 .with_page_index_policy(PageIndexPolicy::Required)
1779 .load_and_finish(&mut file, file_size)
1780 .await
1781 .unwrap();
1782
1783 let page_index = PageIndex::new(None, Some(vec![]));
1784 metadata.set_page_index(Some(page_index));
1785 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1786 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1787 let reader =
1788 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1789 .build()
1790 .unwrap();
1791
1792 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1793 assert_eq!(result.len(), 1);
1794 }
1795
1796 #[tokio::test]
1797 async fn non_empty_offset_index_doesnt_panic_in_read_row_group() {
1798 use tokio::fs::File;
1799 let testdata = arrow::util::test_util::parquet_test_data();
1800 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1801 let mut file = File::open(&path).await.unwrap();
1802 let file_size = file.metadata().await.unwrap().len();
1803 let metadata = ParquetMetaDataReader::new()
1804 .with_page_index_policy(PageIndexPolicy::Required)
1805 .load_and_finish(&mut file, file_size)
1806 .await
1807 .unwrap();
1808
1809 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1810 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1811 let reader =
1812 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1813 .build()
1814 .unwrap();
1815
1816 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1817 assert_eq!(result.len(), 8);
1818 }
1819
1820 #[tokio::test]
1821 async fn empty_offset_index_doesnt_panic_in_column_chunks() {
1822 use tempfile::TempDir;
1823 use tokio::fs::File;
1824 fn write_metadata_to_local_file(
1825 metadata: ParquetMetaData,
1826 file: impl AsRef<std::path::Path>,
1827 ) {
1828 use crate::file::metadata::ParquetMetaDataWriter;
1829 use std::fs::File;
1830 let file = File::create(file).unwrap();
1831 ParquetMetaDataWriter::new(file, &metadata)
1832 .finish()
1833 .unwrap()
1834 }
1835
1836 fn read_metadata_from_local_file(file: impl AsRef<std::path::Path>) -> ParquetMetaData {
1837 use std::fs::File;
1838 let file = File::open(file).unwrap();
1839 ParquetMetaDataReader::new()
1840 .with_page_index_policy(PageIndexPolicy::Required)
1841 .parse_and_finish(&file)
1842 .unwrap()
1843 }
1844
1845 let testdata = arrow::util::test_util::parquet_test_data();
1846 let path = format!("{testdata}/alltypes_plain.parquet");
1847 let mut file = File::open(&path).await.unwrap();
1848 let file_size = file.metadata().await.unwrap().len();
1849 let metadata = ParquetMetaDataReader::new()
1850 .with_page_index_policy(PageIndexPolicy::Required)
1851 .load_and_finish(&mut file, file_size)
1852 .await
1853 .unwrap();
1854
1855 let tempdir = TempDir::new().unwrap();
1856 let metadata_path = tempdir.path().join("thrift_metadata.dat");
1857 write_metadata_to_local_file(metadata, &metadata_path);
1858 let metadata = read_metadata_from_local_file(&metadata_path);
1859
1860 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1861 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1862 let reader =
1863 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1864 .build()
1865 .unwrap();
1866
1867 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1869 assert_eq!(result.len(), 1);
1870 }
1871
1872 #[tokio::test]
1873 async fn test_cached_array_reader_sparse_offset_error() {
1874 use futures::TryStreamExt;
1875
1876 use crate::arrow::arrow_reader::{ArrowPredicateFn, RowFilter, RowSelection, RowSelector};
1877 use arrow_array::{BooleanArray, RecordBatch};
1878
1879 let testdata = arrow::util::test_util::parquet_test_data();
1880 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1881 let data = Bytes::from(std::fs::read(path).unwrap());
1882
1883 let async_reader = TestReader::new(data);
1884
1885 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1887 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1888 .await
1889 .unwrap();
1890
1891 let selection = RowSelection::from(vec![RowSelector::skip(22), RowSelector::select(3)]);
1895
1896 let parquet_schema = builder.parquet_schema();
1900 let proj = ProjectionMask::leaves(parquet_schema, vec![0]);
1901 let always_true = ArrowPredicateFn::new(proj.clone(), |batch: RecordBatch| {
1902 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
1903 });
1904 let filter = RowFilter::new(vec![Box::new(always_true)]);
1905
1906 let stream = builder
1909 .with_batch_size(8)
1910 .with_projection(proj)
1911 .with_row_selection(selection)
1912 .with_row_filter(filter)
1913 .build()
1914 .unwrap();
1915
1916 let _result: Vec<_> = stream.try_collect().await.unwrap();
1919 }
1920
1921 #[tokio::test]
1922 async fn test_predicate_cache_disabled() {
1923 let k = Int32Array::from_iter_values(0..10);
1924 let data = RecordBatch::try_from_iter([("k", Arc::new(k) as ArrayRef)]).unwrap();
1925
1926 let mut buf = Vec::new();
1927 let props = WriterProperties::builder()
1929 .set_data_page_row_count_limit(1)
1930 .set_write_batch_size(1)
1931 .set_max_row_group_row_count(Some(10))
1932 .set_write_page_header_statistics(true)
1933 .build();
1934 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1935 writer.write(&data).unwrap();
1936 writer.close().unwrap();
1937
1938 let data = Bytes::from(buf);
1939 let metadata = ParquetMetaDataReader::new()
1940 .with_page_index_policy(PageIndexPolicy::Required)
1941 .parse_and_finish(&data)
1942 .unwrap();
1943 let parquet_schema = metadata.file_metadata().schema_descr_ptr();
1944
1945 let build_filter = || {
1947 let scalar = Int32Array::from_iter_values([5]);
1948 let predicate = ArrowPredicateFn::new(
1949 ProjectionMask::leaves(&parquet_schema, vec![0]),
1950 move |batch| eq(batch.column(0), &Scalar::new(&scalar)),
1951 );
1952 RowFilter::new(vec![Box::new(predicate)])
1953 };
1954
1955 let selection = RowSelection::from(vec![RowSelector::skip(5), RowSelector::select(1)]);
1957
1958 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1959 let reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1960
1961 let reader_with_cache = TestReader::new(data.clone());
1963 let requests_with_cache = reader_with_cache.requests.clone();
1964 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1965 reader_with_cache,
1966 reader_metadata.clone(),
1967 )
1968 .with_batch_size(1000)
1969 .with_row_selection(selection.clone())
1970 .with_row_filter(build_filter())
1971 .build()
1972 .unwrap();
1973 let batches_with_cache: Vec<_> = stream.try_collect().await.unwrap();
1974
1975 let reader_without_cache = TestReader::new(data);
1977 let requests_without_cache = reader_without_cache.requests.clone();
1978 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1979 reader_without_cache,
1980 reader_metadata,
1981 )
1982 .with_batch_size(1000)
1983 .with_row_selection(selection)
1984 .with_row_filter(build_filter())
1985 .with_max_predicate_cache_size(0) .build()
1987 .unwrap();
1988 let batches_without_cache: Vec<_> = stream.try_collect().await.unwrap();
1989
1990 assert_eq!(batches_with_cache, batches_without_cache);
1991
1992 let requests_with_cache = requests_with_cache.lock().unwrap();
1993 let requests_without_cache = requests_without_cache.lock().unwrap();
1994
1995 assert_eq!(requests_with_cache.len(), 11);
1997 assert_eq!(requests_without_cache.len(), 2);
1998
1999 assert_eq!(
2001 requests_with_cache.iter().map(|r| r.len()).sum::<usize>(),
2002 433
2003 );
2004 assert_eq!(
2005 requests_without_cache
2006 .iter()
2007 .map(|r| r.len())
2008 .sum::<usize>(),
2009 92
2010 );
2011 }
2012
2013 #[test]
2014 fn test_row_numbers_with_multiple_row_groups() {
2015 test_row_numbers_with_multiple_row_groups_helper(
2016 false,
2017 |path, selection, _row_filter, batch_size| {
2018 let runtime = tokio::runtime::Builder::new_current_thread()
2019 .enable_all()
2020 .build()
2021 .expect("Could not create runtime");
2022 runtime.block_on(async move {
2023 let file = tokio::fs::File::open(path).await.unwrap();
2024 let row_number_field = Arc::new(
2025 Field::new("row_number", DataType::Int64, false)
2026 .with_extension_type(RowNumber),
2027 );
2028 let options = ArrowReaderOptions::new()
2029 .with_virtual_columns(vec![row_number_field])
2030 .unwrap();
2031 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2032 .await
2033 .unwrap()
2034 .with_row_selection(selection)
2035 .with_batch_size(batch_size)
2036 .build()
2037 .expect("Could not create reader");
2038 reader.try_collect::<Vec<_>>().await.unwrap()
2039 })
2040 },
2041 );
2042 }
2043
2044 #[test]
2045 fn test_row_numbers_with_multiple_row_groups_and_filter() {
2046 test_row_numbers_with_multiple_row_groups_helper(
2047 true,
2048 |path, selection, row_filter, batch_size| {
2049 let runtime = tokio::runtime::Builder::new_current_thread()
2050 .enable_all()
2051 .build()
2052 .expect("Could not create runtime");
2053 runtime.block_on(async move {
2054 let file = tokio::fs::File::open(path).await.unwrap();
2055 let row_number_field = Arc::new(
2056 Field::new("row_number", DataType::Int64, false)
2057 .with_extension_type(RowNumber),
2058 );
2059 let options = ArrowReaderOptions::new()
2060 .with_virtual_columns(vec![row_number_field])
2061 .unwrap();
2062 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2063 .await
2064 .unwrap()
2065 .with_row_selection(selection)
2066 .with_row_filter(row_filter.expect("No row filter"))
2067 .with_batch_size(batch_size)
2068 .build()
2069 .expect("Could not create reader");
2070 reader.try_collect::<Vec<_>>().await.unwrap()
2071 })
2072 },
2073 );
2074 }
2075
2076 #[tokio::test]
2077 async fn test_nested_lists() -> Result<()> {
2078 let list_inner_field = Arc::new(Field::new("item", DataType::Float32, true));
2080 let table_schema = Arc::new(Schema::new(vec![
2081 Field::new("id", DataType::Int32, false),
2082 Field::new("vector", DataType::List(list_inner_field.clone()), true),
2083 ]));
2084
2085 let mut list_builder =
2086 ListBuilder::new(Float32Builder::new()).with_field(list_inner_field.clone());
2087 list_builder.values().append_slice(&[10.0, 10.0, 10.0]);
2088 list_builder.append(true);
2089 list_builder.values().append_slice(&[20.0, 20.0, 20.0]);
2090 list_builder.append(true);
2091 list_builder.values().append_slice(&[30.0, 30.0, 30.0]);
2092 list_builder.append(true);
2093 list_builder.values().append_slice(&[40.0, 40.0, 40.0]);
2094 list_builder.append(true);
2095 let list_array = list_builder.finish();
2096
2097 let data = vec![RecordBatch::try_new(
2098 table_schema.clone(),
2099 vec![
2100 Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
2101 Arc::new(list_array),
2102 ],
2103 )?];
2104
2105 let mut buffer = Vec::new();
2106 let mut writer = AsyncArrowWriter::try_new(&mut buffer, table_schema, None)?;
2107
2108 for batch in data {
2109 writer.write(&batch).await?;
2110 }
2111
2112 writer.close().await?;
2113
2114 let reader = TestReader::new(Bytes::from(buffer));
2115 let builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
2116
2117 let predicate = ArrowPredicateFn::new(ProjectionMask::all(), |batch| {
2118 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
2119 });
2120
2121 let projection_mask = ProjectionMask::all();
2122
2123 let mut stream = builder
2124 .with_row_filter(RowFilter::new(vec![Box::new(predicate)]))
2125 .with_projection(projection_mask)
2126 .build()?;
2127
2128 while let Some(batch) = stream.next().await {
2129 let _ = batch.unwrap(); }
2131
2132 Ok(())
2133 }
2134}