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