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