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