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(_) => {
634 let bitset_start = bitset_offset
635 .checked_sub(offset)
636 .and_then(|start| usize::try_from(start).ok())
637 .ok_or_else(|| {
638 ParquetError::General("Bloom filter offset is invalid".to_string())
639 })?;
640 buffer.slice(bitset_start..)
641 }
642 None => {
643 let bitset_length: u64 = header.num_bytes.try_into().map_err(|_| {
644 ParquetError::General("Bloom filter length is invalid".to_string())
645 })?;
646 self.input
647 .0
648 .get_bytes(bitset_offset..bitset_offset + bitset_length)
649 .await?
650 }
651 };
652 Ok(Some(Sbbf::new(&bitset)))
653 }
654
655 pub fn with_row_group_selections(
670 mut self,
671 row_group_selections: Vec<RowGroupSelection>,
672 ) -> Self {
673 self.row_group_plan
674 .set_row_group_selections(row_group_selections);
675 self
676 }
677
678 pub fn build(self) -> Result<ParquetRecordBatchStream<T>> {
682 let Self {
683 input,
684 metadata,
685 schema,
686 fields,
687 batch_size,
688 row_group_plan,
689 projection,
690 filter,
691 row_selection_policy: selection_strategy,
692 limit,
693 offset,
694 metrics,
695 max_predicate_cache_size,
696 } = self;
697
698 let projection_len = projection.mask.as_ref().map_or(usize::MAX, |m| m.len());
701 let projected_fields = schema
702 .fields
703 .filter_leaves(|idx, _| idx < projection_len && projection.leaf_included(idx));
704 let projected_schema = Arc::new(Schema::new(projected_fields));
705
706 let decoder = ParquetPushDecoderBuilder {
707 input: PushDecoderInput::default(),
708 metadata,
709 schema,
710 fields,
711 projection,
712 filter,
713 row_group_plan,
714 row_selection_policy: selection_strategy,
715 batch_size,
716 limit,
717 offset,
718 metrics,
719 max_predicate_cache_size,
720 }
721 .build()?;
722
723 let request_state = RequestState::None { input: input.0 };
724
725 Ok(ParquetRecordBatchStream {
726 schema: projected_schema,
727 decoder,
728 request_state,
729 })
730 }
731}
732
733enum RequestState<T> {
737 None {
739 input: T,
740 },
741 Outstanding {
743 ranges: Vec<Range<u64>>,
745 future: BoxFuture<'static, Result<(T, Vec<Bytes>)>>,
750 },
751 Done,
752}
753
754impl<T> RequestState<T>
755where
756 T: AsyncFileReader + Unpin + Send + 'static,
757{
758 fn begin_request(mut input: T, ranges: Vec<Range<u64>>) -> Self {
760 let ranges_captured = ranges.clone();
761
762 let future = async move {
767 let data = input.get_byte_ranges(ranges_captured).await?;
768 Ok((input, data))
769 }
770 .boxed();
771 RequestState::Outstanding { ranges, future }
772 }
773}
774
775impl<T> std::fmt::Debug for RequestState<T> {
776 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
777 match self {
778 RequestState::None { input: _ } => f
779 .debug_struct("RequestState::None")
780 .field("input", &"...")
781 .finish(),
782 RequestState::Outstanding { ranges, .. } => f
783 .debug_struct("RequestState::Outstanding")
784 .field("ranges", &ranges)
785 .finish(),
786 RequestState::Done => {
787 write!(f, "RequestState::Done")
788 }
789 }
790 }
791}
792
793pub struct ParquetRecordBatchStream<T> {
813 schema: SchemaRef,
815 request_state: RequestState<T>,
817 decoder: ParquetPushDecoder,
819}
820
821impl<T> std::fmt::Debug for ParquetRecordBatchStream<T> {
822 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
823 f.debug_struct("ParquetRecordBatchStream")
824 .field("request_state", &self.request_state)
825 .finish()
826 }
827}
828
829impl<T> ParquetRecordBatchStream<T> {
830 pub fn schema(&self) -> &SchemaRef {
835 &self.schema
836 }
837}
838
839impl<T> ParquetRecordBatchStream<T>
840where
841 T: AsyncFileReader + Unpin + Send + 'static,
842{
843 pub async fn next_row_group(&mut self) -> Result<Option<ParquetRecordBatchReader>> {
857 loop {
858 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
861 match request_state {
862 RequestState::None { input } => {
864 match self.decoder.try_next_reader()? {
865 DecodeResult::NeedsData(ranges) => {
866 self.request_state = RequestState::begin_request(input, ranges);
867 }
869 DecodeResult::Data(reader) => {
870 self.request_state = RequestState::None { input };
871 return Ok(Some(reader));
872 }
873 DecodeResult::Finished => return Ok(None),
874 }
875 }
876 RequestState::Outstanding { ranges, future } => {
877 let (input, data) = future.await?;
878 self.decoder.push_ranges(ranges, data)?;
880 self.request_state = RequestState::None { input };
881 }
883 RequestState::Done => {
884 self.request_state = RequestState::Done;
885 return Ok(None);
886 }
887 }
888 }
889 }
890}
891
892impl<T> Stream for ParquetRecordBatchStream<T>
893where
894 T: AsyncFileReader + Unpin + Send + 'static,
895{
896 type Item = Result<RecordBatch>;
897 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
898 match self.poll_next_inner(cx) {
899 Ok(res) => {
900 res.map(|res| Ok(res).transpose())
903 }
904 Err(e) => {
905 self.request_state = RequestState::Done;
906 Poll::Ready(Some(Err(e)))
907 }
908 }
909 }
910}
911
912impl<T> ParquetRecordBatchStream<T>
913where
914 T: AsyncFileReader + Unpin + Send + 'static,
915{
916 fn poll_next_inner(&mut self, cx: &mut Context<'_>) -> Result<Poll<Option<RecordBatch>>> {
921 loop {
922 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
923 match request_state {
924 RequestState::None { input } => {
925 match self.decoder.try_decode()? {
927 DecodeResult::NeedsData(ranges) => {
928 self.request_state = RequestState::begin_request(input, ranges);
929 }
931 DecodeResult::Data(batch) => {
932 self.request_state = RequestState::None { input };
933 return Ok(Poll::Ready(Some(batch)));
934 }
935 DecodeResult::Finished => {
936 self.request_state = RequestState::Done;
937 return Ok(Poll::Ready(None));
938 }
939 }
940 }
941 RequestState::Outstanding { ranges, mut future } => match future.poll_unpin(cx) {
942 Poll::Ready(result) => {
944 let (input, data) = result?;
945 self.decoder.push_ranges(ranges, data)?;
947 self.request_state = RequestState::None { input };
948 }
950 Poll::Pending => {
951 self.request_state = RequestState::Outstanding { ranges, future };
952 return Ok(Poll::Pending);
953 }
954 },
955 RequestState::Done => {
956 self.request_state = RequestState::Done;
958 return Ok(Poll::Ready(None));
959 }
960 }
961 }
962 }
963}
964
965#[cfg(test)]
966mod tests {
967 use super::*;
968 use crate::arrow::arrow_reader::tests::test_row_numbers_with_multiple_row_groups_helper;
969 use crate::arrow::arrow_reader::{
970 ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowFilter, RowSelection, RowSelector,
971 };
972 use crate::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
973 use crate::arrow::schema::virtual_type::RowNumber;
974 use crate::arrow::{ArrowWriter, AsyncArrowWriter, ProjectionMask};
975 use crate::file::metadata::PageIndexPolicy;
976 use crate::file::metadata::ParquetMetaDataReader;
977 use crate::file::metadata::page_index::PageIndex;
978 use crate::file::properties::WriterProperties;
979 use arrow::compute::kernels::cmp::eq;
980 use arrow::error::Result as ArrowResult;
981 use arrow_array::builder::{Float32Builder, ListBuilder, StringBuilder};
982 use arrow_array::cast::AsArray;
983 use arrow_array::types::Int32Type;
984 use arrow_array::{
985 Array, ArrayRef, BooleanArray, Int32Array, RecordBatchReader, Scalar, StringArray,
986 StructArray, UInt64Array,
987 };
988 use arrow_schema::{DataType, Field, Schema};
989 use futures::{StreamExt, TryStreamExt};
990 use rand::{RngExt, rng};
991 use std::collections::HashMap;
992 use std::sync::{Arc, Mutex};
993 use tempfile::tempfile;
994
995 #[derive(Clone)]
996 struct TestReader {
997 data: Bytes,
998 metadata: Option<Arc<ParquetMetaData>>,
999 requests: Arc<Mutex<Vec<Range<usize>>>>,
1000 }
1001
1002 impl TestReader {
1003 fn new(data: Bytes) -> Self {
1004 Self {
1005 data,
1006 metadata: Default::default(),
1007 requests: Default::default(),
1008 }
1009 }
1010 }
1011
1012 impl AsyncFileReader for TestReader {
1013 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1014 let range = range.clone();
1015 self.requests
1016 .lock()
1017 .unwrap()
1018 .push(range.start as usize..range.end as usize);
1019 futures::future::ready(Ok(self
1020 .data
1021 .slice(range.start as usize..range.end as usize)))
1022 .boxed()
1023 }
1024
1025 fn get_metadata<'a>(
1026 &'a mut self,
1027 options: Option<&'a ArrowReaderOptions>,
1028 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
1029 let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
1030 self.metadata = Some(Arc::new(
1031 metadata_reader.parse_and_finish(&self.data).unwrap(),
1032 ));
1033 futures::future::ready(Ok(self.metadata.clone().unwrap().clone())).boxed()
1034 }
1035 }
1036
1037 #[tokio::test]
1038 async fn test_async_reader() {
1039 let testdata = arrow::util::test_util::parquet_test_data();
1040 let path = format!("{testdata}/alltypes_plain.parquet");
1041 let data = Bytes::from(std::fs::read(path).unwrap());
1042
1043 let async_reader = TestReader::new(data.clone());
1044
1045 let requests = async_reader.requests.clone();
1046 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1047 .await
1048 .unwrap();
1049
1050 let metadata = builder.metadata().clone();
1051 assert_eq!(metadata.num_row_groups(), 1);
1052
1053 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1054 let stream = builder
1055 .with_projection(mask.clone())
1056 .with_batch_size(1024)
1057 .build()
1058 .unwrap();
1059
1060 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1061
1062 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1063 .unwrap()
1064 .with_projection(mask)
1065 .with_batch_size(104)
1066 .build()
1067 .unwrap()
1068 .collect::<ArrowResult<Vec<_>>>()
1069 .unwrap();
1070
1071 assert_eq!(async_batches, sync_batches);
1072
1073 let requests = requests.lock().unwrap();
1074 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1075 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1076
1077 assert_eq!(
1078 &requests[..],
1079 &[
1080 offset_1 as usize..(offset_1 + length_1) as usize,
1081 offset_2 as usize..(offset_2 + length_2) as usize
1082 ]
1083 );
1084 }
1085
1086 #[tokio::test]
1087 async fn test_async_reader_row_group_local_selections() {
1088 let batch = RecordBatch::try_from_iter([(
1089 "a",
1090 Arc::new(Int32Array::from_iter_values(0..6)) as ArrayRef,
1091 )])
1092 .unwrap();
1093 let mut data = Vec::new();
1094 let properties = WriterProperties::builder()
1095 .set_max_row_group_row_count(Some(3))
1096 .build();
1097 let mut writer = ArrowWriter::try_new(&mut data, batch.schema(), Some(properties)).unwrap();
1098 writer.write(&batch).unwrap();
1099 writer.close().unwrap();
1100
1101 let stream = ParquetRecordBatchStreamBuilder::new(TestReader::new(data.into()))
1102 .await
1103 .unwrap()
1104 .with_row_group_selections(vec![
1105 RowGroupSelection::new(1, Some(RowSelection::from(vec![RowSelector::select(1)]))),
1106 RowGroupSelection::new(
1107 0,
1108 Some(RowSelection::from(vec![
1109 RowSelector::skip(1),
1110 RowSelector::select(2),
1111 ])),
1112 ),
1113 ])
1114 .build()
1115 .unwrap();
1116
1117 let batches: Vec<_> = stream.try_collect().await.unwrap();
1118 assert_eq!(batches.len(), 2);
1119 assert_eq!(
1120 batches[0].column(0).as_primitive::<Int32Type>().values(),
1121 &[3]
1122 );
1123 assert_eq!(
1124 batches[1].column(0).as_primitive::<Int32Type>().values(),
1125 &[1, 2]
1126 );
1127 }
1128
1129 #[tokio::test]
1130 async fn test_async_reader_with_next_row_group() {
1131 let testdata = arrow::util::test_util::parquet_test_data();
1132 let path = format!("{testdata}/alltypes_plain.parquet");
1133 let data = Bytes::from(std::fs::read(path).unwrap());
1134
1135 let async_reader = TestReader::new(data.clone());
1136
1137 let requests = async_reader.requests.clone();
1138 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1139 .await
1140 .unwrap();
1141
1142 let metadata = builder.metadata().clone();
1143 assert_eq!(metadata.num_row_groups(), 1);
1144
1145 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1146 let mut stream = builder
1147 .with_projection(mask.clone())
1148 .with_batch_size(1024)
1149 .build()
1150 .unwrap();
1151
1152 let mut readers = vec![];
1153 while let Some(reader) = stream.next_row_group().await.unwrap() {
1154 readers.push(reader);
1155 }
1156
1157 let async_batches: Vec<_> = readers
1158 .into_iter()
1159 .flat_map(|r| r.map(|v| v.unwrap()).collect::<Vec<_>>())
1160 .collect();
1161
1162 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1163 .unwrap()
1164 .with_projection(mask)
1165 .with_batch_size(104)
1166 .build()
1167 .unwrap()
1168 .collect::<ArrowResult<Vec<_>>>()
1169 .unwrap();
1170
1171 assert_eq!(async_batches, sync_batches);
1172
1173 let requests = requests.lock().unwrap();
1174 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1175 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1176
1177 assert_eq!(
1178 &requests[..],
1179 &[
1180 offset_1 as usize..(offset_1 + length_1) as usize,
1181 offset_2 as usize..(offset_2 + length_2) as usize
1182 ]
1183 );
1184 }
1185
1186 #[tokio::test]
1187 async fn test_async_reader_with_index() {
1188 let testdata = arrow::util::test_util::parquet_test_data();
1189 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1190 let data = Bytes::from(std::fs::read(path).unwrap());
1191
1192 let async_reader = TestReader::new(data.clone());
1193
1194 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1195 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1196 .await
1197 .unwrap();
1198
1199 let metadata_with_index = builder.metadata();
1201 assert_eq!(metadata_with_index.num_row_groups(), 1);
1202
1203 let page_index = metadata_with_index
1205 .page_index()
1206 .expect("page index should be present");
1207 assert!(page_index.is_complete());
1208 let num_rowgroups = metadata_with_index.num_row_groups();
1209 let num_columns = metadata_with_index
1210 .file_metadata()
1211 .schema_descr()
1212 .num_columns();
1213 for rgidx in 0..num_rowgroups {
1214 for colidx in 0..num_columns {
1216 assert!(page_index.offset_index(rgidx, colidx).is_some());
1217 }
1218 }
1219
1220 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1221 let stream = builder
1222 .with_projection(mask.clone())
1223 .with_batch_size(1024)
1224 .build()
1225 .unwrap();
1226
1227 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1228
1229 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1230 .unwrap()
1231 .with_projection(mask)
1232 .with_batch_size(1024)
1233 .build()
1234 .unwrap()
1235 .collect::<ArrowResult<Vec<_>>>()
1236 .unwrap();
1237
1238 assert_eq!(async_batches, sync_batches);
1239 }
1240
1241 #[tokio::test]
1242 async fn test_async_reader_with_limit() {
1243 let testdata = arrow::util::test_util::parquet_test_data();
1244 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1245 let data = Bytes::from(std::fs::read(path).unwrap());
1246
1247 let metadata = ParquetMetaDataReader::new()
1248 .parse_and_finish(&data)
1249 .unwrap();
1250 let metadata = Arc::new(metadata);
1251
1252 assert_eq!(metadata.num_row_groups(), 1);
1253
1254 let async_reader = TestReader::new(data.clone());
1255
1256 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1257 .await
1258 .unwrap();
1259
1260 assert_eq!(builder.metadata().num_row_groups(), 1);
1261
1262 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1263 let stream = builder
1264 .with_projection(mask.clone())
1265 .with_batch_size(1024)
1266 .with_limit(1)
1267 .build()
1268 .unwrap();
1269
1270 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1271
1272 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1273 .unwrap()
1274 .with_projection(mask)
1275 .with_batch_size(1024)
1276 .with_limit(1)
1277 .build()
1278 .unwrap()
1279 .collect::<ArrowResult<Vec<_>>>()
1280 .unwrap();
1281
1282 assert_eq!(async_batches, sync_batches);
1283 }
1284
1285 #[tokio::test]
1286 async fn test_async_reader_skip_pages() {
1287 let testdata = arrow::util::test_util::parquet_test_data();
1288 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1289 let data = Bytes::from(std::fs::read(path).unwrap());
1290
1291 let async_reader = TestReader::new(data.clone());
1292
1293 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1294 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1295 .await
1296 .unwrap();
1297
1298 assert_eq!(builder.metadata().num_row_groups(), 1);
1299
1300 let selection = RowSelection::from(vec![
1301 RowSelector::skip(21), RowSelector::select(21), RowSelector::skip(41), RowSelector::select(41), RowSelector::skip(25), RowSelector::select(25), RowSelector::skip(7116), RowSelector::select(10), ]);
1310
1311 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![9]);
1312
1313 let stream = builder
1314 .with_projection(mask.clone())
1315 .with_row_selection(selection.clone())
1316 .build()
1317 .expect("building stream");
1318
1319 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1320
1321 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1322 .unwrap()
1323 .with_projection(mask)
1324 .with_batch_size(1024)
1325 .with_row_selection(selection)
1326 .build()
1327 .unwrap()
1328 .collect::<ArrowResult<Vec<_>>>()
1329 .unwrap();
1330
1331 assert_eq!(async_batches, sync_batches);
1332 }
1333
1334 #[tokio::test]
1335 async fn test_fuzz_async_reader_selection() {
1336 let testdata = arrow::util::test_util::parquet_test_data();
1337 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1338 let data = Bytes::from(std::fs::read(path).unwrap());
1339
1340 let mut rand = rng();
1341
1342 for _ in 0..100 {
1343 let mut expected_rows = 0;
1344 let mut total_rows = 0;
1345 let mut skip = false;
1346 let mut selectors = vec![];
1347
1348 while total_rows < 7300 {
1349 let row_count: usize = rand.random_range(1..100);
1350
1351 let row_count = row_count.min(7300 - total_rows);
1352
1353 selectors.push(RowSelector { row_count, skip });
1354
1355 total_rows += row_count;
1356 if !skip {
1357 expected_rows += row_count;
1358 }
1359
1360 skip = !skip;
1361 }
1362
1363 let selection = RowSelection::from(selectors);
1364
1365 let async_reader = TestReader::new(data.clone());
1366
1367 let options =
1368 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1369 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1370 .await
1371 .unwrap();
1372
1373 assert_eq!(builder.metadata().num_row_groups(), 1);
1374
1375 let col_idx: usize = rand.random_range(0..13);
1376 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1377
1378 let stream = builder
1379 .with_projection(mask.clone())
1380 .with_row_selection(selection.clone())
1381 .build()
1382 .expect("building stream");
1383
1384 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1385
1386 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1387
1388 assert_eq!(actual_rows, expected_rows);
1389 }
1390 }
1391
1392 #[tokio::test]
1393 async fn test_async_reader_zero_row_selector() {
1394 let testdata = arrow::util::test_util::parquet_test_data();
1396 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1397 let data = Bytes::from(std::fs::read(path).unwrap());
1398
1399 let mut rand = rng();
1400
1401 let mut expected_rows = 0;
1402 let mut total_rows = 0;
1403 let mut skip = false;
1404 let mut selectors = vec![];
1405
1406 selectors.push(RowSelector {
1407 row_count: 0,
1408 skip: false,
1409 });
1410
1411 while total_rows < 7300 {
1412 let row_count: usize = rand.random_range(1..100);
1413
1414 let row_count = row_count.min(7300 - total_rows);
1415
1416 selectors.push(RowSelector { row_count, skip });
1417
1418 total_rows += row_count;
1419 if !skip {
1420 expected_rows += row_count;
1421 }
1422
1423 skip = !skip;
1424 }
1425
1426 let selection = RowSelection::from(selectors);
1427
1428 let async_reader = TestReader::new(data.clone());
1429
1430 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1431 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1432 .await
1433 .unwrap();
1434
1435 assert_eq!(builder.metadata().num_row_groups(), 1);
1436
1437 let col_idx: usize = rand.random_range(0..13);
1438 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1439
1440 let stream = builder
1441 .with_projection(mask.clone())
1442 .with_row_selection(selection.clone())
1443 .build()
1444 .expect("building stream");
1445
1446 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1447
1448 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1449
1450 assert_eq!(actual_rows, expected_rows);
1451 }
1452
1453 #[tokio::test]
1454 async fn test_limit_multiple_row_groups() {
1455 let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]);
1456 let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]);
1457 let c = Int32Array::from_iter(0..6);
1458 let data = RecordBatch::try_from_iter([
1459 ("a", Arc::new(a) as ArrayRef),
1460 ("b", Arc::new(b) as ArrayRef),
1461 ("c", Arc::new(c) as ArrayRef),
1462 ])
1463 .unwrap();
1464
1465 let mut buf = Vec::with_capacity(1024);
1466 let props = WriterProperties::builder()
1467 .set_max_row_group_row_count(Some(3))
1468 .build();
1469 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1470 writer.write(&data).unwrap();
1471 writer.close().unwrap();
1472
1473 let data: Bytes = buf.into();
1474 let metadata = ParquetMetaDataReader::new()
1475 .parse_and_finish(&data)
1476 .unwrap();
1477
1478 assert_eq!(metadata.num_row_groups(), 2);
1479
1480 let test = TestReader::new(data);
1481
1482 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1483 .await
1484 .unwrap()
1485 .with_batch_size(1024)
1486 .with_limit(4)
1487 .build()
1488 .unwrap();
1489
1490 let batches: Vec<_> = stream.try_collect().await.unwrap();
1491 assert_eq!(batches.len(), 2);
1493
1494 let batch = &batches[0];
1495 assert_eq!(batch.num_rows(), 3);
1497 assert_eq!(batch.num_columns(), 3);
1498 let col2 = batch.column(2).as_primitive::<Int32Type>();
1499 assert_eq!(col2.values(), &[0, 1, 2]);
1500
1501 let batch = &batches[1];
1502 assert_eq!(batch.num_rows(), 1);
1504 assert_eq!(batch.num_columns(), 3);
1505 let col2 = batch.column(2).as_primitive::<Int32Type>();
1506 assert_eq!(col2.values(), &[3]);
1507
1508 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1509 .await
1510 .unwrap()
1511 .with_offset(2)
1512 .with_limit(3)
1513 .build()
1514 .unwrap();
1515
1516 let batches: Vec<_> = stream.try_collect().await.unwrap();
1517 assert_eq!(batches.len(), 2);
1519
1520 let batch = &batches[0];
1521 assert_eq!(batch.num_rows(), 1);
1523 assert_eq!(batch.num_columns(), 3);
1524 let col2 = batch.column(2).as_primitive::<Int32Type>();
1525 assert_eq!(col2.values(), &[2]);
1526
1527 let batch = &batches[1];
1528 assert_eq!(batch.num_rows(), 2);
1530 assert_eq!(batch.num_columns(), 3);
1531 let col2 = batch.column(2).as_primitive::<Int32Type>();
1532 assert_eq!(col2.values(), &[3, 4]);
1533
1534 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1535 .await
1536 .unwrap()
1537 .with_offset(4)
1538 .with_limit(20)
1539 .build()
1540 .unwrap();
1541
1542 let batches: Vec<_> = stream.try_collect().await.unwrap();
1543 assert_eq!(batches.len(), 1);
1545
1546 let batch = &batches[0];
1547 assert_eq!(batch.num_rows(), 2);
1549 assert_eq!(batch.num_columns(), 3);
1550 let col2 = batch.column(2).as_primitive::<Int32Type>();
1551 assert_eq!(col2.values(), &[4, 5]);
1552 }
1553
1554 #[tokio::test]
1555 async fn test_batch_size_overallocate() {
1556 let testdata = arrow::util::test_util::parquet_test_data();
1557 let path = format!("{testdata}/alltypes_plain.parquet");
1559 let data = Bytes::from(std::fs::read(path).unwrap());
1560
1561 let async_reader = TestReader::new(data.clone());
1562
1563 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1564 .await
1565 .unwrap();
1566
1567 let file_rows = builder.metadata().file_metadata().num_rows() as usize;
1568
1569 let builder = builder
1570 .with_projection(ProjectionMask::all())
1571 .with_batch_size(1024);
1572
1573 assert_ne!(1024, file_rows);
1576 assert_eq!(builder.batch_size, file_rows);
1577
1578 let _stream = builder.build().unwrap();
1579 }
1580
1581 #[tokio::test]
1582 async fn test_parquet_record_batch_stream_schema() {
1583 fn get_all_field_names(schema: &Schema) -> Vec<&String> {
1584 schema.flattened_fields().iter().map(|f| f.name()).collect()
1585 }
1586
1587 let mut metadata = HashMap::with_capacity(1);
1596 metadata.insert("key".to_string(), "value".to_string());
1597
1598 let nested_struct_array = StructArray::from(vec![
1599 (
1600 Arc::new(Field::new("d", DataType::Utf8, true)),
1601 Arc::new(StringArray::from(vec!["a", "b"])) as ArrayRef,
1602 ),
1603 (
1604 Arc::new(Field::new("e", DataType::Utf8, true)),
1605 Arc::new(StringArray::from(vec!["c", "d"])) as ArrayRef,
1606 ),
1607 ]);
1608 let struct_array = StructArray::from(vec![
1609 (
1610 Arc::new(Field::new("a", DataType::Int32, true)),
1611 Arc::new(Int32Array::from(vec![-1, 1])) as ArrayRef,
1612 ),
1613 (
1614 Arc::new(Field::new("b", DataType::UInt64, true)),
1615 Arc::new(UInt64Array::from(vec![1, 2])) as ArrayRef,
1616 ),
1617 (
1618 Arc::new(Field::new(
1619 "c",
1620 nested_struct_array.data_type().clone(),
1621 true,
1622 )),
1623 Arc::new(nested_struct_array) as ArrayRef,
1624 ),
1625 ]);
1626
1627 let schema =
1628 Arc::new(Schema::new(struct_array.fields().clone()).with_metadata(metadata.clone()));
1629 let record_batch = RecordBatch::from(struct_array)
1630 .with_schema(schema.clone())
1631 .unwrap();
1632
1633 let mut file = tempfile().unwrap();
1635 let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1636 writer.write(&record_batch).unwrap();
1637 writer.close().unwrap();
1638
1639 let all_fields = ["a", "b", "c", "d", "e"];
1640 let projections = [
1642 (vec![], vec![]),
1643 (vec![0], vec!["a"]),
1644 (vec![0, 1], vec!["a", "b"]),
1645 (vec![0, 1, 2], vec!["a", "b", "c", "d"]),
1646 (vec![0, 1, 2, 3], vec!["a", "b", "c", "d", "e"]),
1647 ];
1648
1649 for (indices, expected_projected_names) in projections {
1651 let assert_schemas = |builder: SchemaRef, reader: SchemaRef, batch: SchemaRef| {
1652 assert_eq!(get_all_field_names(&builder), all_fields);
1654 assert_eq!(builder.metadata, metadata);
1655 assert_eq!(get_all_field_names(&reader), expected_projected_names);
1657 assert_eq!(reader.metadata, HashMap::default());
1658 assert_eq!(get_all_field_names(&batch), expected_projected_names);
1659 assert_eq!(batch.metadata, HashMap::default());
1660 };
1661
1662 let builder =
1663 ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1664 let sync_builder_schema = builder.schema().clone();
1665 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices.clone());
1666 let mut reader = builder.with_projection(mask).build().unwrap();
1667 let sync_reader_schema = reader.schema();
1668 let batch = reader.next().unwrap().unwrap();
1669 let sync_batch_schema = batch.schema();
1670 assert_schemas(sync_builder_schema, sync_reader_schema, sync_batch_schema);
1671
1672 let file = tokio::fs::File::from(file.try_clone().unwrap());
1674 let builder = ParquetRecordBatchStreamBuilder::new(file).await.unwrap();
1675 let async_builder_schema = builder.schema().clone();
1676 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices);
1677 let mut reader = builder.with_projection(mask).build().unwrap();
1678 let async_reader_schema = reader.schema().clone();
1679 let batch = reader.next().await.unwrap().unwrap();
1680 let async_batch_schema = batch.schema();
1681 assert_schemas(
1682 async_builder_schema,
1683 async_reader_schema,
1684 async_batch_schema,
1685 );
1686 }
1687 }
1688
1689 #[tokio::test]
1690 async fn test_nested_skip() {
1691 let schema = Arc::new(Schema::new(vec![
1692 Field::new("col_1", DataType::UInt64, false),
1693 Field::new_list("col_2", Field::new_list_field(DataType::Utf8, true), true),
1694 ]));
1695
1696 let props = WriterProperties::builder()
1698 .set_data_page_row_count_limit(256)
1699 .set_write_batch_size(256)
1700 .set_max_row_group_row_count(Some(1024));
1701
1702 let mut file = tempfile().unwrap();
1704 let mut writer =
1705 ArrowWriter::try_new(&mut file, schema.clone(), Some(props.build())).unwrap();
1706
1707 let mut builder = ListBuilder::new(StringBuilder::new());
1708 for id in 0..1024 {
1709 match id % 3 {
1710 0 => builder.append_value([Some("val_1".to_string()), Some(format!("id_{id}"))]),
1711 1 => builder.append_value([Some(format!("id_{id}"))]),
1712 _ => builder.append_null(),
1713 }
1714 }
1715 let refs = vec![
1716 Arc::new(UInt64Array::from_iter_values(0..1024)) as ArrayRef,
1717 Arc::new(builder.finish()) as ArrayRef,
1718 ];
1719
1720 let batch = RecordBatch::try_new(schema.clone(), refs).unwrap();
1721 writer.write(&batch).unwrap();
1722 writer.close().unwrap();
1723
1724 let selections = [
1725 RowSelection::from(vec![
1726 RowSelector::skip(313),
1727 RowSelector::select(1),
1728 RowSelector::skip(709),
1729 RowSelector::select(1),
1730 ]),
1731 RowSelection::from(vec![
1732 RowSelector::skip(255),
1733 RowSelector::select(1),
1734 RowSelector::skip(767),
1735 RowSelector::select(1),
1736 ]),
1737 RowSelection::from(vec![
1738 RowSelector::select(255),
1739 RowSelector::skip(1),
1740 RowSelector::select(767),
1741 RowSelector::skip(1),
1742 ]),
1743 RowSelection::from(vec![
1744 RowSelector::skip(254),
1745 RowSelector::select(1),
1746 RowSelector::select(1),
1747 RowSelector::skip(767),
1748 RowSelector::select(1),
1749 ]),
1750 ];
1751
1752 for selection in selections {
1753 let expected = selection.row_count();
1754 let mut reader = ParquetRecordBatchStreamBuilder::new_with_options(
1756 tokio::fs::File::from_std(file.try_clone().unwrap()),
1757 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
1758 )
1759 .await
1760 .unwrap();
1761
1762 reader = reader.with_row_selection(selection);
1763
1764 let mut stream = reader.build().unwrap();
1765
1766 let mut total_rows = 0;
1767 while let Some(rb) = stream.next().await {
1768 let rb = rb.unwrap();
1769 total_rows += rb.num_rows();
1770 }
1771 assert_eq!(total_rows, expected);
1772 }
1773 }
1774
1775 #[tokio::test]
1776 async fn empty_offset_index_doesnt_panic_in_read_row_group() {
1777 use tokio::fs::File;
1778 let testdata = arrow::util::test_util::parquet_test_data();
1779 let path = format!("{testdata}/alltypes_plain.parquet");
1780 let mut file = File::open(&path).await.unwrap();
1781 let file_size = file.metadata().await.unwrap().len();
1782 let mut metadata = ParquetMetaDataReader::new()
1783 .with_page_index_policy(PageIndexPolicy::Required)
1784 .load_and_finish(&mut file, file_size)
1785 .await
1786 .unwrap();
1787
1788 let page_index = PageIndex::new(None, Some(vec![]));
1789 metadata.set_page_index(Some(Arc::new(page_index)));
1790 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1791 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1792 let reader =
1793 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1794 .build()
1795 .unwrap();
1796
1797 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1798 assert_eq!(result.len(), 1);
1799 }
1800
1801 #[tokio::test]
1802 async fn non_empty_offset_index_doesnt_panic_in_read_row_group() {
1803 use tokio::fs::File;
1804 let testdata = arrow::util::test_util::parquet_test_data();
1805 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1806 let mut file = File::open(&path).await.unwrap();
1807 let file_size = file.metadata().await.unwrap().len();
1808 let metadata = ParquetMetaDataReader::new()
1809 .with_page_index_policy(PageIndexPolicy::Required)
1810 .load_and_finish(&mut file, file_size)
1811 .await
1812 .unwrap();
1813
1814 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1815 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1816 let reader =
1817 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1818 .build()
1819 .unwrap();
1820
1821 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1822 assert_eq!(result.len(), 8);
1823 }
1824
1825 #[tokio::test]
1826 async fn empty_offset_index_doesnt_panic_in_column_chunks() {
1827 use tempfile::TempDir;
1828 use tokio::fs::File;
1829 fn write_metadata_to_local_file(
1830 metadata: ParquetMetaData,
1831 file: impl AsRef<std::path::Path>,
1832 ) {
1833 use crate::file::metadata::ParquetMetaDataWriter;
1834 use std::fs::File;
1835 let file = File::create(file).unwrap();
1836 ParquetMetaDataWriter::new(file, &metadata)
1837 .finish()
1838 .unwrap()
1839 }
1840
1841 fn read_metadata_from_local_file(file: impl AsRef<std::path::Path>) -> ParquetMetaData {
1842 use std::fs::File;
1843 let file = File::open(file).unwrap();
1844 ParquetMetaDataReader::new()
1845 .with_page_index_policy(PageIndexPolicy::Required)
1846 .parse_and_finish(&file)
1847 .unwrap()
1848 }
1849
1850 let testdata = arrow::util::test_util::parquet_test_data();
1851 let path = format!("{testdata}/alltypes_plain.parquet");
1852 let mut file = File::open(&path).await.unwrap();
1853 let file_size = file.metadata().await.unwrap().len();
1854 let metadata = ParquetMetaDataReader::new()
1855 .with_page_index_policy(PageIndexPolicy::Required)
1856 .load_and_finish(&mut file, file_size)
1857 .await
1858 .unwrap();
1859
1860 let tempdir = TempDir::new().unwrap();
1861 let metadata_path = tempdir.path().join("thrift_metadata.dat");
1862 write_metadata_to_local_file(metadata, &metadata_path);
1863 let metadata = read_metadata_from_local_file(&metadata_path);
1864
1865 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1866 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1867 let reader =
1868 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1869 .build()
1870 .unwrap();
1871
1872 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1874 assert_eq!(result.len(), 1);
1875 }
1876
1877 #[tokio::test]
1878 async fn test_cached_array_reader_sparse_offset_error() {
1879 use futures::TryStreamExt;
1880
1881 use crate::arrow::arrow_reader::{ArrowPredicateFn, RowFilter, RowSelection, RowSelector};
1882 use arrow_array::{BooleanArray, RecordBatch};
1883
1884 let testdata = arrow::util::test_util::parquet_test_data();
1885 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1886 let data = Bytes::from(std::fs::read(path).unwrap());
1887
1888 let async_reader = TestReader::new(data);
1889
1890 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1892 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1893 .await
1894 .unwrap();
1895
1896 let selection = RowSelection::from(vec![RowSelector::skip(22), RowSelector::select(3)]);
1900
1901 let parquet_schema = builder.parquet_schema();
1905 let proj = ProjectionMask::leaves(parquet_schema, vec![0]);
1906 let always_true = ArrowPredicateFn::new(proj.clone(), |batch: RecordBatch| {
1907 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
1908 });
1909 let filter = RowFilter::new(vec![Box::new(always_true)]);
1910
1911 let stream = builder
1914 .with_batch_size(8)
1915 .with_projection(proj)
1916 .with_row_selection(selection)
1917 .with_row_filter(filter)
1918 .build()
1919 .unwrap();
1920
1921 let _result: Vec<_> = stream.try_collect().await.unwrap();
1924 }
1925
1926 #[tokio::test]
1927 async fn test_predicate_cache_disabled() {
1928 let k = Int32Array::from_iter_values(0..10);
1929 let data = RecordBatch::try_from_iter([("k", Arc::new(k) as ArrayRef)]).unwrap();
1930
1931 let mut buf = Vec::new();
1932 let props = WriterProperties::builder()
1934 .set_data_page_row_count_limit(1)
1935 .set_write_batch_size(1)
1936 .set_max_row_group_row_count(Some(10))
1937 .set_write_page_header_statistics(true)
1938 .build();
1939 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1940 writer.write(&data).unwrap();
1941 writer.close().unwrap();
1942
1943 let data = Bytes::from(buf);
1944 let metadata = ParquetMetaDataReader::new()
1945 .with_page_index_policy(PageIndexPolicy::Required)
1946 .parse_and_finish(&data)
1947 .unwrap();
1948 let parquet_schema = metadata.file_metadata().schema_descr_ptr();
1949
1950 let build_filter = || {
1952 let scalar = Int32Array::from_iter_values([5]);
1953 let predicate = ArrowPredicateFn::new(
1954 ProjectionMask::leaves(&parquet_schema, vec![0]),
1955 move |batch| eq(batch.column(0), &Scalar::new(&scalar)),
1956 );
1957 RowFilter::new(vec![Box::new(predicate)])
1958 };
1959
1960 let selection = RowSelection::from(vec![RowSelector::skip(5), RowSelector::select(1)]);
1962
1963 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1964 let reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1965
1966 let reader_with_cache = TestReader::new(data.clone());
1968 let requests_with_cache = reader_with_cache.requests.clone();
1969 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1970 reader_with_cache,
1971 reader_metadata.clone(),
1972 )
1973 .with_batch_size(1000)
1974 .with_row_selection(selection.clone())
1975 .with_row_filter(build_filter())
1976 .build()
1977 .unwrap();
1978 let batches_with_cache: Vec<_> = stream.try_collect().await.unwrap();
1979
1980 let reader_without_cache = TestReader::new(data);
1982 let requests_without_cache = reader_without_cache.requests.clone();
1983 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1984 reader_without_cache,
1985 reader_metadata,
1986 )
1987 .with_batch_size(1000)
1988 .with_row_selection(selection)
1989 .with_row_filter(build_filter())
1990 .with_max_predicate_cache_size(0) .build()
1992 .unwrap();
1993 let batches_without_cache: Vec<_> = stream.try_collect().await.unwrap();
1994
1995 assert_eq!(batches_with_cache, batches_without_cache);
1996
1997 let requests_with_cache = requests_with_cache.lock().unwrap();
1998 let requests_without_cache = requests_without_cache.lock().unwrap();
1999
2000 assert_eq!(requests_with_cache.len(), 11);
2002 assert_eq!(requests_without_cache.len(), 2);
2003
2004 assert_eq!(
2006 requests_with_cache.iter().map(|r| r.len()).sum::<usize>(),
2007 433
2008 );
2009 assert_eq!(
2010 requests_without_cache
2011 .iter()
2012 .map(|r| r.len())
2013 .sum::<usize>(),
2014 92
2015 );
2016 }
2017
2018 #[test]
2019 fn test_row_numbers_with_multiple_row_groups() {
2020 test_row_numbers_with_multiple_row_groups_helper(
2021 false,
2022 |path, selection, _row_filter, batch_size| {
2023 let runtime = tokio::runtime::Builder::new_current_thread()
2024 .enable_all()
2025 .build()
2026 .expect("Could not create runtime");
2027 runtime.block_on(async move {
2028 let file = tokio::fs::File::open(path).await.unwrap();
2029 let row_number_field = Arc::new(
2030 Field::new("row_number", DataType::Int64, false)
2031 .with_extension_type(RowNumber),
2032 );
2033 let options = ArrowReaderOptions::new()
2034 .with_virtual_columns(vec![row_number_field])
2035 .unwrap();
2036 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2037 .await
2038 .unwrap()
2039 .with_row_selection(selection)
2040 .with_batch_size(batch_size)
2041 .build()
2042 .expect("Could not create reader");
2043 reader.try_collect::<Vec<_>>().await.unwrap()
2044 })
2045 },
2046 );
2047 }
2048
2049 #[test]
2050 fn test_row_numbers_with_multiple_row_groups_and_filter() {
2051 test_row_numbers_with_multiple_row_groups_helper(
2052 true,
2053 |path, selection, row_filter, batch_size| {
2054 let runtime = tokio::runtime::Builder::new_current_thread()
2055 .enable_all()
2056 .build()
2057 .expect("Could not create runtime");
2058 runtime.block_on(async move {
2059 let file = tokio::fs::File::open(path).await.unwrap();
2060 let row_number_field = Arc::new(
2061 Field::new("row_number", DataType::Int64, false)
2062 .with_extension_type(RowNumber),
2063 );
2064 let options = ArrowReaderOptions::new()
2065 .with_virtual_columns(vec![row_number_field])
2066 .unwrap();
2067 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2068 .await
2069 .unwrap()
2070 .with_row_selection(selection)
2071 .with_row_filter(row_filter.expect("No row filter"))
2072 .with_batch_size(batch_size)
2073 .build()
2074 .expect("Could not create reader");
2075 reader.try_collect::<Vec<_>>().await.unwrap()
2076 })
2077 },
2078 );
2079 }
2080
2081 #[tokio::test]
2082 async fn test_nested_lists() -> Result<()> {
2083 let list_inner_field = Arc::new(Field::new("item", DataType::Float32, true));
2085 let table_schema = Arc::new(Schema::new(vec![
2086 Field::new("id", DataType::Int32, false),
2087 Field::new("vector", DataType::List(list_inner_field.clone()), true),
2088 ]));
2089
2090 let mut list_builder =
2091 ListBuilder::new(Float32Builder::new()).with_field(list_inner_field.clone());
2092 list_builder.values().append_slice(&[10.0, 10.0, 10.0]);
2093 list_builder.append(true);
2094 list_builder.values().append_slice(&[20.0, 20.0, 20.0]);
2095 list_builder.append(true);
2096 list_builder.values().append_slice(&[30.0, 30.0, 30.0]);
2097 list_builder.append(true);
2098 list_builder.values().append_slice(&[40.0, 40.0, 40.0]);
2099 list_builder.append(true);
2100 let list_array = list_builder.finish();
2101
2102 let data = vec![RecordBatch::try_new(
2103 table_schema.clone(),
2104 vec![
2105 Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
2106 Arc::new(list_array),
2107 ],
2108 )?];
2109
2110 let mut buffer = Vec::new();
2111 let mut writer = AsyncArrowWriter::try_new(&mut buffer, table_schema, None)?;
2112
2113 for batch in data {
2114 writer.write(&batch).await?;
2115 }
2116
2117 writer.close().await?;
2118
2119 let reader = TestReader::new(Bytes::from(buffer));
2120 let builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
2121
2122 let predicate = ArrowPredicateFn::new(ProjectionMask::all(), |batch| {
2123 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
2124 });
2125
2126 let projection_mask = ProjectionMask::all();
2127
2128 let mut stream = builder
2129 .with_row_filter(RowFilter::new(vec![Box::new(predicate)]))
2130 .with_projection(projection_mask)
2131 .build()?;
2132
2133 while let Some(batch) = stream.next().await {
2134 let _ = batch.unwrap(); }
2136
2137 Ok(())
2138 }
2139}