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::{ArrayRef, 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 async fn get_column_chunk_dictionary(
698 &mut self,
699 row_group_idx: usize,
700 column_idx: usize,
701 ) -> Result<Option<ArrayRef>> {
702 ParquetMetaDataReader::read_column_dictionary_async(
703 &mut self.input.0,
704 &self.metadata,
705 row_group_idx,
706 column_idx,
707 )
708 .await
709 }
710
711 pub fn build(self) -> Result<ParquetRecordBatchStream<T>> {
715 let Self {
716 input,
717 metadata,
718 schema,
719 fields,
720 batch_size,
721 row_group_plan,
722 projection,
723 filter,
724 row_selection_policy: selection_strategy,
725 limit,
726 offset,
727 metrics,
728 max_predicate_cache_size,
729 } = self;
730
731 let projection_len = projection.mask.as_ref().map_or(usize::MAX, |m| m.len());
734 let projected_fields = schema
735 .fields
736 .filter_leaves(|idx, _| idx < projection_len && projection.leaf_included(idx));
737 let projected_schema = Arc::new(Schema::new(projected_fields));
738
739 let decoder = ParquetPushDecoderBuilder {
740 input: PushDecoderInput::default(),
741 metadata,
742 schema,
743 fields,
744 projection,
745 filter,
746 row_group_plan,
747 row_selection_policy: selection_strategy,
748 batch_size,
749 limit,
750 offset,
751 metrics,
752 max_predicate_cache_size,
753 }
754 .build()?;
755
756 let request_state = RequestState::None { input: input.0 };
757
758 Ok(ParquetRecordBatchStream {
759 schema: projected_schema,
760 decoder,
761 request_state,
762 })
763 }
764}
765
766enum RequestState<T> {
770 None {
772 input: T,
773 },
774 Outstanding {
776 ranges: Vec<Range<u64>>,
778 future: BoxFuture<'static, Result<(T, Vec<Bytes>)>>,
783 },
784 Done,
785}
786
787impl<T> RequestState<T>
788where
789 T: AsyncFileReader + Unpin + Send + 'static,
790{
791 fn begin_request(mut input: T, ranges: Vec<Range<u64>>) -> Self {
793 let ranges_captured = ranges.clone();
794
795 let future = async move {
800 let data = input.get_byte_ranges(ranges_captured).await?;
801 Ok((input, data))
802 }
803 .boxed();
804 RequestState::Outstanding { ranges, future }
805 }
806}
807
808impl<T> std::fmt::Debug for RequestState<T> {
809 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
810 match self {
811 RequestState::None { input: _ } => f
812 .debug_struct("RequestState::None")
813 .field("input", &"...")
814 .finish(),
815 RequestState::Outstanding { ranges, .. } => f
816 .debug_struct("RequestState::Outstanding")
817 .field("ranges", &ranges)
818 .finish(),
819 RequestState::Done => {
820 write!(f, "RequestState::Done")
821 }
822 }
823 }
824}
825
826pub struct ParquetRecordBatchStream<T> {
846 schema: SchemaRef,
848 request_state: RequestState<T>,
850 decoder: ParquetPushDecoder,
852}
853
854impl<T> std::fmt::Debug for ParquetRecordBatchStream<T> {
855 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
856 f.debug_struct("ParquetRecordBatchStream")
857 .field("request_state", &self.request_state)
858 .finish()
859 }
860}
861
862impl<T> ParquetRecordBatchStream<T> {
863 pub fn schema(&self) -> &SchemaRef {
868 &self.schema
869 }
870}
871
872impl<T> ParquetRecordBatchStream<T>
873where
874 T: AsyncFileReader + Unpin + Send + 'static,
875{
876 pub async fn next_row_group(&mut self) -> Result<Option<ParquetRecordBatchReader>> {
890 loop {
891 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
894 match request_state {
895 RequestState::None { input } => {
897 match self.decoder.try_next_reader()? {
898 DecodeResult::NeedsData(ranges) => {
899 self.request_state = RequestState::begin_request(input, ranges);
900 }
902 DecodeResult::Data(reader) => {
903 self.request_state = RequestState::None { input };
904 return Ok(Some(reader));
905 }
906 DecodeResult::Finished => return Ok(None),
907 }
908 }
909 RequestState::Outstanding { ranges, future } => {
910 let (input, data) = future.await?;
911 self.decoder.push_ranges(ranges, data)?;
913 self.request_state = RequestState::None { input };
914 }
916 RequestState::Done => {
917 self.request_state = RequestState::Done;
918 return Ok(None);
919 }
920 }
921 }
922 }
923}
924
925impl<T> Stream for ParquetRecordBatchStream<T>
926where
927 T: AsyncFileReader + Unpin + Send + 'static,
928{
929 type Item = Result<RecordBatch>;
930 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
931 match self.poll_next_inner(cx) {
932 Ok(res) => {
933 res.map(|res| Ok(res).transpose())
936 }
937 Err(e) => {
938 self.request_state = RequestState::Done;
939 Poll::Ready(Some(Err(e)))
940 }
941 }
942 }
943}
944
945impl<T> ParquetRecordBatchStream<T>
946where
947 T: AsyncFileReader + Unpin + Send + 'static,
948{
949 fn poll_next_inner(&mut self, cx: &mut Context<'_>) -> Result<Poll<Option<RecordBatch>>> {
954 loop {
955 let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
956 match request_state {
957 RequestState::None { input } => {
958 match self.decoder.try_decode()? {
960 DecodeResult::NeedsData(ranges) => {
961 self.request_state = RequestState::begin_request(input, ranges);
962 }
964 DecodeResult::Data(batch) => {
965 self.request_state = RequestState::None { input };
966 return Ok(Poll::Ready(Some(batch)));
967 }
968 DecodeResult::Finished => {
969 self.request_state = RequestState::Done;
970 return Ok(Poll::Ready(None));
971 }
972 }
973 }
974 RequestState::Outstanding { ranges, mut future } => match future.poll_unpin(cx) {
975 Poll::Ready(result) => {
977 let (input, data) = result?;
978 self.decoder.push_ranges(ranges, data)?;
980 self.request_state = RequestState::None { input };
981 }
983 Poll::Pending => {
984 self.request_state = RequestState::Outstanding { ranges, future };
985 return Ok(Poll::Pending);
986 }
987 },
988 RequestState::Done => {
989 self.request_state = RequestState::Done;
991 return Ok(Poll::Ready(None));
992 }
993 }
994 }
995 }
996}
997
998#[cfg(test)]
999mod tests {
1000 use super::*;
1001 use crate::arrow::arrow_reader::tests::test_row_numbers_with_multiple_row_groups_helper;
1002 use crate::arrow::arrow_reader::{
1003 ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowFilter, RowSelection, RowSelector,
1004 };
1005 use crate::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
1006 use crate::arrow::schema::virtual_type::RowNumber;
1007 use crate::arrow::{ArrowWriter, AsyncArrowWriter, ProjectionMask};
1008 use crate::basic::Encoding;
1009 use crate::file::metadata::PageIndexPolicy;
1010 use crate::file::metadata::ParquetMetaDataReader;
1011 use crate::file::metadata::page_index::PageIndex;
1012 use crate::file::properties::WriterProperties;
1013 use arrow::compute::kernels::cmp::eq;
1014 use arrow::error::Result as ArrowResult;
1015 use arrow_array::builder::{Float32Builder, ListBuilder, StringBuilder};
1016 use arrow_array::cast::AsArray;
1017 use arrow_array::types::Int32Type;
1018 use arrow_array::{
1019 Array, ArrayRef, BinaryArray, BooleanArray, Int32Array, RecordBatchReader, Scalar,
1020 StringArray, StructArray, UInt64Array,
1021 };
1022 use arrow_schema::{DataType, Field, Schema};
1023 use futures::{StreamExt, TryStreamExt};
1024 use rand::{RngExt, rng};
1025 use std::collections::HashMap;
1026 use std::sync::{Arc, Mutex};
1027 use tempfile::tempfile;
1028
1029 #[derive(Clone)]
1030 struct TestReader {
1031 data: Bytes,
1032 metadata: Option<Arc<ParquetMetaData>>,
1033 requests: Arc<Mutex<Vec<Range<usize>>>>,
1034 }
1035
1036 impl TestReader {
1037 fn new(data: Bytes) -> Self {
1038 Self {
1039 data,
1040 metadata: Default::default(),
1041 requests: Default::default(),
1042 }
1043 }
1044 }
1045
1046 impl AsyncFileReader for TestReader {
1047 fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1048 let range = range.clone();
1049 self.requests
1050 .lock()
1051 .unwrap()
1052 .push(range.start as usize..range.end as usize);
1053 futures::future::ready(Ok(self
1054 .data
1055 .slice(range.start as usize..range.end as usize)))
1056 .boxed()
1057 }
1058
1059 fn get_metadata<'a>(
1060 &'a mut self,
1061 options: Option<&'a ArrowReaderOptions>,
1062 ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
1063 let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
1064 self.metadata = Some(Arc::new(
1065 metadata_reader.parse_and_finish(&self.data).unwrap(),
1066 ));
1067 futures::future::ready(Ok(self.metadata.clone().unwrap().clone())).boxed()
1068 }
1069 }
1070
1071 #[tokio::test]
1072 async fn test_async_reader() {
1073 let testdata = arrow::util::test_util::parquet_test_data();
1074 let path = format!("{testdata}/alltypes_plain.parquet");
1075 let data = Bytes::from(std::fs::read(path).unwrap());
1076
1077 let async_reader = TestReader::new(data.clone());
1078
1079 let requests = async_reader.requests.clone();
1080 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1081 .await
1082 .unwrap();
1083
1084 let metadata = builder.metadata().clone();
1085 assert_eq!(metadata.num_row_groups(), 1);
1086
1087 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1088 let stream = builder
1089 .with_projection(mask.clone())
1090 .with_batch_size(1024)
1091 .build()
1092 .unwrap();
1093
1094 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1095
1096 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1097 .unwrap()
1098 .with_projection(mask)
1099 .with_batch_size(104)
1100 .build()
1101 .unwrap()
1102 .collect::<ArrowResult<Vec<_>>>()
1103 .unwrap();
1104
1105 assert_eq!(async_batches, sync_batches);
1106
1107 let requests = requests.lock().unwrap();
1108 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1109 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1110
1111 assert_eq!(
1112 &requests[..],
1113 &[
1114 offset_1 as usize..(offset_1 + length_1) as usize,
1115 offset_2 as usize..(offset_2 + length_2) as usize
1116 ]
1117 );
1118 }
1119
1120 #[tokio::test]
1121 async fn test_async_reader_row_group_local_selections() {
1122 let batch = RecordBatch::try_from_iter([(
1123 "a",
1124 Arc::new(Int32Array::from_iter_values(0..6)) as ArrayRef,
1125 )])
1126 .unwrap();
1127 let mut data = Vec::new();
1128 let properties = WriterProperties::builder()
1129 .set_max_row_group_row_count(Some(3))
1130 .build();
1131 let mut writer = ArrowWriter::try_new(&mut data, batch.schema(), Some(properties)).unwrap();
1132 writer.write(&batch).unwrap();
1133 writer.close().unwrap();
1134
1135 let stream = ParquetRecordBatchStreamBuilder::new(TestReader::new(data.into()))
1136 .await
1137 .unwrap()
1138 .with_row_group_selections(vec![
1139 RowGroupSelection::new(1, Some(RowSelection::from(vec![RowSelector::select(1)]))),
1140 RowGroupSelection::new(
1141 0,
1142 Some(RowSelection::from(vec![
1143 RowSelector::skip(1),
1144 RowSelector::select(2),
1145 ])),
1146 ),
1147 ])
1148 .build()
1149 .unwrap();
1150
1151 let batches: Vec<_> = stream.try_collect().await.unwrap();
1152 assert_eq!(batches.len(), 2);
1153 assert_eq!(
1154 batches[0].column(0).as_primitive::<Int32Type>().values(),
1155 &[3]
1156 );
1157 assert_eq!(
1158 batches[1].column(0).as_primitive::<Int32Type>().values(),
1159 &[1, 2]
1160 );
1161 }
1162
1163 #[tokio::test]
1164 async fn test_get_column_chunk_dictionary() {
1165 let schema = Arc::new(Schema::new(vec![Field::new("s", DataType::Utf8, false)]));
1166 let values: Vec<&str> = ["alpha", "beta", "gamma"]
1167 .iter()
1168 .copied()
1169 .cycle()
1170 .take(30)
1171 .collect();
1172 let array: ArrayRef = Arc::new(StringArray::from(values));
1173 let batch = RecordBatch::try_new(schema.clone(), vec![array]).unwrap();
1174
1175 let props = WriterProperties::builder()
1176 .set_dictionary_enabled(true)
1177 .build();
1178 let mut buf = Vec::new();
1179 {
1180 let mut writer = ArrowWriter::try_new(&mut buf, schema, Some(props)).unwrap();
1181 writer.write(&batch).unwrap();
1182 writer.close().unwrap();
1183 }
1184 let data = Bytes::from(buf);
1185
1186 let direct_metadata = ParquetMetaDataReader::new()
1187 .parse_and_finish(&data)
1188 .unwrap();
1189 let mut direct_reader = TestReader::new(data.clone());
1190 let direct = ParquetMetaDataReader::read_column_dictionary_async(
1191 &mut direct_reader,
1192 &direct_metadata,
1193 0,
1194 0,
1195 )
1196 .await
1197 .unwrap()
1198 .unwrap();
1199 let direct = direct.as_any().downcast_ref::<BinaryArray>().unwrap();
1200 assert_eq!(direct.value(0), b"alpha");
1201
1202 let async_reader = TestReader::new(data);
1203 let mut builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1204 .await
1205 .unwrap();
1206
1207 let dictionary = builder
1208 .get_column_chunk_dictionary(0, 0)
1209 .await
1210 .unwrap()
1211 .unwrap();
1212 let dictionary = dictionary.as_any().downcast_ref::<BinaryArray>().unwrap();
1213 let dictionary_values: Vec<&[u8]> = dictionary.iter().map(|v| v.unwrap()).collect();
1214 assert_eq!(
1215 dictionary_values,
1216 vec![b"alpha".as_slice(), b"beta", b"gamma"]
1217 );
1218 }
1219
1220 #[tokio::test]
1223 async fn test_dictionary_selects_row_groups_without_reading_skipped_data() {
1224 let schema = Arc::new(Schema::new(vec![Field::new("s", DataType::Utf8, false)]));
1228 let row_group_values = ["skip", "target"];
1229 let props = WriterProperties::builder()
1230 .set_dictionary_enabled(true)
1231 .build();
1232 let mut buf = Vec::new();
1233 {
1234 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), Some(props)).unwrap();
1235 for value in row_group_values {
1236 let array: ArrayRef = Arc::new(StringArray::from(vec![value; 30]));
1237 let batch = RecordBatch::try_new(schema.clone(), vec![array]).unwrap();
1238 writer.write(&batch).unwrap();
1239 writer.flush().unwrap();
1240 }
1241 writer.close().unwrap();
1242 }
1243
1244 let async_reader = TestReader::new(Bytes::from(buf));
1245 let requests = async_reader.requests.clone();
1248 let mut builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1249 .await
1250 .unwrap();
1251 let metadata = builder.metadata().clone();
1252 assert_eq!(metadata.num_row_groups(), 2);
1253
1254 let mut selected_row_groups = Vec::new();
1257 let mut dictionary_ranges = Vec::new();
1258 for row_group_idx in 0..metadata.num_row_groups() {
1259 let column = metadata.row_group(row_group_idx).column(0);
1260 let encoding_mask = column.page_encoding_stats_mask().unwrap();
1261 assert!(
1262 encoding_mask.is_only(Encoding::PLAIN_DICTIONARY)
1263 || encoding_mask.is_only(Encoding::RLE_DICTIONARY)
1264 );
1265
1266 let dictionary_start = column.dictionary_page_offset().unwrap() as usize;
1267 let data_start = column.data_page_offset() as usize;
1268 dictionary_ranges.push(dictionary_start..data_start);
1269
1270 let dictionary = builder
1271 .get_column_chunk_dictionary(row_group_idx, 0)
1272 .await
1273 .unwrap()
1274 .unwrap();
1275 let dictionary = dictionary.as_binary::<i32>();
1276 if dictionary
1277 .iter()
1278 .any(|value| value == Some(b"target".as_slice()))
1279 {
1280 selected_row_groups.push(row_group_idx);
1281 }
1282 }
1283 assert_eq!(selected_row_groups, vec![1]);
1284
1285 let batches: Vec<_> = builder
1287 .with_row_groups(selected_row_groups)
1288 .build()
1289 .unwrap()
1290 .try_collect()
1291 .await
1292 .unwrap();
1293 let values: Vec<_> = batches
1294 .iter()
1295 .flat_map(|batch| batch.column(0).as_string::<i32>().iter())
1296 .collect();
1297 assert_eq!(values, vec![Some("target"); 30]);
1298
1299 let skipped_column = metadata.row_group(0).column(0);
1302 let (skipped_start, skipped_len) = skipped_column.byte_range();
1303 let skipped_data_range =
1304 skipped_column.data_page_offset() as usize..(skipped_start + skipped_len) as usize;
1305 let requests = requests.lock().unwrap();
1306 for dictionary_range in dictionary_ranges {
1307 assert!(requests.contains(&dictionary_range));
1308 }
1309 assert!(requests.iter().all(|request| {
1310 request.end <= skipped_data_range.start || request.start >= skipped_data_range.end
1311 }));
1312 }
1313
1314 #[tokio::test]
1315 async fn test_async_reader_with_next_row_group() {
1316 let testdata = arrow::util::test_util::parquet_test_data();
1317 let path = format!("{testdata}/alltypes_plain.parquet");
1318 let data = Bytes::from(std::fs::read(path).unwrap());
1319
1320 let async_reader = TestReader::new(data.clone());
1321
1322 let requests = async_reader.requests.clone();
1323 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1324 .await
1325 .unwrap();
1326
1327 let metadata = builder.metadata().clone();
1328 assert_eq!(metadata.num_row_groups(), 1);
1329
1330 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1331 let mut stream = builder
1332 .with_projection(mask.clone())
1333 .with_batch_size(1024)
1334 .build()
1335 .unwrap();
1336
1337 let mut readers = vec![];
1338 while let Some(reader) = stream.next_row_group().await.unwrap() {
1339 readers.push(reader);
1340 }
1341
1342 let async_batches: Vec<_> = readers
1343 .into_iter()
1344 .flat_map(|r| r.map(|v| v.unwrap()).collect::<Vec<_>>())
1345 .collect();
1346
1347 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1348 .unwrap()
1349 .with_projection(mask)
1350 .with_batch_size(104)
1351 .build()
1352 .unwrap()
1353 .collect::<ArrowResult<Vec<_>>>()
1354 .unwrap();
1355
1356 assert_eq!(async_batches, sync_batches);
1357
1358 let requests = requests.lock().unwrap();
1359 let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1360 let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1361
1362 assert_eq!(
1363 &requests[..],
1364 &[
1365 offset_1 as usize..(offset_1 + length_1) as usize,
1366 offset_2 as usize..(offset_2 + length_2) as usize
1367 ]
1368 );
1369 }
1370
1371 #[tokio::test]
1372 async fn test_async_reader_with_index() {
1373 let testdata = arrow::util::test_util::parquet_test_data();
1374 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1375 let data = Bytes::from(std::fs::read(path).unwrap());
1376
1377 let async_reader = TestReader::new(data.clone());
1378
1379 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1380 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1381 .await
1382 .unwrap();
1383
1384 let metadata_with_index = builder.metadata();
1386 assert_eq!(metadata_with_index.num_row_groups(), 1);
1387
1388 let page_index = metadata_with_index
1390 .page_index()
1391 .expect("page index should be present");
1392 assert!(page_index.is_complete());
1393 let num_rowgroups = metadata_with_index.num_row_groups();
1394 let num_columns = metadata_with_index
1395 .file_metadata()
1396 .schema_descr()
1397 .num_columns();
1398 for rgidx in 0..num_rowgroups {
1399 for colidx in 0..num_columns {
1401 assert!(page_index.offset_index(rgidx, colidx).is_some());
1402 }
1403 }
1404
1405 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1406 let stream = builder
1407 .with_projection(mask.clone())
1408 .with_batch_size(1024)
1409 .build()
1410 .unwrap();
1411
1412 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1413
1414 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1415 .unwrap()
1416 .with_projection(mask)
1417 .with_batch_size(1024)
1418 .build()
1419 .unwrap()
1420 .collect::<ArrowResult<Vec<_>>>()
1421 .unwrap();
1422
1423 assert_eq!(async_batches, sync_batches);
1424 }
1425
1426 #[tokio::test]
1427 async fn test_async_reader_with_limit() {
1428 let testdata = arrow::util::test_util::parquet_test_data();
1429 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1430 let data = Bytes::from(std::fs::read(path).unwrap());
1431
1432 let metadata = ParquetMetaDataReader::new()
1433 .parse_and_finish(&data)
1434 .unwrap();
1435 let metadata = Arc::new(metadata);
1436
1437 assert_eq!(metadata.num_row_groups(), 1);
1438
1439 let async_reader = TestReader::new(data.clone());
1440
1441 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1442 .await
1443 .unwrap();
1444
1445 assert_eq!(builder.metadata().num_row_groups(), 1);
1446
1447 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1448 let stream = builder
1449 .with_projection(mask.clone())
1450 .with_batch_size(1024)
1451 .with_limit(1)
1452 .build()
1453 .unwrap();
1454
1455 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1456
1457 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1458 .unwrap()
1459 .with_projection(mask)
1460 .with_batch_size(1024)
1461 .with_limit(1)
1462 .build()
1463 .unwrap()
1464 .collect::<ArrowResult<Vec<_>>>()
1465 .unwrap();
1466
1467 assert_eq!(async_batches, sync_batches);
1468 }
1469
1470 #[tokio::test]
1471 async fn test_async_reader_skip_pages() {
1472 let testdata = arrow::util::test_util::parquet_test_data();
1473 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1474 let data = Bytes::from(std::fs::read(path).unwrap());
1475
1476 let async_reader = TestReader::new(data.clone());
1477
1478 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1479 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1480 .await
1481 .unwrap();
1482
1483 assert_eq!(builder.metadata().num_row_groups(), 1);
1484
1485 let selection = RowSelection::from(vec![
1486 RowSelector::skip(21), RowSelector::select(21), RowSelector::skip(41), RowSelector::select(41), RowSelector::skip(25), RowSelector::select(25), RowSelector::skip(7116), RowSelector::select(10), ]);
1495
1496 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![9]);
1497
1498 let stream = builder
1499 .with_projection(mask.clone())
1500 .with_row_selection(selection.clone())
1501 .build()
1502 .expect("building stream");
1503
1504 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1505
1506 let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1507 .unwrap()
1508 .with_projection(mask)
1509 .with_batch_size(1024)
1510 .with_row_selection(selection)
1511 .build()
1512 .unwrap()
1513 .collect::<ArrowResult<Vec<_>>>()
1514 .unwrap();
1515
1516 assert_eq!(async_batches, sync_batches);
1517 }
1518
1519 #[tokio::test]
1520 async fn test_fuzz_async_reader_selection() {
1521 let testdata = arrow::util::test_util::parquet_test_data();
1522 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1523 let data = Bytes::from(std::fs::read(path).unwrap());
1524
1525 let mut rand = rng();
1526
1527 for _ in 0..100 {
1528 let mut expected_rows = 0;
1529 let mut total_rows = 0;
1530 let mut skip = false;
1531 let mut selectors = vec![];
1532
1533 while total_rows < 7300 {
1534 let row_count: usize = rand.random_range(1..100);
1535
1536 let row_count = row_count.min(7300 - total_rows);
1537
1538 selectors.push(RowSelector { row_count, skip });
1539
1540 total_rows += row_count;
1541 if !skip {
1542 expected_rows += row_count;
1543 }
1544
1545 skip = !skip;
1546 }
1547
1548 let selection = RowSelection::from(selectors);
1549
1550 let async_reader = TestReader::new(data.clone());
1551
1552 let options =
1553 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1554 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1555 .await
1556 .unwrap();
1557
1558 assert_eq!(builder.metadata().num_row_groups(), 1);
1559
1560 let col_idx: usize = rand.random_range(0..13);
1561 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1562
1563 let stream = builder
1564 .with_projection(mask.clone())
1565 .with_row_selection(selection.clone())
1566 .build()
1567 .expect("building stream");
1568
1569 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1570
1571 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1572
1573 assert_eq!(actual_rows, expected_rows);
1574 }
1575 }
1576
1577 #[tokio::test]
1578 async fn test_async_reader_zero_row_selector() {
1579 let testdata = arrow::util::test_util::parquet_test_data();
1581 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1582 let data = Bytes::from(std::fs::read(path).unwrap());
1583
1584 let mut rand = rng();
1585
1586 let mut expected_rows = 0;
1587 let mut total_rows = 0;
1588 let mut skip = false;
1589 let mut selectors = vec![];
1590
1591 selectors.push(RowSelector {
1592 row_count: 0,
1593 skip: false,
1594 });
1595
1596 while total_rows < 7300 {
1597 let row_count: usize = rand.random_range(1..100);
1598
1599 let row_count = row_count.min(7300 - total_rows);
1600
1601 selectors.push(RowSelector { row_count, skip });
1602
1603 total_rows += row_count;
1604 if !skip {
1605 expected_rows += row_count;
1606 }
1607
1608 skip = !skip;
1609 }
1610
1611 let selection = RowSelection::from(selectors);
1612
1613 let async_reader = TestReader::new(data.clone());
1614
1615 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1616 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1617 .await
1618 .unwrap();
1619
1620 assert_eq!(builder.metadata().num_row_groups(), 1);
1621
1622 let col_idx: usize = rand.random_range(0..13);
1623 let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1624
1625 let stream = builder
1626 .with_projection(mask.clone())
1627 .with_row_selection(selection.clone())
1628 .build()
1629 .expect("building stream");
1630
1631 let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1632
1633 let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1634
1635 assert_eq!(actual_rows, expected_rows);
1636 }
1637
1638 #[tokio::test]
1639 async fn test_limit_multiple_row_groups() {
1640 let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]);
1641 let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]);
1642 let c = Int32Array::from_iter(0..6);
1643 let data = RecordBatch::try_from_iter([
1644 ("a", Arc::new(a) as ArrayRef),
1645 ("b", Arc::new(b) as ArrayRef),
1646 ("c", Arc::new(c) as ArrayRef),
1647 ])
1648 .unwrap();
1649
1650 let mut buf = Vec::with_capacity(1024);
1651 let props = WriterProperties::builder()
1652 .set_max_row_group_row_count(Some(3))
1653 .build();
1654 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1655 writer.write(&data).unwrap();
1656 writer.close().unwrap();
1657
1658 let data: Bytes = buf.into();
1659 let metadata = ParquetMetaDataReader::new()
1660 .parse_and_finish(&data)
1661 .unwrap();
1662
1663 assert_eq!(metadata.num_row_groups(), 2);
1664
1665 let test = TestReader::new(data);
1666
1667 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1668 .await
1669 .unwrap()
1670 .with_batch_size(1024)
1671 .with_limit(4)
1672 .build()
1673 .unwrap();
1674
1675 let batches: Vec<_> = stream.try_collect().await.unwrap();
1676 assert_eq!(batches.len(), 2);
1678
1679 let batch = &batches[0];
1680 assert_eq!(batch.num_rows(), 3);
1682 assert_eq!(batch.num_columns(), 3);
1683 let col2 = batch.column(2).as_primitive::<Int32Type>();
1684 assert_eq!(col2.values(), &[0, 1, 2]);
1685
1686 let batch = &batches[1];
1687 assert_eq!(batch.num_rows(), 1);
1689 assert_eq!(batch.num_columns(), 3);
1690 let col2 = batch.column(2).as_primitive::<Int32Type>();
1691 assert_eq!(col2.values(), &[3]);
1692
1693 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1694 .await
1695 .unwrap()
1696 .with_offset(2)
1697 .with_limit(3)
1698 .build()
1699 .unwrap();
1700
1701 let batches: Vec<_> = stream.try_collect().await.unwrap();
1702 assert_eq!(batches.len(), 2);
1704
1705 let batch = &batches[0];
1706 assert_eq!(batch.num_rows(), 1);
1708 assert_eq!(batch.num_columns(), 3);
1709 let col2 = batch.column(2).as_primitive::<Int32Type>();
1710 assert_eq!(col2.values(), &[2]);
1711
1712 let batch = &batches[1];
1713 assert_eq!(batch.num_rows(), 2);
1715 assert_eq!(batch.num_columns(), 3);
1716 let col2 = batch.column(2).as_primitive::<Int32Type>();
1717 assert_eq!(col2.values(), &[3, 4]);
1718
1719 let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1720 .await
1721 .unwrap()
1722 .with_offset(4)
1723 .with_limit(20)
1724 .build()
1725 .unwrap();
1726
1727 let batches: Vec<_> = stream.try_collect().await.unwrap();
1728 assert_eq!(batches.len(), 1);
1730
1731 let batch = &batches[0];
1732 assert_eq!(batch.num_rows(), 2);
1734 assert_eq!(batch.num_columns(), 3);
1735 let col2 = batch.column(2).as_primitive::<Int32Type>();
1736 assert_eq!(col2.values(), &[4, 5]);
1737 }
1738
1739 #[tokio::test]
1740 async fn test_batch_size_overallocate() {
1741 let testdata = arrow::util::test_util::parquet_test_data();
1742 let path = format!("{testdata}/alltypes_plain.parquet");
1744 let data = Bytes::from(std::fs::read(path).unwrap());
1745
1746 let async_reader = TestReader::new(data.clone());
1747
1748 let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1749 .await
1750 .unwrap();
1751
1752 let file_rows = builder.metadata().file_metadata().num_rows() as usize;
1753
1754 let builder = builder
1755 .with_projection(ProjectionMask::all())
1756 .with_batch_size(1024);
1757
1758 assert_ne!(1024, file_rows);
1761 assert_eq!(builder.batch_size, file_rows);
1762
1763 let _stream = builder.build().unwrap();
1764 }
1765
1766 #[tokio::test]
1767 async fn test_parquet_record_batch_stream_schema() {
1768 fn get_all_field_names(schema: &Schema) -> Vec<&String> {
1769 schema.flattened_fields().iter().map(|f| f.name()).collect()
1770 }
1771
1772 let mut metadata = HashMap::with_capacity(1);
1781 metadata.insert("key".to_string(), "value".to_string());
1782
1783 let nested_struct_array = StructArray::from(vec![
1784 (
1785 Arc::new(Field::new("d", DataType::Utf8, true)),
1786 Arc::new(StringArray::from(vec!["a", "b"])) as ArrayRef,
1787 ),
1788 (
1789 Arc::new(Field::new("e", DataType::Utf8, true)),
1790 Arc::new(StringArray::from(vec!["c", "d"])) as ArrayRef,
1791 ),
1792 ]);
1793 let struct_array = StructArray::from(vec![
1794 (
1795 Arc::new(Field::new("a", DataType::Int32, true)),
1796 Arc::new(Int32Array::from(vec![-1, 1])) as ArrayRef,
1797 ),
1798 (
1799 Arc::new(Field::new("b", DataType::UInt64, true)),
1800 Arc::new(UInt64Array::from(vec![1, 2])) as ArrayRef,
1801 ),
1802 (
1803 Arc::new(Field::new(
1804 "c",
1805 nested_struct_array.data_type().clone(),
1806 true,
1807 )),
1808 Arc::new(nested_struct_array) as ArrayRef,
1809 ),
1810 ]);
1811
1812 let schema =
1813 Arc::new(Schema::new(struct_array.fields().clone()).with_metadata(metadata.clone()));
1814 let record_batch = RecordBatch::from(struct_array)
1815 .with_schema(schema.clone())
1816 .unwrap();
1817
1818 let mut file = tempfile().unwrap();
1820 let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1821 writer.write(&record_batch).unwrap();
1822 writer.close().unwrap();
1823
1824 let all_fields = ["a", "b", "c", "d", "e"];
1825 let projections = [
1827 (vec![], vec![]),
1828 (vec![0], vec!["a"]),
1829 (vec![0, 1], vec!["a", "b"]),
1830 (vec![0, 1, 2], vec!["a", "b", "c", "d"]),
1831 (vec![0, 1, 2, 3], vec!["a", "b", "c", "d", "e"]),
1832 ];
1833
1834 for (indices, expected_projected_names) in projections {
1836 let assert_schemas = |builder: SchemaRef, reader: SchemaRef, batch: SchemaRef| {
1837 assert_eq!(get_all_field_names(&builder), all_fields);
1839 assert_eq!(builder.metadata, metadata);
1840 assert_eq!(get_all_field_names(&reader), expected_projected_names);
1842 assert_eq!(reader.metadata, HashMap::default());
1843 assert_eq!(get_all_field_names(&batch), expected_projected_names);
1844 assert_eq!(batch.metadata, HashMap::default());
1845 };
1846
1847 let builder =
1848 ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1849 let sync_builder_schema = builder.schema().clone();
1850 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices.clone());
1851 let mut reader = builder.with_projection(mask).build().unwrap();
1852 let sync_reader_schema = reader.schema();
1853 let batch = reader.next().unwrap().unwrap();
1854 let sync_batch_schema = batch.schema();
1855 assert_schemas(sync_builder_schema, sync_reader_schema, sync_batch_schema);
1856
1857 let file = tokio::fs::File::from(file.try_clone().unwrap());
1859 let builder = ParquetRecordBatchStreamBuilder::new(file).await.unwrap();
1860 let async_builder_schema = builder.schema().clone();
1861 let mask = ProjectionMask::leaves(builder.parquet_schema(), indices);
1862 let mut reader = builder.with_projection(mask).build().unwrap();
1863 let async_reader_schema = reader.schema().clone();
1864 let batch = reader.next().await.unwrap().unwrap();
1865 let async_batch_schema = batch.schema();
1866 assert_schemas(
1867 async_builder_schema,
1868 async_reader_schema,
1869 async_batch_schema,
1870 );
1871 }
1872 }
1873
1874 #[tokio::test]
1875 async fn test_nested_skip() {
1876 let schema = Arc::new(Schema::new(vec![
1877 Field::new("col_1", DataType::UInt64, false),
1878 Field::new_list("col_2", Field::new_list_field(DataType::Utf8, true), true),
1879 ]));
1880
1881 let props = WriterProperties::builder()
1883 .set_data_page_row_count_limit(256)
1884 .set_write_batch_size(256)
1885 .set_max_row_group_row_count(Some(1024));
1886
1887 let mut file = tempfile().unwrap();
1889 let mut writer =
1890 ArrowWriter::try_new(&mut file, schema.clone(), Some(props.build())).unwrap();
1891
1892 let mut builder = ListBuilder::new(StringBuilder::new());
1893 for id in 0..1024 {
1894 match id % 3 {
1895 0 => builder.append_value([Some("val_1".to_string()), Some(format!("id_{id}"))]),
1896 1 => builder.append_value([Some(format!("id_{id}"))]),
1897 _ => builder.append_null(),
1898 }
1899 }
1900 let refs = vec![
1901 Arc::new(UInt64Array::from_iter_values(0..1024)) as ArrayRef,
1902 Arc::new(builder.finish()) as ArrayRef,
1903 ];
1904
1905 let batch = RecordBatch::try_new(schema.clone(), refs).unwrap();
1906 writer.write(&batch).unwrap();
1907 writer.close().unwrap();
1908
1909 let selections = [
1910 RowSelection::from(vec![
1911 RowSelector::skip(313),
1912 RowSelector::select(1),
1913 RowSelector::skip(709),
1914 RowSelector::select(1),
1915 ]),
1916 RowSelection::from(vec![
1917 RowSelector::skip(255),
1918 RowSelector::select(1),
1919 RowSelector::skip(767),
1920 RowSelector::select(1),
1921 ]),
1922 RowSelection::from(vec![
1923 RowSelector::select(255),
1924 RowSelector::skip(1),
1925 RowSelector::select(767),
1926 RowSelector::skip(1),
1927 ]),
1928 RowSelection::from(vec![
1929 RowSelector::skip(254),
1930 RowSelector::select(1),
1931 RowSelector::select(1),
1932 RowSelector::skip(767),
1933 RowSelector::select(1),
1934 ]),
1935 ];
1936
1937 for selection in selections {
1938 let expected = selection.row_count();
1939 let mut reader = ParquetRecordBatchStreamBuilder::new_with_options(
1941 tokio::fs::File::from_std(file.try_clone().unwrap()),
1942 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
1943 )
1944 .await
1945 .unwrap();
1946
1947 reader = reader.with_row_selection(selection);
1948
1949 let mut stream = reader.build().unwrap();
1950
1951 let mut total_rows = 0;
1952 while let Some(rb) = stream.next().await {
1953 let rb = rb.unwrap();
1954 total_rows += rb.num_rows();
1955 }
1956 assert_eq!(total_rows, expected);
1957 }
1958 }
1959
1960 #[tokio::test]
1961 async fn empty_offset_index_doesnt_panic_in_read_row_group() {
1962 use tokio::fs::File;
1963 let testdata = arrow::util::test_util::parquet_test_data();
1964 let path = format!("{testdata}/alltypes_plain.parquet");
1965 let mut file = File::open(&path).await.unwrap();
1966 let file_size = file.metadata().await.unwrap().len();
1967 let mut metadata = ParquetMetaDataReader::new()
1968 .with_page_index_policy(PageIndexPolicy::Required)
1969 .load_and_finish(&mut file, file_size)
1970 .await
1971 .unwrap();
1972
1973 let page_index = PageIndex::new(None, Some(vec![]));
1974 metadata.set_page_index(Some(Arc::new(page_index)));
1975 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1976 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1977 let reader =
1978 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1979 .build()
1980 .unwrap();
1981
1982 let result = reader.try_collect::<Vec<_>>().await.unwrap();
1983 assert_eq!(result.len(), 1);
1984 }
1985
1986 #[tokio::test]
1987 async fn non_empty_offset_index_doesnt_panic_in_read_row_group() {
1988 use tokio::fs::File;
1989 let testdata = arrow::util::test_util::parquet_test_data();
1990 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1991 let mut file = File::open(&path).await.unwrap();
1992 let file_size = file.metadata().await.unwrap().len();
1993 let metadata = ParquetMetaDataReader::new()
1994 .with_page_index_policy(PageIndexPolicy::Required)
1995 .load_and_finish(&mut file, file_size)
1996 .await
1997 .unwrap();
1998
1999 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
2000 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
2001 let reader =
2002 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
2003 .build()
2004 .unwrap();
2005
2006 let result = reader.try_collect::<Vec<_>>().await.unwrap();
2007 assert_eq!(result.len(), 8);
2008 }
2009
2010 #[tokio::test]
2011 async fn empty_offset_index_doesnt_panic_in_column_chunks() {
2012 use tempfile::TempDir;
2013 use tokio::fs::File;
2014 fn write_metadata_to_local_file(
2015 metadata: ParquetMetaData,
2016 file: impl AsRef<std::path::Path>,
2017 ) {
2018 use crate::file::metadata::ParquetMetaDataWriter;
2019 use std::fs::File;
2020 let file = File::create(file).unwrap();
2021 ParquetMetaDataWriter::new(file, &metadata)
2022 .finish()
2023 .unwrap()
2024 }
2025
2026 fn read_metadata_from_local_file(file: impl AsRef<std::path::Path>) -> ParquetMetaData {
2027 use std::fs::File;
2028 let file = File::open(file).unwrap();
2029 ParquetMetaDataReader::new()
2030 .with_page_index_policy(PageIndexPolicy::Required)
2031 .parse_and_finish(&file)
2032 .unwrap()
2033 }
2034
2035 let testdata = arrow::util::test_util::parquet_test_data();
2036 let path = format!("{testdata}/alltypes_plain.parquet");
2037 let mut file = File::open(&path).await.unwrap();
2038 let file_size = file.metadata().await.unwrap().len();
2039 let metadata = ParquetMetaDataReader::new()
2040 .with_page_index_policy(PageIndexPolicy::Required)
2041 .load_and_finish(&mut file, file_size)
2042 .await
2043 .unwrap();
2044
2045 let tempdir = TempDir::new().unwrap();
2046 let metadata_path = tempdir.path().join("thrift_metadata.dat");
2047 write_metadata_to_local_file(metadata, &metadata_path);
2048 let metadata = read_metadata_from_local_file(&metadata_path);
2049
2050 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
2051 let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
2052 let reader =
2053 ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
2054 .build()
2055 .unwrap();
2056
2057 let result = reader.try_collect::<Vec<_>>().await.unwrap();
2059 assert_eq!(result.len(), 1);
2060 }
2061
2062 #[tokio::test]
2063 async fn test_cached_array_reader_sparse_offset_error() {
2064 use futures::TryStreamExt;
2065
2066 use crate::arrow::arrow_reader::{ArrowPredicateFn, RowFilter, RowSelection, RowSelector};
2067 use arrow_array::{BooleanArray, RecordBatch};
2068
2069 let testdata = arrow::util::test_util::parquet_test_data();
2070 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
2071 let data = Bytes::from(std::fs::read(path).unwrap());
2072
2073 let async_reader = TestReader::new(data);
2074
2075 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
2077 let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
2078 .await
2079 .unwrap();
2080
2081 let selection = RowSelection::from(vec![RowSelector::skip(22), RowSelector::select(3)]);
2085
2086 let parquet_schema = builder.parquet_schema();
2090 let proj = ProjectionMask::leaves(parquet_schema, vec![0]);
2091 let always_true = ArrowPredicateFn::new(proj.clone(), |batch: RecordBatch| {
2092 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
2093 });
2094 let filter = RowFilter::new(vec![Box::new(always_true)]);
2095
2096 let stream = builder
2099 .with_batch_size(8)
2100 .with_projection(proj)
2101 .with_row_selection(selection)
2102 .with_row_filter(filter)
2103 .build()
2104 .unwrap();
2105
2106 let _result: Vec<_> = stream.try_collect().await.unwrap();
2109 }
2110
2111 #[tokio::test]
2112 async fn test_predicate_cache_disabled() {
2113 let k = Int32Array::from_iter_values(0..10);
2114 let data = RecordBatch::try_from_iter([("k", Arc::new(k) as ArrayRef)]).unwrap();
2115
2116 let mut buf = Vec::new();
2117 let props = WriterProperties::builder()
2119 .set_data_page_row_count_limit(1)
2120 .set_write_batch_size(1)
2121 .set_max_row_group_row_count(Some(10))
2122 .set_write_page_header_statistics(true)
2123 .build();
2124 let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
2125 writer.write(&data).unwrap();
2126 writer.close().unwrap();
2127
2128 let data = Bytes::from(buf);
2129 let metadata = ParquetMetaDataReader::new()
2130 .with_page_index_policy(PageIndexPolicy::Required)
2131 .parse_and_finish(&data)
2132 .unwrap();
2133 let parquet_schema = metadata.file_metadata().schema_descr_ptr();
2134
2135 let build_filter = || {
2137 let scalar = Int32Array::from_iter_values([5]);
2138 let predicate = ArrowPredicateFn::new(
2139 ProjectionMask::leaves(&parquet_schema, vec![0]),
2140 move |batch| eq(batch.column(0), &Scalar::new(&scalar)),
2141 );
2142 RowFilter::new(vec![Box::new(predicate)])
2143 };
2144
2145 let selection = RowSelection::from(vec![RowSelector::skip(5), RowSelector::select(1)]);
2147
2148 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
2149 let reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
2150
2151 let reader_with_cache = TestReader::new(data.clone());
2153 let requests_with_cache = reader_with_cache.requests.clone();
2154 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
2155 reader_with_cache,
2156 reader_metadata.clone(),
2157 )
2158 .with_batch_size(1000)
2159 .with_row_selection(selection.clone())
2160 .with_row_filter(build_filter())
2161 .build()
2162 .unwrap();
2163 let batches_with_cache: Vec<_> = stream.try_collect().await.unwrap();
2164
2165 let reader_without_cache = TestReader::new(data);
2167 let requests_without_cache = reader_without_cache.requests.clone();
2168 let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
2169 reader_without_cache,
2170 reader_metadata,
2171 )
2172 .with_batch_size(1000)
2173 .with_row_selection(selection)
2174 .with_row_filter(build_filter())
2175 .with_max_predicate_cache_size(0) .build()
2177 .unwrap();
2178 let batches_without_cache: Vec<_> = stream.try_collect().await.unwrap();
2179
2180 assert_eq!(batches_with_cache, batches_without_cache);
2181
2182 let requests_with_cache = requests_with_cache.lock().unwrap();
2183 let requests_without_cache = requests_without_cache.lock().unwrap();
2184
2185 assert_eq!(requests_with_cache.len(), 11);
2187 assert_eq!(requests_without_cache.len(), 2);
2188
2189 assert_eq!(
2191 requests_with_cache.iter().map(|r| r.len()).sum::<usize>(),
2192 433
2193 );
2194 assert_eq!(
2195 requests_without_cache
2196 .iter()
2197 .map(|r| r.len())
2198 .sum::<usize>(),
2199 92
2200 );
2201 }
2202
2203 #[test]
2204 fn test_row_numbers_with_multiple_row_groups() {
2205 test_row_numbers_with_multiple_row_groups_helper(
2206 false,
2207 |path, selection, _row_filter, batch_size| {
2208 let runtime = tokio::runtime::Builder::new_current_thread()
2209 .enable_all()
2210 .build()
2211 .expect("Could not create runtime");
2212 runtime.block_on(async move {
2213 let file = tokio::fs::File::open(path).await.unwrap();
2214 let row_number_field = Arc::new(
2215 Field::new("row_number", DataType::Int64, false)
2216 .with_extension_type(RowNumber),
2217 );
2218 let options = ArrowReaderOptions::new()
2219 .with_virtual_columns(vec![row_number_field])
2220 .unwrap();
2221 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2222 .await
2223 .unwrap()
2224 .with_row_selection(selection)
2225 .with_batch_size(batch_size)
2226 .build()
2227 .expect("Could not create reader");
2228 reader.try_collect::<Vec<_>>().await.unwrap()
2229 })
2230 },
2231 );
2232 }
2233
2234 #[test]
2235 fn test_row_numbers_with_multiple_row_groups_and_filter() {
2236 test_row_numbers_with_multiple_row_groups_helper(
2237 true,
2238 |path, selection, row_filter, batch_size| {
2239 let runtime = tokio::runtime::Builder::new_current_thread()
2240 .enable_all()
2241 .build()
2242 .expect("Could not create runtime");
2243 runtime.block_on(async move {
2244 let file = tokio::fs::File::open(path).await.unwrap();
2245 let row_number_field = Arc::new(
2246 Field::new("row_number", DataType::Int64, false)
2247 .with_extension_type(RowNumber),
2248 );
2249 let options = ArrowReaderOptions::new()
2250 .with_virtual_columns(vec![row_number_field])
2251 .unwrap();
2252 let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2253 .await
2254 .unwrap()
2255 .with_row_selection(selection)
2256 .with_row_filter(row_filter.expect("No row filter"))
2257 .with_batch_size(batch_size)
2258 .build()
2259 .expect("Could not create reader");
2260 reader.try_collect::<Vec<_>>().await.unwrap()
2261 })
2262 },
2263 );
2264 }
2265
2266 #[tokio::test]
2267 async fn test_nested_lists() -> Result<()> {
2268 let list_inner_field = Arc::new(Field::new("item", DataType::Float32, true));
2270 let table_schema = Arc::new(Schema::new(vec![
2271 Field::new("id", DataType::Int32, false),
2272 Field::new("vector", DataType::List(list_inner_field.clone()), true),
2273 ]));
2274
2275 let mut list_builder =
2276 ListBuilder::new(Float32Builder::new()).with_field(list_inner_field.clone());
2277 list_builder.values().append_slice(&[10.0, 10.0, 10.0]);
2278 list_builder.append(true);
2279 list_builder.values().append_slice(&[20.0, 20.0, 20.0]);
2280 list_builder.append(true);
2281 list_builder.values().append_slice(&[30.0, 30.0, 30.0]);
2282 list_builder.append(true);
2283 list_builder.values().append_slice(&[40.0, 40.0, 40.0]);
2284 list_builder.append(true);
2285 let list_array = list_builder.finish();
2286
2287 let data = vec![RecordBatch::try_new(
2288 table_schema.clone(),
2289 vec![
2290 Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
2291 Arc::new(list_array),
2292 ],
2293 )?];
2294
2295 let mut buffer = Vec::new();
2296 let mut writer = AsyncArrowWriter::try_new(&mut buffer, table_schema, None)?;
2297
2298 for batch in data {
2299 writer.write(&batch).await?;
2300 }
2301
2302 writer.close().await?;
2303
2304 let reader = TestReader::new(Bytes::from(buffer));
2305 let builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
2306
2307 let predicate = ArrowPredicateFn::new(ProjectionMask::all(), |batch| {
2308 Ok(BooleanArray::from(vec![true; batch.num_rows()]))
2309 });
2310
2311 let projection_mask = ProjectionMask::all();
2312
2313 let mut stream = builder
2314 .with_row_filter(RowFilter::new(vec![Box::new(predicate)]))
2315 .with_projection(projection_mask)
2316 .build()?;
2317
2318 while let Some(batch) = stream.next().await {
2319 let _ = batch.unwrap(); }
2321
2322 Ok(())
2323 }
2324}