Skip to main content

parquet/arrow/async_reader/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! `async` API for reading Parquet files as [`RecordBatch`]es
19//!
20//! See the [crate-level documentation](crate) for more details.
21//!
22//! See example on [`ParquetRecordBatchStreamBuilder::new`]
23
24use std::fmt::Formatter;
25use std::io::SeekFrom;
26use std::ops::Range;
27use std::pin::Pin;
28use std::sync::Arc;
29use std::task::{Context, Poll};
30
31use bytes::Bytes;
32use futures::future::{BoxFuture, FutureExt};
33use futures::stream::Stream;
34use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeek, AsyncSeekExt};
35
36use arrow_array::RecordBatch;
37use arrow_schema::{Schema, SchemaRef};
38
39use crate::arrow::arrow_reader::{
40    ArrowReaderBuilder, ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReader,
41};
42
43use crate::basic::{BloomFilterAlgorithm, BloomFilterCompression, BloomFilterHash};
44use crate::bloom_filter::{
45    SBBF_HEADER_SIZE_ESTIMATE, Sbbf, chunk_read_bloom_filter_header_and_offset,
46};
47use crate::errors::{ParquetError, Result};
48use crate::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
49
50mod metadata;
51pub use metadata::*;
52
53mod spawn;
54pub use spawn::SpawnedReader;
55
56/// Re-exported so [`ParquetRecordBatchStreamBuilder::with_row_group_selections`]
57/// can be used without importing from another module.
58pub 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
68/// The asynchronous interface used by [`ParquetRecordBatchStream`] to read parquet files
69///
70/// Notes:
71///
72/// 1. There is a default implementation for types that implement [`AsyncRead`]
73///    and [`AsyncSeek`], for example [`tokio::fs::File`].
74///
75/// 2. Implementations for remote storage, such as the `object_store` crate,
76///    can implement this interface directly, typically by pairing a store
77///    handle with an object path and delegating [`Self::get_bytes`] and
78///    [`Self::get_byte_ranges`] to ranged reads. [`SpawnedReader`] can wrap
79///    such a reader to perform its I/O on a dedicated runtime, and
80///    [`ParquetMetaDataReader::with_arrow_reader_options`] simplifies
81///    implementing [`Self::get_metadata`].
82///
83/// # Example: implementing `AsyncFileReader` for the `object_store` crate
84///
85/// ```no_run
86/// # use std::ops::Range;
87/// # use std::sync::Arc;
88/// use bytes::Bytes;
89/// use futures::future::BoxFuture;
90/// use futures::{FutureExt, TryFutureExt};
91/// use object_store::path::Path;
92/// use object_store::{GetOptions, GetRange, ObjectStore, ObjectStoreExt};
93/// use parquet::arrow::arrow_reader::ArrowReaderOptions;
94/// use parquet::arrow::async_reader::{AsyncFileReader, MetadataSuffixFetch};
95/// use parquet::errors::{ParquetError, Result};
96/// use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
97///
98/// fn to_parquet_err(e: object_store::Error) -> ParquetError {
99///     ParquetError::External(Box::new(e))
100/// }
101///
102/// #[derive(Clone)]
103/// struct ObjectStoreReader {
104///     store: Arc<dyn ObjectStore>,
105///     path: Path,
106/// }
107///
108/// impl AsyncFileReader for ObjectStoreReader {
109///     fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
110///         self.store
111///             .get_range(&self.path, range)
112///             .map_err(to_parquet_err)
113///             .boxed()
114///     }
115///
116///     fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_, Result<Vec<Bytes>>> {
117///         async move {
118///             self.store
119///                 .get_ranges(&self.path, &ranges)
120///                 .await
121///                 .map_err(to_parquet_err)
122///         }
123///         .boxed()
124///     }
125///
126///     fn get_metadata<'a>(
127///         &'a mut self,
128///         options: Option<&'a ArrowReaderOptions>,
129///     ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
130///         async move {
131///             let metadata = ParquetMetaDataReader::new()
132///                 .with_arrow_reader_options(options)
133///                 .load_via_suffix_and_finish(self)
134///                 .await?;
135///             Ok(Arc::new(metadata))
136///         }
137///         .boxed()
138///     }
139/// }
140///
141/// /// Supports fetching the parquet footer without knowing the file size,
142/// /// via suffix range requests
143/// impl MetadataSuffixFetch for &mut ObjectStoreReader {
144///     fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
145///         let options = GetOptions {
146///             range: Some(GetRange::Suffix(suffix as u64)),
147///             ..Default::default()
148///         };
149///         async move {
150///             let resp = self
151///                 .store
152///                 .get_opts(&self.path, options)
153///                 .await
154///                 .map_err(to_parquet_err)?;
155///             resp.bytes().await.map_err(to_parquet_err)
156///         }
157///         .boxed()
158///     }
159/// }
160/// ```
161///
162/// [`ParquetMetaDataReader::with_arrow_reader_options`]: crate::file::metadata::ParquetMetaDataReader::with_arrow_reader_options
163///
164/// [`tokio::fs::File`]: https://docs.rs/tokio/latest/tokio/fs/struct.File.html
165pub trait AsyncFileReader: Send {
166    /// Retrieve the bytes in `range`
167    fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>>;
168
169    /// Retrieve multiple byte ranges. The default implementation will call `get_bytes` sequentially
170    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    /// Return a future which results in the [`ParquetMetaData`] for this Parquet file.
185    ///
186    /// This is an asynchronous operation as it may involve reading the file
187    /// footer and potentially other metadata from disk or a remote source.
188    ///
189    /// Reading data from Parquet requires the metadata to understand the
190    /// schema, row groups, and location of pages within the file. This metadata
191    /// is stored primarily in the footer of the Parquet file, and can be read using
192    /// [`ParquetMetaDataReader`].
193    ///
194    /// However, implementations can significantly speed up reading Parquet by
195    /// supplying cached metadata or pre-fetched metadata via this API.
196    ///
197    /// # Parameters
198    /// * `options`: Optional [`ArrowReaderOptions`] that may contain decryption
199    ///   and other options that affect how the metadata is read.
200    fn get_metadata<'a>(
201        &'a mut self,
202        options: Option<&'a ArrowReaderOptions>,
203    ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>>;
204}
205
206/// This allows Box<dyn AsyncFileReader + '_> to be used as an AsyncFileReader,
207impl 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    /// Returns a new [`ArrowReaderMetadata`] for this builder
268    ///
269    /// See [`ParquetRecordBatchStreamBuilder::new_with_metadata`] for how this can be used
270    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)]
280/// Newtype (wrapper) used within [`ArrowReaderBuilder`] to distinguish sync readers from async
281///
282/// Allows sharing the same builder for different readers while keeping the same
283/// ParquetRecordBatchStreamBuilder API
284pub struct AsyncReader<T>(T);
285
286/// A builder for reading parquet files from an `async` source as  [`ParquetRecordBatchStream`]
287///
288/// This can be used to decode a Parquet file in streaming fashion (without
289/// downloading the whole file at once) from a remote source, such as an object store.
290///
291/// This builder handles reading the parquet file metadata, allowing consumers
292/// to use this information to select what specific columns, row groups, etc.
293/// they wish to be read by the resulting stream.
294///
295/// See examples on [`ParquetRecordBatchStreamBuilder::new`], including how to
296/// issue multiple I/O requests in parallel using multiple streams.
297///
298/// # See also:
299/// * [`ParquetPushDecoderBuilder`] for lower level control over buffering and
300///   decoding.
301/// * [`ParquetRecordBatchStream::next_row_group`] for I/O prefetching
302///
303///
304/// See [`ArrowReaderBuilder`] for additional member functions
305pub type ParquetRecordBatchStreamBuilder<T> = ArrowReaderBuilder<AsyncReader<T>>;
306
307impl<T: AsyncFileReader + Send + 'static> ParquetRecordBatchStreamBuilder<T> {
308    /// Create a new [`ParquetRecordBatchStreamBuilder`] for reading from the
309    /// specified source.
310    ///
311    /// # Examples:
312    /// * [Basic example reading from an async source](#example)
313    /// * [Configuring options and reading metadata](#example-configuring-options-and-reading-metadata)
314    /// * [Reading Row Groups in Parallel](#example-reading-row-groups-in-parallel)
315    ///
316    /// # Example
317    /// ```
318    /// # #[tokio::main(flavor="current_thread")]
319    /// # async fn main() {
320    /// #
321    /// # use arrow_array::RecordBatch;
322    /// # use arrow::util::pretty::pretty_format_batches;
323    /// # use futures::TryStreamExt;
324    /// #
325    /// # use parquet::arrow::{ParquetRecordBatchStreamBuilder, ProjectionMask};
326    /// #
327    /// # fn assert_batches_eq(batches: &[RecordBatch], expected_lines: &[&str]) {
328    /// #     let formatted = pretty_format_batches(batches).unwrap().to_string();
329    /// #     let actual_lines: Vec<_> = formatted.trim().lines().collect();
330    /// #     assert_eq!(
331    /// #          &actual_lines, expected_lines,
332    /// #          "\n\nexpected:\n\n{:#?}\nactual:\n\n{:#?}\n\n",
333    /// #          expected_lines, actual_lines
334    /// #      );
335    /// #  }
336    /// #
337    /// # let testdata = arrow::util::test_util::parquet_test_data();
338    /// # let path = format!("{}/alltypes_plain.parquet", testdata);
339    /// // Use tokio::fs::File to read data using an async I/O. This can be replaced with
340    /// // another async I/O reader such as a reader from an object store.
341    /// let file = tokio::fs::File::open(path).await.unwrap();
342    ///
343    /// // Configure options for reading from the async source
344    /// let builder = ParquetRecordBatchStreamBuilder::new(file)
345    ///     .await
346    ///     .unwrap();
347    /// // Building the stream opens the parquet file (reads metadata, etc) and returns
348    /// // a stream that can be used to incrementally read the data in batches
349    /// let stream = builder.build().unwrap();
350    /// // In this example, we collect the stream into a Vec<RecordBatch>
351    /// // but real applications would likely process the batches as they are read
352    /// let results = stream.try_collect::<Vec<_>>().await.unwrap();
353    /// // Demonstrate the results are as expected
354    /// assert_batches_eq(
355    ///     &results,
356    ///     &[
357    ///       "+----+----------+-------------+--------------+---------+------------+-----------+------------+------------------+------------+---------------------+",
358    ///       "| id | bool_col | tinyint_col | smallint_col | int_col | bigint_col | float_col | double_col | date_string_col  | string_col | timestamp_col       |",
359    ///       "+----+----------+-------------+--------------+---------+------------+-----------+------------+------------------+------------+---------------------+",
360    ///       "| 4  | true     | 0           | 0            | 0       | 0          | 0.0       | 0.0        | 30332f30312f3039 | 30         | 2009-03-01T00:00:00 |",
361    ///       "| 5  | false    | 1           | 1            | 1       | 10         | 1.1       | 10.1       | 30332f30312f3039 | 31         | 2009-03-01T00:01:00 |",
362    ///       "| 6  | true     | 0           | 0            | 0       | 0          | 0.0       | 0.0        | 30342f30312f3039 | 30         | 2009-04-01T00:00:00 |",
363    ///       "| 7  | false    | 1           | 1            | 1       | 10         | 1.1       | 10.1       | 30342f30312f3039 | 31         | 2009-04-01T00:01:00 |",
364    ///       "| 2  | true     | 0           | 0            | 0       | 0          | 0.0       | 0.0        | 30322f30312f3039 | 30         | 2009-02-01T00:00:00 |",
365    ///       "| 3  | false    | 1           | 1            | 1       | 10         | 1.1       | 10.1       | 30322f30312f3039 | 31         | 2009-02-01T00:01:00 |",
366    ///       "| 0  | true     | 0           | 0            | 0       | 0          | 0.0       | 0.0        | 30312f30312f3039 | 30         | 2009-01-01T00:00:00 |",
367    ///       "| 1  | false    | 1           | 1            | 1       | 10         | 1.1       | 10.1       | 30312f30312f3039 | 31         | 2009-01-01T00:01:00 |",
368    ///       "+----+----------+-------------+--------------+---------+------------+-----------+------------+------------------+------------+---------------------+",
369    ///      ],
370    ///  );
371    /// # }
372    /// ```
373    ///
374    /// # Example Configuring Options and Reading Metadata
375    ///
376    /// There are many options that control the behavior of the reader, such as
377    /// `with_batch_size`, `with_projection`, `with_filter`, etc...
378    ///
379    /// ```
380    /// # #[tokio::main(flavor="current_thread")]
381    /// # async fn main() {
382    /// #
383    /// # use arrow_array::RecordBatch;
384    /// # use arrow::util::pretty::pretty_format_batches;
385    /// # use futures::TryStreamExt;
386    /// #
387    /// # use parquet::arrow::{ParquetRecordBatchStreamBuilder, ProjectionMask};
388    /// #
389    /// # fn assert_batches_eq(batches: &[RecordBatch], expected_lines: &[&str]) {
390    /// #     let formatted = pretty_format_batches(batches).unwrap().to_string();
391    /// #     let actual_lines: Vec<_> = formatted.trim().lines().collect();
392    /// #     assert_eq!(
393    /// #          &actual_lines, expected_lines,
394    /// #          "\n\nexpected:\n\n{:#?}\nactual:\n\n{:#?}\n\n",
395    /// #          expected_lines, actual_lines
396    /// #      );
397    /// #  }
398    /// #
399    /// # let testdata = arrow::util::test_util::parquet_test_data();
400    /// # let path = format!("{}/alltypes_plain.parquet", testdata);
401    /// // As before, use tokio::fs::File to read data using an async I/O.
402    /// let file = tokio::fs::File::open(path).await.unwrap();
403    ///
404    /// // Configure options for reading from the async source, in this case we set the batch size
405    /// // to 3 which produces 3 rows at a time.
406    /// let builder = ParquetRecordBatchStreamBuilder::new(file)
407    ///     .await
408    ///     .unwrap()
409    ///     .with_batch_size(3);
410    ///
411    /// // We can also read the metadata to inspect the schema and other metadata
412    /// // before actually reading the data
413    /// let file_metadata = builder.metadata().file_metadata();
414    /// // Specify that we only want to read the 1st, 2nd, and 6th columns
415    /// let mask = ProjectionMask::roots(file_metadata.schema_descr(), [1, 2, 6]);
416    ///
417    /// let stream = builder.with_projection(mask).build().unwrap();
418    /// let results = stream.try_collect::<Vec<_>>().await.unwrap();
419    /// // Print out the results
420    /// assert_batches_eq(
421    ///     &results,
422    ///     &[
423    ///         "+----------+-------------+-----------+",
424    ///         "| bool_col | tinyint_col | float_col |",
425    ///         "+----------+-------------+-----------+",
426    ///         "| true     | 0           | 0.0       |",
427    ///         "| false    | 1           | 1.1       |",
428    ///         "| true     | 0           | 0.0       |",
429    ///         "| false    | 1           | 1.1       |",
430    ///         "| true     | 0           | 0.0       |",
431    ///         "| false    | 1           | 1.1       |",
432    ///         "| true     | 0           | 0.0       |",
433    ///         "| false    | 1           | 1.1       |",
434    ///         "+----------+-------------+-----------+",
435    ///      ],
436    ///  );
437    ///
438    /// // The results has 8 rows, so since we set the batch size to 3, we expect
439    /// // 3 batches, two with 3 rows each and the last batch with 2 rows.
440    /// assert_eq!(results.len(), 3);
441    /// # }
442    /// ```
443    ///
444    /// # Example reading Row Groups in Parallel
445    ///
446    /// Each [`ParquetRecordBatchStream`] is independent and can be used to read
447    /// from the same underlying source in parallel. Use
448    /// [`ParquetRecordBatchStream::next_row_group`] with a single stream to
449    /// begin prefetching the next Row Group. To read a file in parallel, create
450    /// a stream for each subset of the file. For example, you can read each
451    /// row group in parallel by creating a stream for each row group using the
452    /// [`ParquetRecordBatchStreamBuilder::with_row_groups`] API as shown below
453    ///
454    /// ```
455    /// # use std::sync::Arc;
456    /// # use arrow_array::{ArrayRef, Int32Array, RecordBatch};
457    /// # use arrow::util::pretty::pretty_format_batches;
458    /// # use futures::{StreamExt, TryStreamExt};
459    /// # use tempfile::NamedTempFile;
460    /// # use parquet::arrow::{ArrowWriter, ParquetRecordBatchStreamBuilder, ProjectionMask};
461    /// # use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
462    /// # use parquet::file::metadata::ParquetMetaDataReader;
463    /// # use parquet::file::properties::{WriterProperties};
464    /// # // write to a temporary file with 10 RowGroups and read back with async API
465    /// # fn write_file() -> parquet::errors::Result<NamedTempFile> {
466    /// #   let mut file = NamedTempFile::new().unwrap();
467    /// #   let small_batch = RecordBatch::try_from_iter([
468    /// #      ("id", Arc::new(Int32Array::from(vec![0, 1, 2, 3, 4])) as ArrayRef),
469    /// #   ]).unwrap();
470    /// #   let props = WriterProperties::builder()
471    /// #     .set_max_row_group_row_count(Some(5))
472    /// #     .set_write_batch_size(5)
473    /// #     .build();
474    /// #   let mut writer = ArrowWriter::try_new(&mut file, small_batch.schema(), Some(props))?;
475    /// #   for i in 0..10 {
476    /// #     writer.write(&small_batch)?
477    /// #   };
478    /// #   writer.close()?;
479    /// #   Ok(file)
480    /// # }
481    /// # #[tokio::main(flavor="current_thread")]
482    /// # async fn main() -> parquet::errors::Result<()> {
483    /// # let t = write_file()?;
484    /// # let path = t.path();
485    /// // This example uses a tokio::fs::File as the async source, but it
486    /// // could be any async source such as an object store reader)
487    /// let mut file = tokio::fs::File::open(path).await?;
488    /// // To read Row Groups in parallel, create a separate stream builder for each Row Group.
489    /// // First get the metadata to find the row group information
490    /// let file_size = file.metadata().await?.len();
491    /// let metadata = ParquetMetaDataReader::new().load_and_finish(&mut file, file_size).await?;
492    /// assert_eq!(metadata.num_row_groups(), 10); // file has 10 row groups with 5 rows each
493    /// // Create a stream reader for each row group
494    /// let reader_metadata = ArrowReaderMetadata::try_new(
495    ///   Arc::new(metadata),
496    ///   ArrowReaderOptions::new()
497    /// )?;
498    /// let mut streams = vec![];
499    ///  for row_group_index in 0..10 {
500    ///   // Each stream needs its own source instance to issue
501    ///   // parallel IO requests, so clone the file for each stream
502    ///   let this_file = file.try_clone().await?;
503    ///   let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
504    ///        this_file,
505    ///        reader_metadata.clone()
506    ///      )
507    ///      .with_row_groups(vec![row_group_index]) // read only this row group
508    ///      .build()?;
509    ///     streams.push(stream);
510    /// }
511    /// // Each reader can now be polled independently and in parallel, for
512    /// // example using StreamExt::buffered to read from 3 at a time
513    /// let results = futures::stream::iter(streams)
514    ///  .map(|stream| async move { stream })
515    ///  .buffered(3)
516    ///  .flatten()
517    ///  .try_collect::<Vec<_>>().await?;
518    /// // read all 50 rows (10 row groups x 5 rows per group)
519    /// assert_eq!(50, results.iter().map(|s| s.num_rows()).sum::<usize>());
520    /// # Ok(())
521    /// # }
522    /// ```
523    pub async fn new(input: T) -> Result<Self> {
524        Self::new_with_options(input, Default::default()).await
525    }
526
527    /// Create a new [`ParquetRecordBatchStreamBuilder`] with the provided async source
528    /// and [`ArrowReaderOptions`].
529    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    /// Create a [`ParquetRecordBatchStreamBuilder`] from the provided [`ArrowReaderMetadata`]
535    ///
536    /// This allows loading metadata once and using it to create multiple builders with
537    /// potentially different settings, that can be read in parallel.
538    ///
539    /// # Example of reading from multiple streams in parallel
540    ///
541    /// ```
542    /// # use std::fs::metadata;
543    /// # use std::sync::Arc;
544    /// # use bytes::Bytes;
545    /// # use arrow_array::{Int32Array, RecordBatch};
546    /// # use arrow_schema::{DataType, Field, Schema};
547    /// # use parquet::arrow::arrow_reader::ArrowReaderMetadata;
548    /// # use parquet::arrow::{ArrowWriter, ParquetRecordBatchStreamBuilder};
549    /// # use tempfile::tempfile;
550    /// # use futures::StreamExt;
551    /// # #[tokio::main(flavor="current_thread")]
552    /// # async fn main() {
553    /// #
554    /// # let mut file = tempfile().unwrap();
555    /// # let schema = Arc::new(Schema::new(vec![Field::new("i32", DataType::Int32, false)]));
556    /// # let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
557    /// # let batch = RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1, 2, 3]))]).unwrap();
558    /// # writer.write(&batch).unwrap();
559    /// # writer.close().unwrap();
560    /// // open file with parquet data
561    /// let mut file = tokio::fs::File::from_std(file);
562    /// // load metadata once
563    /// let meta = ArrowReaderMetadata::load_async(&mut file, Default::default()).await.unwrap();
564    /// // create two readers, a and b, from the same underlying file
565    /// // without reading the metadata again
566    /// let mut a = ParquetRecordBatchStreamBuilder::new_with_metadata(
567    ///     file.try_clone().await.unwrap(),
568    ///     meta.clone()
569    /// ).build().unwrap();
570    /// let mut b = ParquetRecordBatchStreamBuilder::new_with_metadata(file, meta).build().unwrap();
571    ///
572    /// // Can read batches from both readers in parallel
573    /// assert_eq!(
574    ///   a.next().await.unwrap().unwrap(),
575    ///   b.next().await.unwrap().unwrap(),
576    /// );
577    /// # }
578    /// ```
579    pub fn new_with_metadata(input: T, metadata: ArrowReaderMetadata) -> Self {
580        Self::new_builder(AsyncReader(input), metadata)
581    }
582
583    /// Read bloom filter for a column in a row group
584    ///
585    /// Returns `None` if the column does not have a bloom filter
586    ///
587    /// We should call this function after other forms pruning, such as projection and predicate pushdown.
588    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                // this match exists to future proof the singleton algorithm enum
619            }
620        }
621        match header.compression {
622            BloomFilterCompression::UNCOMPRESSED => {
623                // this match exists to future proof the singleton compression enum
624            }
625        }
626        match header.hash {
627            BloomFilterHash::XXHASH => {
628                // this match exists to future proof the singleton hash enum
629            }
630        }
631
632        let bitset = match column_metadata.bloom_filter_length() {
633            Some(_) => buffer.slice(
634                (TryInto::<usize>::try_into(bitset_offset).unwrap()
635                    - TryInto::<usize>::try_into(offset).unwrap())..,
636            ),
637            None => {
638                let bitset_length: u64 = header.num_bytes.try_into().map_err(|_| {
639                    ParquetError::General("Bloom filter length is invalid".to_string())
640                })?;
641                self.input
642                    .0
643                    .get_bytes(bitset_offset..bitset_offset + bitset_length)
644                    .await?
645            }
646        };
647        Ok(Some(Sbbf::new(&bitset)))
648    }
649
650    /// Select row groups and rows using row-group-local coordinates.
651    ///
652    /// Entries are decoded in the supplied order, omitted row groups are
653    /// skipped, and a `None` selection reads the whole row group. This is
654    /// mutually exclusive with [`ArrowReaderBuilder::with_row_groups`] and
655    /// [`ArrowReaderBuilder::with_row_selection`]; combining them returns an
656    /// error from [`Self::build`].
657    ///
658    /// See [`ParquetPushDecoderBuilder::with_row_group_selections`] for the
659    /// full semantics and a worked example. This builder supports the same API
660    /// because the async stream is implemented using the push decoder; the
661    /// synchronous reader does not support row-group-local selections.
662    ///
663    /// [`ParquetPushDecoderBuilder::with_row_group_selections`]: crate::arrow::push_decoder::ParquetPushDecoderBuilder::with_row_group_selections
664    pub fn with_row_group_selections(
665        mut self,
666        row_group_selections: Vec<RowGroupSelection>,
667    ) -> Self {
668        self.row_group_plan
669            .set_row_group_selections(row_group_selections);
670        self
671    }
672
673    /// Build a new [`ParquetRecordBatchStream`]
674    ///
675    /// See examples on [`ParquetRecordBatchStreamBuilder::new`]
676    pub fn build(self) -> Result<ParquetRecordBatchStream<T>> {
677        let Self {
678            input,
679            metadata,
680            schema,
681            fields,
682            batch_size,
683            row_group_plan,
684            projection,
685            filter,
686            row_selection_policy: selection_strategy,
687            limit,
688            offset,
689            metrics,
690            max_predicate_cache_size,
691        } = self;
692
693        // Ensure schema of ParquetRecordBatchStream respects projection, and does
694        // not store metadata (same as for ParquetRecordBatchReader and emitted RecordBatches)
695        let projection_len = projection.mask.as_ref().map_or(usize::MAX, |m| m.len());
696        let projected_fields = schema
697            .fields
698            .filter_leaves(|idx, _| idx < projection_len && projection.leaf_included(idx));
699        let projected_schema = Arc::new(Schema::new(projected_fields));
700
701        let decoder = ParquetPushDecoderBuilder {
702            input: PushDecoderInput::default(),
703            metadata,
704            schema,
705            fields,
706            projection,
707            filter,
708            row_group_plan,
709            row_selection_policy: selection_strategy,
710            batch_size,
711            limit,
712            offset,
713            metrics,
714            max_predicate_cache_size,
715        }
716        .build()?;
717
718        let request_state = RequestState::None { input: input.0 };
719
720        Ok(ParquetRecordBatchStream {
721            schema: projected_schema,
722            decoder,
723            request_state,
724        })
725    }
726}
727
728/// State machine that tracks outstanding requests to fetch data
729///
730/// The parameter `T` is the input, typically an `AsyncFileReader`
731enum RequestState<T> {
732    /// No outstanding requests
733    None {
734        input: T,
735    },
736    /// There is an outstanding request for data
737    Outstanding {
738        /// Ranges that have been requested
739        ranges: Vec<Range<u64>>,
740        /// Future that will resolve (input, requested_ranges)
741        ///
742        /// Note the future owns the reader while the request is outstanding
743        /// and returns it upon completion
744        future: BoxFuture<'static, Result<(T, Vec<Bytes>)>>,
745    },
746    Done,
747}
748
749impl<T> RequestState<T>
750where
751    T: AsyncFileReader + Unpin + Send + 'static,
752{
753    /// Issue a request to fetch `ranges`, returning the Outstanding state
754    fn begin_request(mut input: T, ranges: Vec<Range<u64>>) -> Self {
755        let ranges_captured = ranges.clone();
756
757        // Note this must move the input *into* the future
758        // because the get_byte_ranges future has a lifetime
759        // (aka can have references internally) and thus must
760        // own the input while the request is outstanding.
761        let future = async move {
762            let data = input.get_byte_ranges(ranges_captured).await?;
763            Ok((input, data))
764        }
765        .boxed();
766        RequestState::Outstanding { ranges, future }
767    }
768}
769
770impl<T> std::fmt::Debug for RequestState<T> {
771    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
772        match self {
773            RequestState::None { input: _ } => f
774                .debug_struct("RequestState::None")
775                .field("input", &"...")
776                .finish(),
777            RequestState::Outstanding { ranges, .. } => f
778                .debug_struct("RequestState::Outstanding")
779                .field("ranges", &ranges)
780                .finish(),
781            RequestState::Done => {
782                write!(f, "RequestState::Done")
783            }
784        }
785    }
786}
787
788/// An asynchronous [`Stream`]of [`RecordBatch`] constructed using [`ParquetRecordBatchStreamBuilder`] to read parquet files.
789///
790/// `ParquetRecordBatchStream` also provides [`ParquetRecordBatchStream::next_row_group`] for fetching row groups,
791/// allowing users to decode record batches separately from I/O.
792///
793/// # I/O Buffering
794///
795/// `ParquetRecordBatchStream` buffers *all* data pages selected after predicates
796/// (projection + filtering, etc) and decodes the rows from those buffered pages.
797///
798/// For example, if all rows and columns are selected, the entire row group is
799/// buffered in memory during decode. This minimizes the number of IO operations
800/// required, which is especially important for object stores, where IO operations
801/// have latencies in the hundreds of milliseconds
802///
803/// See [`ParquetPushDecoderBuilder`] for an API with lower level control over
804/// buffering.
805///
806/// [`Stream`]: https://docs.rs/futures/latest/futures/stream/trait.Stream.html
807pub struct ParquetRecordBatchStream<T> {
808    /// Output schema of the stream
809    schema: SchemaRef,
810    /// Input and Outstanding IO request, if any
811    request_state: RequestState<T>,
812    /// Decoding state machine (no IO)
813    decoder: ParquetPushDecoder,
814}
815
816impl<T> std::fmt::Debug for ParquetRecordBatchStream<T> {
817    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
818        f.debug_struct("ParquetRecordBatchStream")
819            .field("request_state", &self.request_state)
820            .finish()
821    }
822}
823
824impl<T> ParquetRecordBatchStream<T> {
825    /// Returns the projected [`SchemaRef`] for reading the parquet file.
826    ///
827    /// Note that the schema metadata will be stripped here. See
828    /// [`ParquetRecordBatchStreamBuilder::schema`] if the metadata is desired.
829    pub fn schema(&self) -> &SchemaRef {
830        &self.schema
831    }
832}
833
834impl<T> ParquetRecordBatchStream<T>
835where
836    T: AsyncFileReader + Unpin + Send + 'static,
837{
838    /// Fetches the next row group from the stream.
839    ///
840    /// Users can continue to call this function to get row groups and decode them concurrently.
841    ///
842    /// ## Notes
843    ///
844    /// ParquetRecordBatchStream should be used either as a `Stream` or with `next_row_group`; they should not be used simultaneously.
845    ///
846    /// ## Returns
847    ///
848    /// - `Ok(None)` if the stream has ended.
849    /// - `Err(error)` if the stream has errored. All subsequent calls will return `Ok(None)`.
850    /// - `Ok(Some(reader))` which holds all the data for the row group.
851    pub async fn next_row_group(&mut self) -> Result<Option<ParquetRecordBatchReader>> {
852        loop {
853            // Take ownership of request state to process, leaving self in a
854            // valid state
855            let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
856            match request_state {
857                // No outstanding requests, proceed to setup next row group
858                RequestState::None { input } => {
859                    match self.decoder.try_next_reader()? {
860                        DecodeResult::NeedsData(ranges) => {
861                            self.request_state = RequestState::begin_request(input, ranges);
862                            // Will loop again: the input might be ready immediately.
863                        }
864                        DecodeResult::Data(reader) => {
865                            self.request_state = RequestState::None { input };
866                            return Ok(Some(reader));
867                        }
868                        DecodeResult::Finished => return Ok(None),
869                    }
870                }
871                RequestState::Outstanding { ranges, future } => {
872                    let (input, data) = future.await?;
873                    // Push the requested data to the decoder and try again
874                    self.decoder.push_ranges(ranges, data)?;
875                    self.request_state = RequestState::None { input };
876                    // Will try and decode on the next iteration.
877                }
878                RequestState::Done => {
879                    self.request_state = RequestState::Done;
880                    return Ok(None);
881                }
882            }
883        }
884    }
885}
886
887impl<T> Stream for ParquetRecordBatchStream<T>
888where
889    T: AsyncFileReader + Unpin + Send + 'static,
890{
891    type Item = Result<RecordBatch>;
892    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
893        match self.poll_next_inner(cx) {
894            Ok(res) => {
895                // Successfully decoded a batch, or reached end of stream.
896                // convert Option<RecordBatch> to Option<Result<RecordBatch>>
897                res.map(|res| Ok(res).transpose())
898            }
899            Err(e) => {
900                self.request_state = RequestState::Done;
901                Poll::Ready(Some(Err(e)))
902            }
903        }
904    }
905}
906
907impl<T> ParquetRecordBatchStream<T>
908where
909    T: AsyncFileReader + Unpin + Send + 'static,
910{
911    /// Inner state machine
912    ///
913    /// Note this is separate from poll_next so we can use ? operator to check for errors
914    /// as it returns `Result<Poll<Option<RecordBatch>>>`
915    fn poll_next_inner(&mut self, cx: &mut Context<'_>) -> Result<Poll<Option<RecordBatch>>> {
916        loop {
917            let request_state = std::mem::replace(&mut self.request_state, RequestState::Done);
918            match request_state {
919                RequestState::None { input } => {
920                    // No outstanding requests, proceed to decode the next batch
921                    match self.decoder.try_decode()? {
922                        DecodeResult::NeedsData(ranges) => {
923                            self.request_state = RequestState::begin_request(input, ranges);
924                            // Will loop again: the input might be ready immediately.
925                        }
926                        DecodeResult::Data(batch) => {
927                            self.request_state = RequestState::None { input };
928                            return Ok(Poll::Ready(Some(batch)));
929                        }
930                        DecodeResult::Finished => {
931                            self.request_state = RequestState::Done;
932                            return Ok(Poll::Ready(None));
933                        }
934                    }
935                }
936                RequestState::Outstanding { ranges, mut future } => match future.poll_unpin(cx) {
937                    // Data was ready, push it to the decoder and continue
938                    Poll::Ready(result) => {
939                        let (input, data) = result?;
940                        // Push the requested data to the decoder
941                        self.decoder.push_ranges(ranges, data)?;
942                        self.request_state = RequestState::None { input };
943                        // The next iteration will try to decode the next batch.
944                    }
945                    Poll::Pending => {
946                        self.request_state = RequestState::Outstanding { ranges, future };
947                        return Ok(Poll::Pending);
948                    }
949                },
950                RequestState::Done => {
951                    // Stream is done (error or end), return None
952                    self.request_state = RequestState::Done;
953                    return Ok(Poll::Ready(None));
954                }
955            }
956        }
957    }
958}
959
960#[cfg(test)]
961mod tests {
962    use super::*;
963    use crate::arrow::arrow_reader::tests::test_row_numbers_with_multiple_row_groups_helper;
964    use crate::arrow::arrow_reader::{
965        ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowFilter, RowSelection, RowSelector,
966    };
967    use crate::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
968    use crate::arrow::schema::virtual_type::RowNumber;
969    use crate::arrow::{ArrowWriter, AsyncArrowWriter, ProjectionMask};
970    use crate::file::metadata::ParquetMetaDataReader;
971    use crate::file::metadata::{PageIndex, PageIndexPolicy};
972    use crate::file::properties::WriterProperties;
973    use arrow::compute::kernels::cmp::eq;
974    use arrow::error::Result as ArrowResult;
975    use arrow_array::builder::{Float32Builder, ListBuilder, StringBuilder};
976    use arrow_array::cast::AsArray;
977    use arrow_array::types::Int32Type;
978    use arrow_array::{
979        Array, ArrayRef, BooleanArray, Int32Array, RecordBatchReader, Scalar, StringArray,
980        StructArray, UInt64Array,
981    };
982    use arrow_schema::{DataType, Field, Schema};
983    use futures::{StreamExt, TryStreamExt};
984    use rand::{RngExt, rng};
985    use std::collections::HashMap;
986    use std::sync::{Arc, Mutex};
987    use tempfile::tempfile;
988
989    #[derive(Clone)]
990    struct TestReader {
991        data: Bytes,
992        metadata: Option<Arc<ParquetMetaData>>,
993        requests: Arc<Mutex<Vec<Range<usize>>>>,
994    }
995
996    impl TestReader {
997        fn new(data: Bytes) -> Self {
998            Self {
999                data,
1000                metadata: Default::default(),
1001                requests: Default::default(),
1002            }
1003        }
1004    }
1005
1006    impl AsyncFileReader for TestReader {
1007        fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1008            let range = range.clone();
1009            self.requests
1010                .lock()
1011                .unwrap()
1012                .push(range.start as usize..range.end as usize);
1013            futures::future::ready(Ok(self
1014                .data
1015                .slice(range.start as usize..range.end as usize)))
1016            .boxed()
1017        }
1018
1019        fn get_metadata<'a>(
1020            &'a mut self,
1021            options: Option<&'a ArrowReaderOptions>,
1022        ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
1023            let metadata_reader = ParquetMetaDataReader::new().with_arrow_reader_options(options);
1024            self.metadata = Some(Arc::new(
1025                metadata_reader.parse_and_finish(&self.data).unwrap(),
1026            ));
1027            futures::future::ready(Ok(self.metadata.clone().unwrap().clone())).boxed()
1028        }
1029    }
1030
1031    #[tokio::test]
1032    async fn test_async_reader() {
1033        let testdata = arrow::util::test_util::parquet_test_data();
1034        let path = format!("{testdata}/alltypes_plain.parquet");
1035        let data = Bytes::from(std::fs::read(path).unwrap());
1036
1037        let async_reader = TestReader::new(data.clone());
1038
1039        let requests = async_reader.requests.clone();
1040        let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1041            .await
1042            .unwrap();
1043
1044        let metadata = builder.metadata().clone();
1045        assert_eq!(metadata.num_row_groups(), 1);
1046
1047        let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1048        let stream = builder
1049            .with_projection(mask.clone())
1050            .with_batch_size(1024)
1051            .build()
1052            .unwrap();
1053
1054        let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1055
1056        let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1057            .unwrap()
1058            .with_projection(mask)
1059            .with_batch_size(104)
1060            .build()
1061            .unwrap()
1062            .collect::<ArrowResult<Vec<_>>>()
1063            .unwrap();
1064
1065        assert_eq!(async_batches, sync_batches);
1066
1067        let requests = requests.lock().unwrap();
1068        let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1069        let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1070
1071        assert_eq!(
1072            &requests[..],
1073            &[
1074                offset_1 as usize..(offset_1 + length_1) as usize,
1075                offset_2 as usize..(offset_2 + length_2) as usize
1076            ]
1077        );
1078    }
1079
1080    #[tokio::test]
1081    async fn test_async_reader_row_group_local_selections() {
1082        let batch = RecordBatch::try_from_iter([(
1083            "a",
1084            Arc::new(Int32Array::from_iter_values(0..6)) as ArrayRef,
1085        )])
1086        .unwrap();
1087        let mut data = Vec::new();
1088        let properties = WriterProperties::builder()
1089            .set_max_row_group_row_count(Some(3))
1090            .build();
1091        let mut writer = ArrowWriter::try_new(&mut data, batch.schema(), Some(properties)).unwrap();
1092        writer.write(&batch).unwrap();
1093        writer.close().unwrap();
1094
1095        let stream = ParquetRecordBatchStreamBuilder::new(TestReader::new(data.into()))
1096            .await
1097            .unwrap()
1098            .with_row_group_selections(vec![
1099                RowGroupSelection::new(1, Some(RowSelection::from(vec![RowSelector::select(1)]))),
1100                RowGroupSelection::new(
1101                    0,
1102                    Some(RowSelection::from(vec![
1103                        RowSelector::skip(1),
1104                        RowSelector::select(2),
1105                    ])),
1106                ),
1107            ])
1108            .build()
1109            .unwrap();
1110
1111        let batches: Vec<_> = stream.try_collect().await.unwrap();
1112        assert_eq!(batches.len(), 2);
1113        assert_eq!(
1114            batches[0].column(0).as_primitive::<Int32Type>().values(),
1115            &[3]
1116        );
1117        assert_eq!(
1118            batches[1].column(0).as_primitive::<Int32Type>().values(),
1119            &[1, 2]
1120        );
1121    }
1122
1123    #[tokio::test]
1124    async fn test_async_reader_with_next_row_group() {
1125        let testdata = arrow::util::test_util::parquet_test_data();
1126        let path = format!("{testdata}/alltypes_plain.parquet");
1127        let data = Bytes::from(std::fs::read(path).unwrap());
1128
1129        let async_reader = TestReader::new(data.clone());
1130
1131        let requests = async_reader.requests.clone();
1132        let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1133            .await
1134            .unwrap();
1135
1136        let metadata = builder.metadata().clone();
1137        assert_eq!(metadata.num_row_groups(), 1);
1138
1139        let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1140        let mut stream = builder
1141            .with_projection(mask.clone())
1142            .with_batch_size(1024)
1143            .build()
1144            .unwrap();
1145
1146        let mut readers = vec![];
1147        while let Some(reader) = stream.next_row_group().await.unwrap() {
1148            readers.push(reader);
1149        }
1150
1151        let async_batches: Vec<_> = readers
1152            .into_iter()
1153            .flat_map(|r| r.map(|v| v.unwrap()).collect::<Vec<_>>())
1154            .collect();
1155
1156        let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1157            .unwrap()
1158            .with_projection(mask)
1159            .with_batch_size(104)
1160            .build()
1161            .unwrap()
1162            .collect::<ArrowResult<Vec<_>>>()
1163            .unwrap();
1164
1165        assert_eq!(async_batches, sync_batches);
1166
1167        let requests = requests.lock().unwrap();
1168        let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range();
1169        let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range();
1170
1171        assert_eq!(
1172            &requests[..],
1173            &[
1174                offset_1 as usize..(offset_1 + length_1) as usize,
1175                offset_2 as usize..(offset_2 + length_2) as usize
1176            ]
1177        );
1178    }
1179
1180    #[tokio::test]
1181    async fn test_async_reader_with_index() {
1182        let testdata = arrow::util::test_util::parquet_test_data();
1183        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1184        let data = Bytes::from(std::fs::read(path).unwrap());
1185
1186        let async_reader = TestReader::new(data.clone());
1187
1188        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1189        let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1190            .await
1191            .unwrap();
1192
1193        // The builder should have page and offset indexes loaded now
1194        let metadata_with_index = builder.metadata();
1195        assert_eq!(metadata_with_index.num_row_groups(), 1);
1196
1197        // Check offset indexes are present for all columns of all row groups
1198        let page_index = metadata_with_index.page_index().unwrap();
1199        let num_rowgroups = metadata_with_index.num_row_groups();
1200        let num_columns = metadata_with_index
1201            .file_metadata()
1202            .schema_descr()
1203            .num_columns();
1204        for rgidx in 0..num_rowgroups {
1205            let column_index = page_index.column_indexes_for_rowgroup(rgidx);
1206            let offset_index = page_index.offset_indexes_for_rowgroup(rgidx);
1207            assert!(column_index.is_some_and(|ci| ci.len() == num_columns));
1208            assert!(offset_index.is_some_and(|oi| oi.len() == num_columns));
1209            // some column indexes are not defined, but all offset indexes should be
1210            for colidx in 0..num_columns {
1211                assert!(page_index.offset_index(rgidx, colidx).is_some());
1212            }
1213        }
1214
1215        let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1216        let stream = builder
1217            .with_projection(mask.clone())
1218            .with_batch_size(1024)
1219            .build()
1220            .unwrap();
1221
1222        let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1223
1224        let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1225            .unwrap()
1226            .with_projection(mask)
1227            .with_batch_size(1024)
1228            .build()
1229            .unwrap()
1230            .collect::<ArrowResult<Vec<_>>>()
1231            .unwrap();
1232
1233        assert_eq!(async_batches, sync_batches);
1234    }
1235
1236    #[tokio::test]
1237    async fn test_async_reader_with_limit() {
1238        let testdata = arrow::util::test_util::parquet_test_data();
1239        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1240        let data = Bytes::from(std::fs::read(path).unwrap());
1241
1242        let metadata = ParquetMetaDataReader::new()
1243            .parse_and_finish(&data)
1244            .unwrap();
1245        let metadata = Arc::new(metadata);
1246
1247        assert_eq!(metadata.num_row_groups(), 1);
1248
1249        let async_reader = TestReader::new(data.clone());
1250
1251        let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1252            .await
1253            .unwrap();
1254
1255        assert_eq!(builder.metadata().num_row_groups(), 1);
1256
1257        let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]);
1258        let stream = builder
1259            .with_projection(mask.clone())
1260            .with_batch_size(1024)
1261            .with_limit(1)
1262            .build()
1263            .unwrap();
1264
1265        let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1266
1267        let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1268            .unwrap()
1269            .with_projection(mask)
1270            .with_batch_size(1024)
1271            .with_limit(1)
1272            .build()
1273            .unwrap()
1274            .collect::<ArrowResult<Vec<_>>>()
1275            .unwrap();
1276
1277        assert_eq!(async_batches, sync_batches);
1278    }
1279
1280    #[tokio::test]
1281    async fn test_async_reader_skip_pages() {
1282        let testdata = arrow::util::test_util::parquet_test_data();
1283        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1284        let data = Bytes::from(std::fs::read(path).unwrap());
1285
1286        let async_reader = TestReader::new(data.clone());
1287
1288        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1289        let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1290            .await
1291            .unwrap();
1292
1293        assert_eq!(builder.metadata().num_row_groups(), 1);
1294
1295        let selection = RowSelection::from(vec![
1296            RowSelector::skip(21),   // Skip first page
1297            RowSelector::select(21), // Select page to boundary
1298            RowSelector::skip(41),   // Skip multiple pages
1299            RowSelector::select(41), // Select multiple pages
1300            RowSelector::skip(25),   // Skip page across boundary
1301            RowSelector::select(25), // Select across page boundary
1302            RowSelector::skip(7116), // Skip to final page boundary
1303            RowSelector::select(10), // Select final page
1304        ]);
1305
1306        let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![9]);
1307
1308        let stream = builder
1309            .with_projection(mask.clone())
1310            .with_row_selection(selection.clone())
1311            .build()
1312            .expect("building stream");
1313
1314        let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1315
1316        let sync_batches = ParquetRecordBatchReaderBuilder::try_new(data)
1317            .unwrap()
1318            .with_projection(mask)
1319            .with_batch_size(1024)
1320            .with_row_selection(selection)
1321            .build()
1322            .unwrap()
1323            .collect::<ArrowResult<Vec<_>>>()
1324            .unwrap();
1325
1326        assert_eq!(async_batches, sync_batches);
1327    }
1328
1329    #[tokio::test]
1330    async fn test_fuzz_async_reader_selection() {
1331        let testdata = arrow::util::test_util::parquet_test_data();
1332        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1333        let data = Bytes::from(std::fs::read(path).unwrap());
1334
1335        let mut rand = rng();
1336
1337        for _ in 0..100 {
1338            let mut expected_rows = 0;
1339            let mut total_rows = 0;
1340            let mut skip = false;
1341            let mut selectors = vec![];
1342
1343            while total_rows < 7300 {
1344                let row_count: usize = rand.random_range(1..100);
1345
1346                let row_count = row_count.min(7300 - total_rows);
1347
1348                selectors.push(RowSelector { row_count, skip });
1349
1350                total_rows += row_count;
1351                if !skip {
1352                    expected_rows += row_count;
1353                }
1354
1355                skip = !skip;
1356            }
1357
1358            let selection = RowSelection::from(selectors);
1359
1360            let async_reader = TestReader::new(data.clone());
1361
1362            let options =
1363                ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1364            let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1365                .await
1366                .unwrap();
1367
1368            assert_eq!(builder.metadata().num_row_groups(), 1);
1369
1370            let col_idx: usize = rand.random_range(0..13);
1371            let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1372
1373            let stream = builder
1374                .with_projection(mask.clone())
1375                .with_row_selection(selection.clone())
1376                .build()
1377                .expect("building stream");
1378
1379            let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1380
1381            let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1382
1383            assert_eq!(actual_rows, expected_rows);
1384        }
1385    }
1386
1387    #[tokio::test]
1388    async fn test_async_reader_zero_row_selector() {
1389        //See https://github.com/apache/arrow-rs/issues/2669
1390        let testdata = arrow::util::test_util::parquet_test_data();
1391        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1392        let data = Bytes::from(std::fs::read(path).unwrap());
1393
1394        let mut rand = rng();
1395
1396        let mut expected_rows = 0;
1397        let mut total_rows = 0;
1398        let mut skip = false;
1399        let mut selectors = vec![];
1400
1401        selectors.push(RowSelector {
1402            row_count: 0,
1403            skip: false,
1404        });
1405
1406        while total_rows < 7300 {
1407            let row_count: usize = rand.random_range(1..100);
1408
1409            let row_count = row_count.min(7300 - total_rows);
1410
1411            selectors.push(RowSelector { row_count, skip });
1412
1413            total_rows += row_count;
1414            if !skip {
1415                expected_rows += row_count;
1416            }
1417
1418            skip = !skip;
1419        }
1420
1421        let selection = RowSelection::from(selectors);
1422
1423        let async_reader = TestReader::new(data.clone());
1424
1425        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1426        let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1427            .await
1428            .unwrap();
1429
1430        assert_eq!(builder.metadata().num_row_groups(), 1);
1431
1432        let col_idx: usize = rand.random_range(0..13);
1433        let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]);
1434
1435        let stream = builder
1436            .with_projection(mask.clone())
1437            .with_row_selection(selection.clone())
1438            .build()
1439            .expect("building stream");
1440
1441        let async_batches: Vec<_> = stream.try_collect().await.unwrap();
1442
1443        let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum();
1444
1445        assert_eq!(actual_rows, expected_rows);
1446    }
1447
1448    #[tokio::test]
1449    async fn test_limit_multiple_row_groups() {
1450        let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]);
1451        let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]);
1452        let c = Int32Array::from_iter(0..6);
1453        let data = RecordBatch::try_from_iter([
1454            ("a", Arc::new(a) as ArrayRef),
1455            ("b", Arc::new(b) as ArrayRef),
1456            ("c", Arc::new(c) as ArrayRef),
1457        ])
1458        .unwrap();
1459
1460        let mut buf = Vec::with_capacity(1024);
1461        let props = WriterProperties::builder()
1462            .set_max_row_group_row_count(Some(3))
1463            .build();
1464        let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1465        writer.write(&data).unwrap();
1466        writer.close().unwrap();
1467
1468        let data: Bytes = buf.into();
1469        let metadata = ParquetMetaDataReader::new()
1470            .parse_and_finish(&data)
1471            .unwrap();
1472
1473        assert_eq!(metadata.num_row_groups(), 2);
1474
1475        let test = TestReader::new(data);
1476
1477        let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1478            .await
1479            .unwrap()
1480            .with_batch_size(1024)
1481            .with_limit(4)
1482            .build()
1483            .unwrap();
1484
1485        let batches: Vec<_> = stream.try_collect().await.unwrap();
1486        // Expect one batch for each row group
1487        assert_eq!(batches.len(), 2);
1488
1489        let batch = &batches[0];
1490        // First batch should contain all rows
1491        assert_eq!(batch.num_rows(), 3);
1492        assert_eq!(batch.num_columns(), 3);
1493        let col2 = batch.column(2).as_primitive::<Int32Type>();
1494        assert_eq!(col2.values(), &[0, 1, 2]);
1495
1496        let batch = &batches[1];
1497        // Second batch should trigger the limit and only have one row
1498        assert_eq!(batch.num_rows(), 1);
1499        assert_eq!(batch.num_columns(), 3);
1500        let col2 = batch.column(2).as_primitive::<Int32Type>();
1501        assert_eq!(col2.values(), &[3]);
1502
1503        let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1504            .await
1505            .unwrap()
1506            .with_offset(2)
1507            .with_limit(3)
1508            .build()
1509            .unwrap();
1510
1511        let batches: Vec<_> = stream.try_collect().await.unwrap();
1512        // Expect one batch for each row group
1513        assert_eq!(batches.len(), 2);
1514
1515        let batch = &batches[0];
1516        // First batch should contain one row
1517        assert_eq!(batch.num_rows(), 1);
1518        assert_eq!(batch.num_columns(), 3);
1519        let col2 = batch.column(2).as_primitive::<Int32Type>();
1520        assert_eq!(col2.values(), &[2]);
1521
1522        let batch = &batches[1];
1523        // Second batch should contain two rows
1524        assert_eq!(batch.num_rows(), 2);
1525        assert_eq!(batch.num_columns(), 3);
1526        let col2 = batch.column(2).as_primitive::<Int32Type>();
1527        assert_eq!(col2.values(), &[3, 4]);
1528
1529        let stream = ParquetRecordBatchStreamBuilder::new(test.clone())
1530            .await
1531            .unwrap()
1532            .with_offset(4)
1533            .with_limit(20)
1534            .build()
1535            .unwrap();
1536
1537        let batches: Vec<_> = stream.try_collect().await.unwrap();
1538        // Should skip first row group
1539        assert_eq!(batches.len(), 1);
1540
1541        let batch = &batches[0];
1542        // First batch should contain two rows
1543        assert_eq!(batch.num_rows(), 2);
1544        assert_eq!(batch.num_columns(), 3);
1545        let col2 = batch.column(2).as_primitive::<Int32Type>();
1546        assert_eq!(col2.values(), &[4, 5]);
1547    }
1548
1549    #[tokio::test]
1550    async fn test_batch_size_overallocate() {
1551        let testdata = arrow::util::test_util::parquet_test_data();
1552        // `alltypes_plain.parquet` only have 8 rows
1553        let path = format!("{testdata}/alltypes_plain.parquet");
1554        let data = Bytes::from(std::fs::read(path).unwrap());
1555
1556        let async_reader = TestReader::new(data.clone());
1557
1558        let builder = ParquetRecordBatchStreamBuilder::new(async_reader)
1559            .await
1560            .unwrap();
1561
1562        let file_rows = builder.metadata().file_metadata().num_rows() as usize;
1563
1564        let builder = builder
1565            .with_projection(ProjectionMask::all())
1566            .with_batch_size(1024);
1567
1568        // even though the batch size is set to 1024, it should adjust to the max
1569        // number of rows in the file (8)
1570        assert_ne!(1024, file_rows);
1571        assert_eq!(builder.batch_size, file_rows);
1572
1573        let _stream = builder.build().unwrap();
1574    }
1575
1576    #[tokio::test]
1577    async fn test_parquet_record_batch_stream_schema() {
1578        fn get_all_field_names(schema: &Schema) -> Vec<&String> {
1579            schema.flattened_fields().iter().map(|f| f.name()).collect()
1580        }
1581
1582        // ParquetRecordBatchReaderBuilder::schema differs from
1583        // ParquetRecordBatchReader::schema and RecordBatch::schema in the returned
1584        // schema contents (in terms of custom metadata attached to schema, and fields
1585        // returned). Test to ensure this remains consistent behaviour.
1586        //
1587        // Ensure same for asynchronous versions of the above.
1588
1589        // Prep data, for a schema with nested fields, with custom metadata
1590        let mut metadata = HashMap::with_capacity(1);
1591        metadata.insert("key".to_string(), "value".to_string());
1592
1593        let nested_struct_array = StructArray::from(vec![
1594            (
1595                Arc::new(Field::new("d", DataType::Utf8, true)),
1596                Arc::new(StringArray::from(vec!["a", "b"])) as ArrayRef,
1597            ),
1598            (
1599                Arc::new(Field::new("e", DataType::Utf8, true)),
1600                Arc::new(StringArray::from(vec!["c", "d"])) as ArrayRef,
1601            ),
1602        ]);
1603        let struct_array = StructArray::from(vec![
1604            (
1605                Arc::new(Field::new("a", DataType::Int32, true)),
1606                Arc::new(Int32Array::from(vec![-1, 1])) as ArrayRef,
1607            ),
1608            (
1609                Arc::new(Field::new("b", DataType::UInt64, true)),
1610                Arc::new(UInt64Array::from(vec![1, 2])) as ArrayRef,
1611            ),
1612            (
1613                Arc::new(Field::new(
1614                    "c",
1615                    nested_struct_array.data_type().clone(),
1616                    true,
1617                )),
1618                Arc::new(nested_struct_array) as ArrayRef,
1619            ),
1620        ]);
1621
1622        let schema =
1623            Arc::new(Schema::new(struct_array.fields().clone()).with_metadata(metadata.clone()));
1624        let record_batch = RecordBatch::from(struct_array)
1625            .with_schema(schema.clone())
1626            .unwrap();
1627
1628        // Write parquet with custom metadata in schema
1629        let mut file = tempfile().unwrap();
1630        let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1631        writer.write(&record_batch).unwrap();
1632        writer.close().unwrap();
1633
1634        let all_fields = ["a", "b", "c", "d", "e"];
1635        // (leaf indices in mask, expected names in output schema all fields)
1636        let projections = [
1637            (vec![], vec![]),
1638            (vec![0], vec!["a"]),
1639            (vec![0, 1], vec!["a", "b"]),
1640            (vec![0, 1, 2], vec!["a", "b", "c", "d"]),
1641            (vec![0, 1, 2, 3], vec!["a", "b", "c", "d", "e"]),
1642        ];
1643
1644        // Ensure we're consistent for each of these projections
1645        for (indices, expected_projected_names) in projections {
1646            let assert_schemas = |builder: SchemaRef, reader: SchemaRef, batch: SchemaRef| {
1647                // Builder schema should preserve all fields and metadata
1648                assert_eq!(get_all_field_names(&builder), all_fields);
1649                assert_eq!(builder.metadata, metadata);
1650                // Reader & batch schema should show only projected fields, and no metadata
1651                assert_eq!(get_all_field_names(&reader), expected_projected_names);
1652                assert_eq!(reader.metadata, HashMap::default());
1653                assert_eq!(get_all_field_names(&batch), expected_projected_names);
1654                assert_eq!(batch.metadata, HashMap::default());
1655            };
1656
1657            let builder =
1658                ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1659            let sync_builder_schema = builder.schema().clone();
1660            let mask = ProjectionMask::leaves(builder.parquet_schema(), indices.clone());
1661            let mut reader = builder.with_projection(mask).build().unwrap();
1662            let sync_reader_schema = reader.schema();
1663            let batch = reader.next().unwrap().unwrap();
1664            let sync_batch_schema = batch.schema();
1665            assert_schemas(sync_builder_schema, sync_reader_schema, sync_batch_schema);
1666
1667            // asynchronous should be same
1668            let file = tokio::fs::File::from(file.try_clone().unwrap());
1669            let builder = ParquetRecordBatchStreamBuilder::new(file).await.unwrap();
1670            let async_builder_schema = builder.schema().clone();
1671            let mask = ProjectionMask::leaves(builder.parquet_schema(), indices);
1672            let mut reader = builder.with_projection(mask).build().unwrap();
1673            let async_reader_schema = reader.schema().clone();
1674            let batch = reader.next().await.unwrap().unwrap();
1675            let async_batch_schema = batch.schema();
1676            assert_schemas(
1677                async_builder_schema,
1678                async_reader_schema,
1679                async_batch_schema,
1680            );
1681        }
1682    }
1683
1684    #[tokio::test]
1685    async fn test_nested_skip() {
1686        let schema = Arc::new(Schema::new(vec![
1687            Field::new("col_1", DataType::UInt64, false),
1688            Field::new_list("col_2", Field::new_list_field(DataType::Utf8, true), true),
1689        ]));
1690
1691        // Default writer properties
1692        let props = WriterProperties::builder()
1693            .set_data_page_row_count_limit(256)
1694            .set_write_batch_size(256)
1695            .set_max_row_group_row_count(Some(1024));
1696
1697        // Write data
1698        let mut file = tempfile().unwrap();
1699        let mut writer =
1700            ArrowWriter::try_new(&mut file, schema.clone(), Some(props.build())).unwrap();
1701
1702        let mut builder = ListBuilder::new(StringBuilder::new());
1703        for id in 0..1024 {
1704            match id % 3 {
1705                0 => builder.append_value([Some("val_1".to_string()), Some(format!("id_{id}"))]),
1706                1 => builder.append_value([Some(format!("id_{id}"))]),
1707                _ => builder.append_null(),
1708            }
1709        }
1710        let refs = vec![
1711            Arc::new(UInt64Array::from_iter_values(0..1024)) as ArrayRef,
1712            Arc::new(builder.finish()) as ArrayRef,
1713        ];
1714
1715        let batch = RecordBatch::try_new(schema.clone(), refs).unwrap();
1716        writer.write(&batch).unwrap();
1717        writer.close().unwrap();
1718
1719        let selections = [
1720            RowSelection::from(vec![
1721                RowSelector::skip(313),
1722                RowSelector::select(1),
1723                RowSelector::skip(709),
1724                RowSelector::select(1),
1725            ]),
1726            RowSelection::from(vec![
1727                RowSelector::skip(255),
1728                RowSelector::select(1),
1729                RowSelector::skip(767),
1730                RowSelector::select(1),
1731            ]),
1732            RowSelection::from(vec![
1733                RowSelector::select(255),
1734                RowSelector::skip(1),
1735                RowSelector::select(767),
1736                RowSelector::skip(1),
1737            ]),
1738            RowSelection::from(vec![
1739                RowSelector::skip(254),
1740                RowSelector::select(1),
1741                RowSelector::select(1),
1742                RowSelector::skip(767),
1743                RowSelector::select(1),
1744            ]),
1745        ];
1746
1747        for selection in selections {
1748            let expected = selection.row_count();
1749            // Read data
1750            let mut reader = ParquetRecordBatchStreamBuilder::new_with_options(
1751                tokio::fs::File::from_std(file.try_clone().unwrap()),
1752                ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
1753            )
1754            .await
1755            .unwrap();
1756
1757            reader = reader.with_row_selection(selection);
1758
1759            let mut stream = reader.build().unwrap();
1760
1761            let mut total_rows = 0;
1762            while let Some(rb) = stream.next().await {
1763                let rb = rb.unwrap();
1764                total_rows += rb.num_rows();
1765            }
1766            assert_eq!(total_rows, expected);
1767        }
1768    }
1769
1770    #[tokio::test]
1771    async fn empty_offset_index_doesnt_panic_in_read_row_group() {
1772        use tokio::fs::File;
1773        let testdata = arrow::util::test_util::parquet_test_data();
1774        let path = format!("{testdata}/alltypes_plain.parquet");
1775        let mut file = File::open(&path).await.unwrap();
1776        let file_size = file.metadata().await.unwrap().len();
1777        let mut metadata = ParquetMetaDataReader::new()
1778            .with_page_index_policy(PageIndexPolicy::Required)
1779            .load_and_finish(&mut file, file_size)
1780            .await
1781            .unwrap();
1782
1783        let page_index = PageIndex::new(None, Some(vec![]));
1784        metadata.set_page_index(Some(page_index));
1785        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1786        let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1787        let reader =
1788            ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1789                .build()
1790                .unwrap();
1791
1792        let result = reader.try_collect::<Vec<_>>().await.unwrap();
1793        assert_eq!(result.len(), 1);
1794    }
1795
1796    #[tokio::test]
1797    async fn non_empty_offset_index_doesnt_panic_in_read_row_group() {
1798        use tokio::fs::File;
1799        let testdata = arrow::util::test_util::parquet_test_data();
1800        let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1801        let mut file = File::open(&path).await.unwrap();
1802        let file_size = file.metadata().await.unwrap().len();
1803        let metadata = ParquetMetaDataReader::new()
1804            .with_page_index_policy(PageIndexPolicy::Required)
1805            .load_and_finish(&mut file, file_size)
1806            .await
1807            .unwrap();
1808
1809        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1810        let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1811        let reader =
1812            ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1813                .build()
1814                .unwrap();
1815
1816        let result = reader.try_collect::<Vec<_>>().await.unwrap();
1817        assert_eq!(result.len(), 8);
1818    }
1819
1820    #[tokio::test]
1821    async fn empty_offset_index_doesnt_panic_in_column_chunks() {
1822        use tempfile::TempDir;
1823        use tokio::fs::File;
1824        fn write_metadata_to_local_file(
1825            metadata: ParquetMetaData,
1826            file: impl AsRef<std::path::Path>,
1827        ) {
1828            use crate::file::metadata::ParquetMetaDataWriter;
1829            use std::fs::File;
1830            let file = File::create(file).unwrap();
1831            ParquetMetaDataWriter::new(file, &metadata)
1832                .finish()
1833                .unwrap()
1834        }
1835
1836        fn read_metadata_from_local_file(file: impl AsRef<std::path::Path>) -> ParquetMetaData {
1837            use std::fs::File;
1838            let file = File::open(file).unwrap();
1839            ParquetMetaDataReader::new()
1840                .with_page_index_policy(PageIndexPolicy::Required)
1841                .parse_and_finish(&file)
1842                .unwrap()
1843        }
1844
1845        let testdata = arrow::util::test_util::parquet_test_data();
1846        let path = format!("{testdata}/alltypes_plain.parquet");
1847        let mut file = File::open(&path).await.unwrap();
1848        let file_size = file.metadata().await.unwrap().len();
1849        let metadata = ParquetMetaDataReader::new()
1850            .with_page_index_policy(PageIndexPolicy::Required)
1851            .load_and_finish(&mut file, file_size)
1852            .await
1853            .unwrap();
1854
1855        let tempdir = TempDir::new().unwrap();
1856        let metadata_path = tempdir.path().join("thrift_metadata.dat");
1857        write_metadata_to_local_file(metadata, &metadata_path);
1858        let metadata = read_metadata_from_local_file(&metadata_path);
1859
1860        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1861        let arrow_reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1862        let reader =
1863            ParquetRecordBatchStreamBuilder::new_with_metadata(file, arrow_reader_metadata)
1864                .build()
1865                .unwrap();
1866
1867        // Panics here
1868        let result = reader.try_collect::<Vec<_>>().await.unwrap();
1869        assert_eq!(result.len(), 1);
1870    }
1871
1872    #[tokio::test]
1873    async fn test_cached_array_reader_sparse_offset_error() {
1874        use futures::TryStreamExt;
1875
1876        use crate::arrow::arrow_reader::{ArrowPredicateFn, RowFilter, RowSelection, RowSelector};
1877        use arrow_array::{BooleanArray, RecordBatch};
1878
1879        let testdata = arrow::util::test_util::parquet_test_data();
1880        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
1881        let data = Bytes::from(std::fs::read(path).unwrap());
1882
1883        let async_reader = TestReader::new(data);
1884
1885        // Enable page index so the fetch logic loads only required pages
1886        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1887        let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options)
1888            .await
1889            .unwrap();
1890
1891        // Skip the first 22 rows (entire first Parquet page) and then select the
1892        // next 3 rows (22, 23, 24). This means the fetch step will not include
1893        // the first page starting at file offset 0.
1894        let selection = RowSelection::from(vec![RowSelector::skip(22), RowSelector::select(3)]);
1895
1896        // Trivial predicate on column 0 that always returns `true`. Using the
1897        // same column in both predicate and projection activates the caching
1898        // layer (Producer/Consumer pattern).
1899        let parquet_schema = builder.parquet_schema();
1900        let proj = ProjectionMask::leaves(parquet_schema, vec![0]);
1901        let always_true = ArrowPredicateFn::new(proj.clone(), |batch: RecordBatch| {
1902            Ok(BooleanArray::from(vec![true; batch.num_rows()]))
1903        });
1904        let filter = RowFilter::new(vec![Box::new(always_true)]);
1905
1906        // Build the stream with batch size 8 so the cache reads whole batches
1907        // that straddle the requested row range (rows 0-7, 8-15, 16-23, …).
1908        let stream = builder
1909            .with_batch_size(8)
1910            .with_projection(proj)
1911            .with_row_selection(selection)
1912            .with_row_filter(filter)
1913            .build()
1914            .unwrap();
1915
1916        // Collecting the stream should fail with the sparse column chunk offset
1917        // error we want to reproduce.
1918        let _result: Vec<_> = stream.try_collect().await.unwrap();
1919    }
1920
1921    #[tokio::test]
1922    async fn test_predicate_cache_disabled() {
1923        let k = Int32Array::from_iter_values(0..10);
1924        let data = RecordBatch::try_from_iter([("k", Arc::new(k) as ArrayRef)]).unwrap();
1925
1926        let mut buf = Vec::new();
1927        // both the page row limit and batch size are set to 1 to create one page per row
1928        let props = WriterProperties::builder()
1929            .set_data_page_row_count_limit(1)
1930            .set_write_batch_size(1)
1931            .set_max_row_group_row_count(Some(10))
1932            .set_write_page_header_statistics(true)
1933            .build();
1934        let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap();
1935        writer.write(&data).unwrap();
1936        writer.close().unwrap();
1937
1938        let data = Bytes::from(buf);
1939        let metadata = ParquetMetaDataReader::new()
1940            .with_page_index_policy(PageIndexPolicy::Required)
1941            .parse_and_finish(&data)
1942            .unwrap();
1943        let parquet_schema = metadata.file_metadata().schema_descr_ptr();
1944
1945        // the filter is not clone-able, so we use a lambda to simplify
1946        let build_filter = || {
1947            let scalar = Int32Array::from_iter_values([5]);
1948            let predicate = ArrowPredicateFn::new(
1949                ProjectionMask::leaves(&parquet_schema, vec![0]),
1950                move |batch| eq(batch.column(0), &Scalar::new(&scalar)),
1951            );
1952            RowFilter::new(vec![Box::new(predicate)])
1953        };
1954
1955        // select only one of the pages
1956        let selection = RowSelection::from(vec![RowSelector::skip(5), RowSelector::select(1)]);
1957
1958        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
1959        let reader_metadata = ArrowReaderMetadata::try_new(metadata.into(), options).unwrap();
1960
1961        // using the predicate cache (default)
1962        let reader_with_cache = TestReader::new(data.clone());
1963        let requests_with_cache = reader_with_cache.requests.clone();
1964        let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1965            reader_with_cache,
1966            reader_metadata.clone(),
1967        )
1968        .with_batch_size(1000)
1969        .with_row_selection(selection.clone())
1970        .with_row_filter(build_filter())
1971        .build()
1972        .unwrap();
1973        let batches_with_cache: Vec<_> = stream.try_collect().await.unwrap();
1974
1975        // disabling the predicate cache
1976        let reader_without_cache = TestReader::new(data);
1977        let requests_without_cache = reader_without_cache.requests.clone();
1978        let stream = ParquetRecordBatchStreamBuilder::new_with_metadata(
1979            reader_without_cache,
1980            reader_metadata,
1981        )
1982        .with_batch_size(1000)
1983        .with_row_selection(selection)
1984        .with_row_filter(build_filter())
1985        .with_max_predicate_cache_size(0) // disabling it by setting the limit to 0
1986        .build()
1987        .unwrap();
1988        let batches_without_cache: Vec<_> = stream.try_collect().await.unwrap();
1989
1990        assert_eq!(batches_with_cache, batches_without_cache);
1991
1992        let requests_with_cache = requests_with_cache.lock().unwrap();
1993        let requests_without_cache = requests_without_cache.lock().unwrap();
1994
1995        // less requests will be made without the predicate cache
1996        assert_eq!(requests_with_cache.len(), 11);
1997        assert_eq!(requests_without_cache.len(), 2);
1998
1999        // less bytes will be retrieved without the predicate cache
2000        assert_eq!(
2001            requests_with_cache.iter().map(|r| r.len()).sum::<usize>(),
2002            433
2003        );
2004        assert_eq!(
2005            requests_without_cache
2006                .iter()
2007                .map(|r| r.len())
2008                .sum::<usize>(),
2009            92
2010        );
2011    }
2012
2013    #[test]
2014    fn test_row_numbers_with_multiple_row_groups() {
2015        test_row_numbers_with_multiple_row_groups_helper(
2016            false,
2017            |path, selection, _row_filter, batch_size| {
2018                let runtime = tokio::runtime::Builder::new_current_thread()
2019                    .enable_all()
2020                    .build()
2021                    .expect("Could not create runtime");
2022                runtime.block_on(async move {
2023                    let file = tokio::fs::File::open(path).await.unwrap();
2024                    let row_number_field = Arc::new(
2025                        Field::new("row_number", DataType::Int64, false)
2026                            .with_extension_type(RowNumber),
2027                    );
2028                    let options = ArrowReaderOptions::new()
2029                        .with_virtual_columns(vec![row_number_field])
2030                        .unwrap();
2031                    let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2032                        .await
2033                        .unwrap()
2034                        .with_row_selection(selection)
2035                        .with_batch_size(batch_size)
2036                        .build()
2037                        .expect("Could not create reader");
2038                    reader.try_collect::<Vec<_>>().await.unwrap()
2039                })
2040            },
2041        );
2042    }
2043
2044    #[test]
2045    fn test_row_numbers_with_multiple_row_groups_and_filter() {
2046        test_row_numbers_with_multiple_row_groups_helper(
2047            true,
2048            |path, selection, row_filter, batch_size| {
2049                let runtime = tokio::runtime::Builder::new_current_thread()
2050                    .enable_all()
2051                    .build()
2052                    .expect("Could not create runtime");
2053                runtime.block_on(async move {
2054                    let file = tokio::fs::File::open(path).await.unwrap();
2055                    let row_number_field = Arc::new(
2056                        Field::new("row_number", DataType::Int64, false)
2057                            .with_extension_type(RowNumber),
2058                    );
2059                    let options = ArrowReaderOptions::new()
2060                        .with_virtual_columns(vec![row_number_field])
2061                        .unwrap();
2062                    let reader = ParquetRecordBatchStreamBuilder::new_with_options(file, options)
2063                        .await
2064                        .unwrap()
2065                        .with_row_selection(selection)
2066                        .with_row_filter(row_filter.expect("No row filter"))
2067                        .with_batch_size(batch_size)
2068                        .build()
2069                        .expect("Could not create reader");
2070                    reader.try_collect::<Vec<_>>().await.unwrap()
2071                })
2072            },
2073        );
2074    }
2075
2076    #[tokio::test]
2077    async fn test_nested_lists() -> Result<()> {
2078        // Test case for https://github.com/apache/arrow-rs/issues/8657
2079        let list_inner_field = Arc::new(Field::new("item", DataType::Float32, true));
2080        let table_schema = Arc::new(Schema::new(vec![
2081            Field::new("id", DataType::Int32, false),
2082            Field::new("vector", DataType::List(list_inner_field.clone()), true),
2083        ]));
2084
2085        let mut list_builder =
2086            ListBuilder::new(Float32Builder::new()).with_field(list_inner_field.clone());
2087        list_builder.values().append_slice(&[10.0, 10.0, 10.0]);
2088        list_builder.append(true);
2089        list_builder.values().append_slice(&[20.0, 20.0, 20.0]);
2090        list_builder.append(true);
2091        list_builder.values().append_slice(&[30.0, 30.0, 30.0]);
2092        list_builder.append(true);
2093        list_builder.values().append_slice(&[40.0, 40.0, 40.0]);
2094        list_builder.append(true);
2095        let list_array = list_builder.finish();
2096
2097        let data = vec![RecordBatch::try_new(
2098            table_schema.clone(),
2099            vec![
2100                Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
2101                Arc::new(list_array),
2102            ],
2103        )?];
2104
2105        let mut buffer = Vec::new();
2106        let mut writer = AsyncArrowWriter::try_new(&mut buffer, table_schema, None)?;
2107
2108        for batch in data {
2109            writer.write(&batch).await?;
2110        }
2111
2112        writer.close().await?;
2113
2114        let reader = TestReader::new(Bytes::from(buffer));
2115        let builder = ParquetRecordBatchStreamBuilder::new(reader).await?;
2116
2117        let predicate = ArrowPredicateFn::new(ProjectionMask::all(), |batch| {
2118            Ok(BooleanArray::from(vec![true; batch.num_rows()]))
2119        });
2120
2121        let projection_mask = ProjectionMask::all();
2122
2123        let mut stream = builder
2124            .with_row_filter(RowFilter::new(vec![Box::new(predicate)]))
2125            .with_projection(projection_mask)
2126            .build()?;
2127
2128        while let Some(batch) = stream.next().await {
2129            let _ = batch.unwrap(); // ensure there is no panic
2130        }
2131
2132        Ok(())
2133    }
2134}