Skip to main content

parquet/arrow/arrow_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//! Contains reader which reads parquet data into arrow [`RecordBatch`]
19
20use arrow_array::cast::AsArray;
21use arrow_array::{BooleanArray, RecordBatch, RecordBatchReader};
22use arrow_buffer::{BooleanBuffer, BooleanBufferBuilder};
23use arrow_schema::{ArrowError, DataType as ArrowType, FieldRef, Schema, SchemaRef};
24use arrow_select::filter::filter_record_batch;
25pub use filter::{ArrowPredicate, ArrowPredicateFn, RowFilter};
26use selection::MaskCursor;
27pub use selection::{
28    MaskRunIter, RowSelection, RowSelectionCursor, RowSelectionPolicy, RowSelector,
29};
30use std::fmt::{Debug, Formatter};
31use std::sync::Arc;
32
33pub use crate::arrow::array_reader::RowGroups;
34use crate::arrow::array_reader::{ArrayReader, ArrayReaderBuilder};
35use crate::arrow::schema::{
36    ParquetField, parquet_to_arrow_schema_and_fields, virtual_type::is_virtual_column,
37};
38use crate::arrow::{FieldLevels, ProjectionMask, parquet_to_arrow_field_levels_with_virtual};
39use crate::basic::{BloomFilterAlgorithm, BloomFilterCompression, BloomFilterHash};
40use crate::bloom_filter::{
41    SBBF_HEADER_SIZE_ESTIMATE, Sbbf, chunk_read_bloom_filter_header_and_offset,
42};
43use crate::column::page::{PageIterator, PageReader};
44#[cfg(feature = "encryption")]
45use crate::encryption::decrypt::FileDecryptionProperties;
46use crate::errors::{ParquetError, Result};
47use crate::file::metadata::{
48    PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions, ParquetMetaDataReader,
49    ParquetStatisticsPolicy, RowGroupMetaData,
50};
51use crate::file::reader::{ChunkReader, SerializedPageReader};
52use crate::schema::types::SchemaDescriptor;
53
54use crate::arrow::arrow_reader::metrics::ArrowReaderMetrics;
55// Exposed so integration tests and benchmarks can temporarily override the threshold.
56pub use read_plan::{PredicateOptions, ReadPlan, ReadPlanBuilder};
57
58mod filter;
59pub mod metrics;
60mod read_plan;
61pub(crate) mod selection;
62pub mod statistics;
63
64/// Default batch size for reading parquet files
65pub const DEFAULT_BATCH_SIZE: usize = 1024;
66
67/// A row group and its optional row-group-local [`RowSelection`].
68///
69/// A row-group-local selection is relative to the rows in this row group. For
70/// example, an offset of 100 refers to the row at offset 100 within the row
71/// group, not within the Parquet file.
72///
73/// A `None` selection reads the entire row group. Omitting a row group skips
74/// it. Entries are decoded in the supplied order.
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct RowGroupSelection {
77    pub(crate) row_group_index: usize,
78    pub(crate) selection: Option<RowSelection>,
79}
80
81impl RowGroupSelection {
82    /// Creates a row-group-local selection.
83    pub fn new(row_group_index: usize, selection: Option<RowSelection>) -> Self {
84        Self {
85            row_group_index,
86            selection,
87        }
88    }
89
90    /// The index of the row group this selection applies to.
91    pub fn row_group_index(&self) -> usize {
92        self.row_group_index
93    }
94
95    /// The row-group-local selection, or `None` if the entire row group is
96    /// read.
97    pub fn selection(&self) -> Option<&RowSelection> {
98        self.selection.as_ref()
99    }
100}
101
102/// Row-selection configuration shared by the Arrow reader builders.
103#[derive(Debug)]
104pub(crate) enum RowGroupPlan {
105    /// First select `row_groups`, if provided, and then apply `selection`
106    /// across the concatenated rows from those row groups.
107    ///
108    /// This is formed by [`ArrowReaderBuilder::with_row_groups`] and
109    /// [`ArrowReaderBuilder::with_row_selection`].
110    Global {
111        row_groups: Option<Vec<usize>>,
112        selection: Option<RowSelection>,
113    },
114    /// Apply each row-group-local selection independently, in the order
115    /// supplied.
116    ///
117    /// This is formed by `with_row_group_selections` on the push decoder and
118    /// async stream builders.
119    PerRowGroup(Vec<RowGroupSelection>),
120    /// Mutually exclusive global and per-row-group configuration was supplied.
121    /// This is reported as an error when the reader is built.
122    Conflicting,
123}
124
125impl RowGroupPlan {
126    fn set_row_groups(&mut self, new_row_groups: Vec<usize>) {
127        match self {
128            Self::Global { row_groups, .. } => *row_groups = Some(new_row_groups),
129            Self::PerRowGroup(_) => *self = Self::Conflicting,
130            Self::Conflicting => {}
131        }
132    }
133
134    fn set_row_selection(&mut self, new_selection: RowSelection) {
135        match self {
136            Self::Global { selection, .. } => *selection = Some(new_selection),
137            Self::PerRowGroup(_) => *self = Self::Conflicting,
138            Self::Conflicting => {}
139        }
140    }
141
142    pub(crate) fn set_row_group_selections(
143        &mut self,
144        row_group_selections: Vec<RowGroupSelection>,
145    ) {
146        match self {
147            Self::Global {
148                row_groups: None,
149                selection: None,
150            }
151            | Self::PerRowGroup(_) => {
152                *self = Self::PerRowGroup(row_group_selections);
153            }
154            Self::Global { .. } => *self = Self::Conflicting,
155            Self::Conflicting => {}
156        }
157    }
158
159    pub(crate) fn conflict_error() -> ParquetError {
160        ParquetError::General(
161            "with_row_group_selections cannot be combined with with_row_groups or with_row_selection"
162                .to_string(),
163        )
164    }
165
166    fn into_global(self) -> Result<(Option<Vec<usize>>, Option<RowSelection>)> {
167        match self {
168            Self::Global {
169                row_groups,
170                selection,
171            } => Ok((row_groups, selection)),
172            Self::PerRowGroup(_) => Err(ParquetError::General(
173                "Row-group-local selections are not supported by the synchronous reader"
174                    .to_string(),
175            )),
176            Self::Conflicting => Err(Self::conflict_error()),
177        }
178    }
179}
180
181/// Builder for constructing Parquet readers that decode into [Apache Arrow]
182/// arrays.
183///
184/// Most users should use one of the following specializations:
185///
186/// * synchronous API: [`ParquetRecordBatchReaderBuilder`]
187/// * `async` API: [`ParquetRecordBatchStreamBuilder`]
188/// * decoder API: [`ParquetPushDecoderBuilder`]
189///
190/// # Features
191/// * Projection pushdown: [`Self::with_projection`]
192/// * Cached metadata: [`ArrowReaderMetadata::load`]
193/// * Offset skipping: [`Self::with_offset`] and [`Self::with_limit`]
194/// * Row group filtering: [`Self::with_row_groups`]
195/// * Range filtering: [`Self::with_row_selection`]
196/// * Row level filtering: [`Self::with_row_filter`]
197///
198/// # Implementing Predicate Pushdown
199///
200/// [`Self::with_row_filter`] permits filter evaluation *during* the decoding
201/// process, which is efficient and allows the most low level optimizations.
202///
203/// However, most Parquet based systems will apply filters at many steps prior
204/// to decoding such as pruning files, row groups and data pages. This crate
205/// provides the low level APIs needed to implement such filtering, but does not
206/// include any logic to actually evaluate predicates. For example:
207///
208/// * [`Self::with_row_groups`] for Row Group pruning
209/// * [`Self::with_row_selection`] for data page pruning
210/// * [`StatisticsConverter`] to convert Parquet statistics to Arrow arrays
211///
212/// The rationale for this design is that implementing predicate pushdown is a
213/// complex topic and varies significantly from system to system. For example
214///
215/// 1. Predicates supported (do you support predicates like prefix matching, user defined functions, etc)
216/// 2. Evaluating predicates on multiple files (with potentially different but compatible schemas)
217/// 3. Evaluating predicates using information from an external metadata catalog (e.g. Apache Iceberg or similar)
218/// 4. Interleaving fetching metadata, evaluating predicates, and decoding files
219///
220/// You can read more about this design in the [Querying Parquet with
221/// Millisecond Latency] Arrow blog post.
222///
223/// [`ParquetRecordBatchStreamBuilder`]: crate::arrow::async_reader::ParquetRecordBatchStreamBuilder
224/// [`ParquetPushDecoderBuilder`]: crate::arrow::push_decoder::ParquetPushDecoderBuilder
225/// [Apache Arrow]: https://arrow.apache.org/
226/// [`StatisticsConverter`]: statistics::StatisticsConverter
227/// [Querying Parquet with Millisecond Latency]: https://arrow.apache.org/blog/2022/12/26/querying-parquet-with-millisecond-latency/
228pub struct ArrowReaderBuilder<T> {
229    /// The "input" to read parquet data from.
230    ///
231    /// Note in the case of the [`ParquetPushDecoderBuilder`] there is no
232    /// underlying reader; the input is instead [`PushDecoderInput`], the buffer that
233    /// caller-pushed bytes accumulate in.
234    ///
235    /// [`ParquetPushDecoderBuilder`]: crate::arrow::push_decoder::ParquetPushDecoderBuilder
236    /// [`PushDecoderInput`]: crate::arrow::push_decoder::PushDecoderInput
237    pub(crate) input: T,
238
239    pub(crate) metadata: Arc<ParquetMetaData>,
240
241    pub(crate) schema: SchemaRef,
242
243    pub(crate) fields: Option<Arc<ParquetField>>,
244
245    pub(crate) batch_size: usize,
246
247    pub(crate) row_group_plan: RowGroupPlan,
248
249    pub(crate) projection: ProjectionMask,
250
251    pub(crate) filter: Option<RowFilter>,
252
253    pub(crate) row_selection_policy: RowSelectionPolicy,
254
255    pub(crate) limit: Option<usize>,
256
257    pub(crate) offset: Option<usize>,
258
259    pub(crate) metrics: ArrowReaderMetrics,
260
261    pub(crate) max_predicate_cache_size: usize,
262}
263
264impl<T: Debug> Debug for ArrowReaderBuilder<T> {
265    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
266        f.debug_struct("ArrowReaderBuilder<T>")
267            .field("input", &self.input)
268            .field("metadata", &self.metadata)
269            .field("schema", &self.schema)
270            .field("fields", &self.fields)
271            .field("batch_size", &self.batch_size)
272            .field("row_group_plan", &self.row_group_plan)
273            .field("projection", &self.projection)
274            .field("filter", &self.filter)
275            .field("row_selection_policy", &self.row_selection_policy)
276            .field("limit", &self.limit)
277            .field("offset", &self.offset)
278            .field("metrics", &self.metrics)
279            .finish()
280    }
281}
282
283impl<T> ArrowReaderBuilder<T> {
284    pub(crate) fn new_builder(input: T, metadata: ArrowReaderMetadata) -> Self {
285        Self {
286            input,
287            metadata: metadata.metadata,
288            schema: metadata.schema,
289            fields: metadata.fields,
290            batch_size: DEFAULT_BATCH_SIZE,
291            row_group_plan: RowGroupPlan::Global {
292                row_groups: None,
293                selection: None,
294            },
295            projection: ProjectionMask::all(),
296            filter: None,
297            row_selection_policy: RowSelectionPolicy::default(),
298            limit: None,
299            offset: None,
300            metrics: ArrowReaderMetrics::Disabled,
301            max_predicate_cache_size: 100 * 1024 * 1024, // 100MB default cache size
302        }
303    }
304
305    /// Returns a reference to the [`ParquetMetaData`] for this parquet file
306    pub fn metadata(&self) -> &Arc<ParquetMetaData> {
307        &self.metadata
308    }
309
310    /// Returns the parquet [`SchemaDescriptor`] for this parquet file
311    pub fn parquet_schema(&self) -> &SchemaDescriptor {
312        self.metadata.file_metadata().schema_descr()
313    }
314
315    /// Returns the arrow [`SchemaRef`] for this parquet file
316    pub fn schema(&self) -> &SchemaRef {
317        &self.schema
318    }
319
320    /// Set the size of [`RecordBatch`] to produce. Defaults to [`DEFAULT_BATCH_SIZE`].
321    ///
322    /// This may be used as a hint for internal allocations, but does not
323    /// guarantee exact internal buffer capacities.
324    ///
325    /// If `batch_size` is more than the file row count, use the file row count.
326    pub fn with_batch_size(self, batch_size: usize) -> Self {
327        // Try to avoid allocate large buffer
328        let batch_size = batch_size.min(self.metadata.file_metadata().num_rows() as usize);
329        Self { batch_size, ..self }
330    }
331
332    /// Only read data from the provided row group indexes
333    ///
334    /// This is also called row group filtering
335    ///
336    /// On [`ParquetPushDecoderBuilder`] and [`ParquetRecordBatchStreamBuilder`],
337    /// which additionally offer `with_row_group_selections`, this cannot be
338    /// combined with that method; attempting to do so returns an error from
339    /// `build`.
340    ///
341    /// [`ParquetPushDecoderBuilder`]: crate::arrow::push_decoder::ParquetPushDecoderBuilder
342    /// [`ParquetRecordBatchStreamBuilder`]: crate::arrow::async_reader::ParquetRecordBatchStreamBuilder
343    pub fn with_row_groups(mut self, row_groups: Vec<usize>) -> Self {
344        self.row_group_plan.set_row_groups(row_groups);
345        self
346    }
347
348    /// Only read data from the provided column indexes
349    pub fn with_projection(self, mask: ProjectionMask) -> Self {
350        Self {
351            projection: mask,
352            ..self
353        }
354    }
355
356    /// Configure how row selections should be materialised during execution
357    ///
358    /// See [`RowSelectionPolicy`] for more details
359    pub fn with_row_selection_policy(self, policy: RowSelectionPolicy) -> Self {
360        Self {
361            row_selection_policy: policy,
362            ..self
363        }
364    }
365
366    /// Provide a [`RowSelection`] to filter out rows, and avoid fetching their
367    /// data into memory.
368    ///
369    /// This feature is used to restrict which rows are decoded within row
370    /// groups, skipping ranges of rows that are not needed. Such selections
371    /// could be determined by evaluating predicates against the parquet page
372    /// [`Index`] or some other external information available to a query
373    /// engine.
374    ///
375    /// # Notes
376    ///
377    /// Row group filtering (see [`Self::with_row_groups`]) is applied prior to
378    /// applying the row selection, and therefore rows from skipped row groups
379    /// should not be included in the [`RowSelection`] (see example below)
380    ///
381    /// On [`ParquetPushDecoderBuilder`] and [`ParquetRecordBatchStreamBuilder`],
382    /// which additionally offer `with_row_group_selections`, this cannot be
383    /// combined with that method; attempting to do so returns an error from
384    /// `build`.
385    ///
386    /// [`ParquetPushDecoderBuilder`]: crate::arrow::push_decoder::ParquetPushDecoderBuilder
387    /// [`ParquetRecordBatchStreamBuilder`]: crate::arrow::async_reader::ParquetRecordBatchStreamBuilder
388    ///
389    /// It is recommended to enable writing the page index if using this
390    /// functionality, to allow more efficient skipping over data pages. See
391    /// [`ArrowReaderOptions::with_page_index_policy`].
392    ///
393    /// # Example
394    ///
395    /// Given a parquet file with 4 row groups, and a row group filter of `[0,
396    /// 2, 3]`, in order to scan rows 50-100 in row group 2 and rows 200-300 in
397    /// row group 3:
398    ///
399    /// ```text
400    ///   Row Group 0, 1000 rows (selected)
401    ///   Row Group 1, 1000 rows (skipped)
402    ///   Row Group 2, 1000 rows (selected, but want to only scan rows 50-100)
403    ///   Row Group 3, 1000 rows (selected, but want to only scan rows 200-300)
404    /// ```
405    ///
406    /// You could pass the following [`RowSelection`]:
407    ///
408    /// ```text
409    ///  Select 1000    (scan all rows in row group 0)
410    ///  Skip 50        (skip the first 50 rows in row group 2)
411    ///  Select 50      (scan rows 50-100 in row group 2)
412    ///  Skip 900       (skip the remaining rows in row group 2)
413    ///  Skip 200       (skip the first 200 rows in row group 3)
414    ///  Select 100     (scan rows 200-300 in row group 3)
415    ///  Skip 700       (skip the remaining rows in row group 3)
416    /// ```
417    /// Note there is no entry for the (entirely) skipped row group 1.
418    ///
419    /// Note you can represent the same selection with fewer entries. Instead of
420    ///
421    /// ```text
422    ///  Skip 900       (skip the remaining rows in row group 2)
423    ///  Skip 200       (skip the first 200 rows in row group 3)
424    /// ```
425    ///
426    /// you could use
427    ///
428    /// ```text
429    /// Skip 1100      (skip the remaining 900 rows in row group 2 and the first 200 rows in row group 3)
430    /// ```
431    ///
432    /// [`Index`]: crate::file::page_index::column_index::ColumnIndexMetaData
433    pub fn with_row_selection(mut self, selection: RowSelection) -> Self {
434        self.row_group_plan.set_row_selection(selection);
435        self
436    }
437
438    /// Provide a [`RowFilter`] to skip decoding rows
439    ///
440    /// Row filters are applied after row group selection and row selection
441    ///
442    /// It is recommended to enable reading the page index if using this functionality, to allow
443    /// more efficient skipping over data pages. See [`ArrowReaderOptions::with_page_index_policy`].
444    ///
445    /// See the [blog post on late materialization] for a more technical explanation.
446    ///
447    /// [blog post on late materialization]: https://arrow.apache.org/blog/2025/12/11/parquet-late-materialization-deep-dive
448    ///
449    /// # Example
450    /// ```rust
451    /// # use std::fs::File;
452    /// # use arrow_array::Int32Array;
453    /// # use parquet::arrow::ProjectionMask;
454    /// # use parquet::arrow::arrow_reader::{ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowFilter};
455    /// # fn main() -> Result<(), parquet::errors::ParquetError> {
456    /// # let testdata = arrow::util::test_util::parquet_test_data();
457    /// # let path = format!("{testdata}/alltypes_plain.parquet");
458    /// # let file = File::open(&path)?;
459    /// let builder = ParquetRecordBatchReaderBuilder::try_new(file)?;
460    /// let schema_desc = builder.metadata().file_metadata().schema_descr_ptr();
461    /// // Create predicate that evaluates `int_col != 1`.
462    /// // `int_col` column has index 4 (zero based) in the schema
463    /// let projection = ProjectionMask::leaves(&schema_desc, [4]);
464    /// // Only the projection columns are passed to the predicate so
465    /// // int_col is column 0 in the predicate
466    /// let predicate = ArrowPredicateFn::new(projection, |batch| {
467    ///     let int_col = batch.column(0);
468    ///     arrow::compute::kernels::cmp::neq(int_col, &Int32Array::new_scalar(1))
469    /// });
470    /// let row_filter = RowFilter::new(vec![Box::new(predicate)]);
471    /// // The filter will be invoked during the reading process
472    /// let reader = builder.with_row_filter(row_filter).build()?;
473    /// # for b in reader { let _ = b?; }
474    /// # Ok(())
475    /// # }
476    /// ```
477    pub fn with_row_filter(self, filter: RowFilter) -> Self {
478        Self {
479            filter: Some(filter),
480            ..self
481        }
482    }
483
484    /// Provide a limit to the number of rows to be read
485    ///
486    /// The limit will be applied after any [`Self::with_row_selection`] and [`Self::with_row_filter`]
487    /// allowing it to limit the final set of rows decoded after any pushed down predicates
488    ///
489    /// It is recommended to enable reading the page index if using this functionality, to allow
490    /// more efficient skipping over data pages. See [`ArrowReaderOptions::with_page_index_policy`]
491    pub fn with_limit(self, limit: usize) -> Self {
492        Self {
493            limit: Some(limit),
494            ..self
495        }
496    }
497
498    /// Provide an offset to skip over the given number of rows
499    ///
500    /// The offset will be applied after any [`Self::with_row_selection`] and [`Self::with_row_filter`]
501    /// allowing it to skip rows after any pushed down predicates
502    ///
503    /// It is recommended to enable reading the page index if using this functionality, to allow
504    /// more efficient skipping over data pages. See [`ArrowReaderOptions::with_page_index_policy`]
505    pub fn with_offset(self, offset: usize) -> Self {
506        Self {
507            offset: Some(offset),
508            ..self
509        }
510    }
511
512    /// Specify metrics collection during reading
513    ///
514    /// To access the metrics, create an [`ArrowReaderMetrics`] and pass a
515    /// clone of the provided metrics to the builder.
516    ///
517    /// For example:
518    ///
519    /// ```rust
520    /// # use std::sync::Arc;
521    /// # use bytes::Bytes;
522    /// # use arrow_array::{Int32Array, RecordBatch};
523    /// # use arrow_schema::{DataType, Field, Schema};
524    /// # use parquet::arrow::arrow_reader::{ParquetRecordBatchReader, ParquetRecordBatchReaderBuilder};
525    /// use parquet::arrow::arrow_reader::metrics::ArrowReaderMetrics;
526    /// # use parquet::arrow::ArrowWriter;
527    /// # let mut file: Vec<u8> = Vec::with_capacity(1024);
528    /// # let schema = Arc::new(Schema::new(vec![Field::new("i32", DataType::Int32, false)]));
529    /// # let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
530    /// # let batch = RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1, 2, 3]))]).unwrap();
531    /// # writer.write(&batch).unwrap();
532    /// # writer.close().unwrap();
533    /// # let file = Bytes::from(file);
534    /// // Create metrics object to pass into the reader
535    /// let metrics = ArrowReaderMetrics::enabled();
536    /// let reader = ParquetRecordBatchReaderBuilder::try_new(file).unwrap()
537    ///   // Configure the builder to use the metrics by passing a clone
538    ///   .with_metrics(metrics.clone())
539    ///   // Build the reader
540    ///   .build().unwrap();
541    /// // .. read data from the reader ..
542    ///
543    /// // check the metrics
544    /// assert!(metrics.records_read_from_inner().is_some());
545    /// ```
546    pub fn with_metrics(self, metrics: ArrowReaderMetrics) -> Self {
547        Self { metrics, ..self }
548    }
549
550    /// Set the maximum size (per row group) of the predicate cache in bytes for
551    /// the async decoder.
552    ///
553    /// Defaults to 100MB (across all columns). Set to `usize::MAX` to use
554    /// unlimited cache size.
555    ///
556    /// This cache is used to store decoded arrays that are used in
557    /// predicate evaluation ([`Self::with_row_filter`]).
558    ///
559    /// This cache is only used for the "async" decoder, [`ParquetRecordBatchStream`]. See
560    /// [this ticket] for more details and alternatives.
561    ///
562    /// [`ParquetRecordBatchStream`]: https://docs.rs/parquet/latest/parquet/arrow/async_reader/struct.ParquetRecordBatchStream.html
563    /// [this ticket]: https://github.com/apache/arrow-rs/issues/8000
564    pub fn with_max_predicate_cache_size(self, max_predicate_cache_size: usize) -> Self {
565        Self {
566            max_predicate_cache_size,
567            ..self
568        }
569    }
570}
571
572/// Options that control how [`ParquetMetaData`] is read when constructing
573/// an Arrow reader.
574///
575/// To use these options, pass them to one of the following methods:
576/// * [`ParquetRecordBatchReaderBuilder::try_new_with_options`]
577/// * [`ParquetRecordBatchStreamBuilder::new_with_options`]
578///
579/// For fine-grained control over metadata loading, use
580/// [`ArrowReaderMetadata::load`] to load metadata with these options,
581///
582/// See [`ArrowReaderBuilder`] for how to configure how the column data
583/// is then read from the file, including projection and filter pushdown
584///
585/// [`ParquetRecordBatchStreamBuilder::new_with_options`]: crate::arrow::async_reader::ParquetRecordBatchStreamBuilder::new_with_options
586#[derive(Debug, Clone, Default)]
587pub struct ArrowReaderOptions {
588    /// Should the reader strip any user defined metadata from the Arrow schema
589    skip_arrow_metadata: bool,
590    /// If provided, used as the schema hint when determining the Arrow schema,
591    /// otherwise the schema hint is read from the [ARROW_SCHEMA_META_KEY]
592    ///
593    /// [ARROW_SCHEMA_META_KEY]: crate::arrow::ARROW_SCHEMA_META_KEY
594    supplied_schema: Option<SchemaRef>,
595
596    pub(crate) column_index: PageIndexPolicy,
597    pub(crate) offset_index: PageIndexPolicy,
598
599    /// Options to control reading of Parquet metadata
600    metadata_options: ParquetMetaDataOptions,
601    /// If encryption is enabled, the file decryption properties can be provided
602    #[cfg(feature = "encryption")]
603    pub(crate) file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
604
605    virtual_columns: Vec<FieldRef>,
606}
607
608impl ArrowReaderOptions {
609    /// Create a new [`ArrowReaderOptions`] with the default settings
610    pub fn new() -> Self {
611        Self::default()
612    }
613
614    /// Skip decoding the embedded arrow metadata (defaults to `false`)
615    ///
616    /// Parquet files generated by some writers may contain embedded arrow
617    /// schema and metadata.
618    /// This may not be correct or compatible with your system,
619    /// for example, see [ARROW-16184](https://issues.apache.org/jira/browse/ARROW-16184)
620    pub fn with_skip_arrow_metadata(self, skip_arrow_metadata: bool) -> Self {
621        Self {
622            skip_arrow_metadata,
623            ..self
624        }
625    }
626
627    /// Provide a schema hint to use when reading the Parquet file.
628    ///
629    /// If provided, this schema takes precedence over any arrow schema embedded
630    /// in the metadata (see the [`arrow`] documentation for more details).
631    ///
632    /// If the provided schema is not compatible with the data stored in the
633    /// parquet file schema, an error will be returned when constructing the
634    /// builder.
635    ///
636    /// This option is only required if you want to explicitly control the
637    /// conversion of Parquet types to Arrow types, such as casting a column to
638    /// a different type. For example, if you wanted to read an Int64 in
639    /// a Parquet file to a [`TimestampMicrosecondArray`] in the Arrow schema.
640    ///
641    /// [`arrow`]: crate::arrow
642    /// [`TimestampMicrosecondArray`]: arrow_array::TimestampMicrosecondArray
643    ///
644    /// # Notes
645    ///
646    /// The provided schema must have the same number of columns as the parquet schema and
647    /// the column names must be the same.
648    ///
649    /// # Example
650    /// ```
651    /// # use std::sync::Arc;
652    /// # use bytes::Bytes;
653    /// # use arrow_array::{ArrayRef, Int32Array, RecordBatch};
654    /// # use arrow_schema::{DataType, Field, Schema, TimeUnit};
655    /// # use parquet::arrow::arrow_reader::{ArrowReaderOptions, ParquetRecordBatchReaderBuilder};
656    /// # use parquet::arrow::ArrowWriter;
657    /// // Write data - schema is inferred from the data to be Int32
658    /// let mut file = Vec::new();
659    /// let batch = RecordBatch::try_from_iter(vec![
660    ///     ("col_1", Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef),
661    /// ]).unwrap();
662    /// let mut writer = ArrowWriter::try_new(&mut file, batch.schema(), None).unwrap();
663    /// writer.write(&batch).unwrap();
664    /// writer.close().unwrap();
665    /// let file = Bytes::from(file);
666    ///
667    /// // Read the file back.
668    /// // Supply a schema that interprets the Int32 column as a Timestamp.
669    /// let supplied_schema = Arc::new(Schema::new(vec![
670    ///     Field::new("col_1", DataType::Timestamp(TimeUnit::Nanosecond, None), false)
671    /// ]));
672    /// let options = ArrowReaderOptions::new().with_schema(supplied_schema.clone());
673    /// let mut builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
674    ///     file.clone(),
675    ///     options
676    /// ).expect("Error if the schema is not compatible with the parquet file schema.");
677    ///
678    /// // Create the reader and read the data using the supplied schema.
679    /// let mut reader = builder.build().unwrap();
680    /// let _batch = reader.next().unwrap().unwrap();
681    /// ```
682    ///
683    /// # Example: Preserving Dictionary Encoding
684    ///
685    /// By default, Parquet string columns are read as `Utf8Array` (or `LargeUtf8Array`),
686    /// even if the underlying Parquet data uses dictionary encoding. You can preserve
687    /// the dictionary encoding by specifying a `Dictionary` type in the schema hint:
688    ///
689    /// ```
690    /// # use std::sync::Arc;
691    /// # use bytes::Bytes;
692    /// # use arrow_array::{ArrayRef, RecordBatch, StringArray};
693    /// # use arrow_schema::{DataType, Field, Schema};
694    /// # use parquet::arrow::arrow_reader::{ArrowReaderOptions, ParquetRecordBatchReaderBuilder};
695    /// # use parquet::arrow::ArrowWriter;
696    /// // Write a Parquet file with string data
697    /// let mut file = Vec::new();
698    /// let schema = Arc::new(Schema::new(vec![
699    ///     Field::new("city", DataType::Utf8, false)
700    /// ]));
701    /// let cities = StringArray::from(vec!["Berlin", "Berlin", "Paris", "Berlin", "Paris"]);
702    /// let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(cities)]).unwrap();
703    ///
704    /// let mut writer = ArrowWriter::try_new(&mut file, batch.schema(), None).unwrap();
705    /// writer.write(&batch).unwrap();
706    /// writer.close().unwrap();
707    /// let file = Bytes::from(file);
708    ///
709    /// // Read the file back, requesting dictionary encoding preservation
710    /// let dict_schema = Arc::new(Schema::new(vec![
711    ///     Field::new("city", DataType::Dictionary(
712    ///         Box::new(DataType::Int32),
713    ///         Box::new(DataType::Utf8)
714    ///     ), false)
715    /// ]));
716    /// let options = ArrowReaderOptions::new().with_schema(dict_schema);
717    /// let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
718    ///     file.clone(),
719    ///     options
720    /// ).unwrap();
721    ///
722    /// let mut reader = builder.build().unwrap();
723    /// let batch = reader.next().unwrap().unwrap();
724    ///
725    /// // The column is now a DictionaryArray
726    /// assert!(matches!(
727    ///     batch.column(0).data_type(),
728    ///     DataType::Dictionary(_, _)
729    /// ));
730    /// ```
731    ///
732    /// **Note**: Dictionary encoding preservation works best when:
733    /// 1. The original column was dictionary encoded (the default for string columns)
734    /// 2. There are a small number of distinct values
735    pub fn with_schema(self, schema: SchemaRef) -> Self {
736        Self {
737            supplied_schema: Some(schema),
738            skip_arrow_metadata: true,
739            ..self
740        }
741    }
742
743    /// Sets the [`PageIndexPolicy`] for both the column and offset indexes.
744    ///
745    /// The `PageIndex` can be used to push down predicates to the parquet scan,
746    /// potentially eliminating unnecessary IO, by some query engines.
747    /// The `PageIndex` consists of two structures: the `ColumnIndex` and `OffsetIndex`.
748    /// This method sets the same policy for both. For fine-grained control, use
749    /// [`Self::with_column_index_policy`] and [`Self::with_offset_index_policy`].
750    pub fn with_page_index_policy(self, policy: PageIndexPolicy) -> Self {
751        self.with_column_index_policy(policy)
752            .with_offset_index_policy(policy)
753    }
754
755    /// Sets the [`PageIndexPolicy`] for the Parquet [ColumnIndex] structure.
756    ///
757    /// The `ColumnIndex` contains min/max statistics for each page, which can be used
758    /// for predicate pushdown and page-level pruning.
759    ///
760    /// [ColumnIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md
761    pub fn with_column_index_policy(mut self, policy: PageIndexPolicy) -> Self {
762        self.column_index = policy;
763        self
764    }
765
766    /// Sets the [`PageIndexPolicy`] for the Parquet [OffsetIndex] structure.
767    ///
768    /// The `OffsetIndex` contains the locations and sizes of each page, which enables
769    /// efficient page-level skipping and random access within column chunks.
770    ///
771    /// [OffsetIndex]: https://github.com/apache/parquet-format/blob/master/PageIndex.md
772    pub fn with_offset_index_policy(mut self, policy: PageIndexPolicy) -> Self {
773        self.offset_index = policy;
774        self
775    }
776
777    /// Provide a Parquet schema to use when decoding the metadata. The schema in the Parquet
778    /// footer will be skipped.
779    ///
780    /// This can be used to avoid reparsing the schema from the file when it is
781    /// already known.
782    pub fn with_parquet_schema(mut self, schema: Arc<SchemaDescriptor>) -> Self {
783        self.metadata_options.set_schema(schema);
784        self
785    }
786
787    /// Set whether to convert the [`encoding_stats`] in the Parquet `ColumnMetaData` to a bitmask
788    /// (defaults to `false`).
789    ///
790    /// See [`ColumnChunkMetaData::page_encoding_stats_mask`] for an explanation of why this
791    /// might be desirable.
792    ///
793    /// [`ColumnChunkMetaData::page_encoding_stats_mask`]:
794    /// crate::file::metadata::ColumnChunkMetaData::page_encoding_stats_mask
795    /// [`encoding_stats`]:
796    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L917
797    pub fn with_encoding_stats_as_mask(mut self, val: bool) -> Self {
798        self.metadata_options.set_encoding_stats_as_mask(val);
799        self
800    }
801
802    /// Sets the decoding policy for [`encoding_stats`] in the Parquet `ColumnMetaData`.
803    ///
804    /// [`encoding_stats`]:
805    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L917
806    pub fn with_encoding_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
807        self.metadata_options.set_encoding_stats_policy(policy);
808        self
809    }
810
811    /// Sets the decoding policy for [`statistics`] in the Parquet `ColumnMetaData`.
812    ///
813    /// [`statistics`]:
814    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L912
815    pub fn with_column_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
816        self.metadata_options.set_column_stats_policy(policy);
817        self
818    }
819
820    /// Sets the decoding policy for [`size_statistics`] in the Parquet `ColumnMetaData`.
821    ///
822    /// [`size_statistics`]:
823    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L936
824    pub fn with_size_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
825        self.metadata_options.set_size_stats_policy(policy);
826        self
827    }
828
829    /// Provide the file decryption properties to use when reading encrypted parquet files.
830    ///
831    /// If encryption is enabled and the file is encrypted, the `file_decryption_properties` must be provided.
832    #[cfg(feature = "encryption")]
833    pub fn with_file_decryption_properties(
834        self,
835        file_decryption_properties: Arc<FileDecryptionProperties>,
836    ) -> Self {
837        Self {
838            file_decryption_properties: Some(file_decryption_properties),
839            ..self
840        }
841    }
842
843    /// Include virtual columns in the output.
844    ///
845    /// Virtual columns are columns that are not part of the Parquet schema, but are added to the output by the reader such as row numbers and row group indices.
846    ///
847    /// # Example
848    /// ```
849    /// # use std::sync::Arc;
850    /// # use bytes::Bytes;
851    /// # use arrow_array::{ArrayRef, Int64Array, RecordBatch};
852    /// # use arrow_schema::{DataType, Field, Schema};
853    /// # use parquet::arrow::{ArrowWriter, RowNumber};
854    /// # use parquet::arrow::arrow_reader::{ArrowReaderOptions, ParquetRecordBatchReaderBuilder};
855    /// #
856    /// # fn main() -> Result<(), Box<dyn std::error::Error>> {
857    /// // Create a simple record batch with some data
858    /// let values = Arc::new(Int64Array::from(vec![1, 2, 3])) as ArrayRef;
859    /// let batch = RecordBatch::try_from_iter(vec![("value", values)])?;
860    ///
861    /// // Write the batch to an in-memory buffer
862    /// let mut file = Vec::new();
863    /// let mut writer = ArrowWriter::try_new(
864    ///     &mut file,
865    ///     batch.schema(),
866    ///     None
867    /// )?;
868    /// writer.write(&batch)?;
869    /// writer.close()?;
870    /// let file = Bytes::from(file);
871    ///
872    /// // Create a virtual column for row numbers
873    /// let row_number_field = Arc::new(Field::new("row_number", DataType::Int64, false)
874    ///     .with_extension_type(RowNumber));
875    ///
876    /// // Configure options with virtual columns
877    /// let options = ArrowReaderOptions::new()
878    ///     .with_virtual_columns(vec![row_number_field])?;
879    ///
880    /// // Create a reader with the options
881    /// let mut reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
882    ///     file,
883    ///     options
884    /// )?
885    /// .build()?;
886    ///
887    /// // Read the batch - it will include both the original column and the virtual row_number column
888    /// let result_batch = reader.next().unwrap()?;
889    /// assert_eq!(result_batch.num_columns(), 2); // "value" + "row_number"
890    /// assert_eq!(result_batch.num_rows(), 3);
891    /// #
892    /// # Ok(())
893    /// # }
894    /// ```
895    pub fn with_virtual_columns(self, virtual_columns: Vec<FieldRef>) -> Result<Self> {
896        // Validate that all fields are virtual columns
897        for field in &virtual_columns {
898            if !is_virtual_column(field) {
899                return Err(ParquetError::General(format!(
900                    "Field '{}' is not a virtual column. Virtual columns must have extension type names starting with 'arrow.virtual.'",
901                    field.name()
902                )));
903            }
904        }
905        Ok(Self {
906            virtual_columns,
907            ..self
908        })
909    }
910
911    /// Retrieve the currently set [`PageIndexPolicy`] for the offset index.
912    ///
913    /// This can be set via [`with_offset_index_policy`][Self::with_offset_index_policy]
914    /// or [`with_page_index_policy`][Self::with_page_index_policy].
915    pub fn offset_index_policy(&self) -> PageIndexPolicy {
916        self.offset_index
917    }
918
919    /// Retrieve the currently set [`PageIndexPolicy`] for the column index.
920    ///
921    /// This can be set via [`with_column_index_policy`][Self::with_column_index_policy]
922    /// or [`with_page_index_policy`][Self::with_page_index_policy].
923    pub fn column_index_policy(&self) -> PageIndexPolicy {
924        self.column_index
925    }
926
927    /// Retrieve the currently set metadata decoding options.
928    pub fn metadata_options(&self) -> &ParquetMetaDataOptions {
929        &self.metadata_options
930    }
931
932    /// Retrieve the currently set file decryption properties.
933    ///
934    /// This can be set via
935    /// [`file_decryption_properties`][Self::with_file_decryption_properties].
936    #[cfg(feature = "encryption")]
937    pub fn file_decryption_properties(&self) -> Option<&Arc<FileDecryptionProperties>> {
938        self.file_decryption_properties.as_ref()
939    }
940}
941
942impl ParquetMetaDataReader {
943    /// Applies the metadata related settings from [`ArrowReaderOptions`],
944    /// such as the [`ParquetMetaDataOptions`], decryption properties, and
945    /// [`PageIndexPolicy`] to this reader.
946    ///
947    /// The page index policies are only applied if at least one of them is not
948    /// [`PageIndexPolicy::Skip`], so policies previously configured on this
949    /// reader (e.g. from a preload setting) are preserved when the options do
950    /// not request the page index.
951    ///
952    /// This encodes the canonical way to construct a `ParquetMetaDataReader`
953    /// inside `AsyncFileReader::get_metadata` (available with the `async`
954    /// feature), so implementations outside this crate do not need to
955    /// duplicate it.
956    pub fn with_arrow_reader_options(mut self, options: Option<&ArrowReaderOptions>) -> Self {
957        let Some(options) = options else { return self };
958
959        self = self.with_metadata_options(Some(options.metadata_options().clone()));
960
961        #[cfg(feature = "encryption")]
962        {
963            self = self.with_decryption_properties(
964                options.file_decryption_properties.as_ref().map(Arc::clone),
965            );
966        }
967
968        if options.column_index_policy() != PageIndexPolicy::Skip
969            || options.offset_index_policy() != PageIndexPolicy::Skip
970        {
971            self = self
972                .with_column_index_policy(options.column_index_policy())
973                .with_offset_index_policy(options.offset_index_policy());
974        }
975
976        self
977    }
978}
979
980/// The metadata necessary to construct a [`ArrowReaderBuilder`]
981///
982/// Note this structure is cheaply clone-able as it consists of several arcs.
983///
984/// This structure allows
985///
986/// 1. Loading metadata for a file once and then using that same metadata to
987///    construct multiple separate readers, for example, to distribute readers
988///    across multiple threads
989///
990/// 2. Using a cached copy of the [`ParquetMetadata`] rather than reading it
991///    from the file each time a reader is constructed.
992///
993/// [`ParquetMetadata`]: crate::file::metadata::ParquetMetaData
994#[derive(Debug, Clone)]
995pub struct ArrowReaderMetadata {
996    /// The Parquet Metadata, if known aprior
997    pub(crate) metadata: Arc<ParquetMetaData>,
998    /// The Arrow Schema
999    pub(crate) schema: SchemaRef,
1000    /// The Parquet schema (root field)
1001    pub(crate) fields: Option<Arc<ParquetField>>,
1002}
1003
1004impl ArrowReaderMetadata {
1005    /// Create [`ArrowReaderMetadata`] from the provided [`ArrowReaderOptions`]
1006    /// and [`ChunkReader`]
1007    ///
1008    /// See [`ParquetRecordBatchReaderBuilder::new_with_metadata`] for an
1009    /// example of how this can be used
1010    ///
1011    /// # Notes
1012    ///
1013    /// If `options` indicates the page index should be read, but
1014    /// `Self::metadata` is missing the page index, this function will attempt
1015    /// to load the page index by making an object store request.
1016    ///
1017    /// See [`ArrowReaderOptions::with_page_index_policy`] for more information on the page index.
1018    pub fn load<T: ChunkReader>(reader: &T, options: ArrowReaderOptions) -> Result<Self> {
1019        let metadata = ParquetMetaDataReader::new()
1020            .with_column_index_policy(options.column_index)
1021            .with_offset_index_policy(options.offset_index)
1022            .with_metadata_options(Some(options.metadata_options.clone()));
1023        #[cfg(feature = "encryption")]
1024        let metadata = metadata.with_decryption_properties(
1025            options.file_decryption_properties.as_ref().map(Arc::clone),
1026        );
1027        let metadata = metadata.parse_and_finish(reader)?;
1028        Self::try_new(Arc::new(metadata), options)
1029    }
1030
1031    /// Create a new [`ArrowReaderMetadata`] from a pre-existing
1032    /// [`ParquetMetaData`] and [`ArrowReaderOptions`].
1033    ///
1034    /// # Notes
1035    ///
1036    /// This function will not attempt to load the PageIndex if not present in the metadata, regardless
1037    /// of the settings in `options`. See [`Self::load`] to load metadata including the page index if needed.
1038    pub fn try_new(metadata: Arc<ParquetMetaData>, options: ArrowReaderOptions) -> Result<Self> {
1039        match options.supplied_schema {
1040            Some(supplied_schema) => Self::with_supplied_schema(
1041                metadata,
1042                supplied_schema.clone(),
1043                &options.virtual_columns,
1044            ),
1045            None => {
1046                let kv_metadata = match options.skip_arrow_metadata {
1047                    true => None,
1048                    false => metadata.file_metadata().key_value_metadata(),
1049                };
1050
1051                let (schema, fields) = parquet_to_arrow_schema_and_fields(
1052                    metadata.file_metadata().schema_descr(),
1053                    ProjectionMask::all(),
1054                    kv_metadata,
1055                    &options.virtual_columns,
1056                )?;
1057
1058                Ok(Self {
1059                    metadata,
1060                    schema: Arc::new(schema),
1061                    fields: fields.map(Arc::new),
1062                })
1063            }
1064        }
1065    }
1066
1067    fn with_supplied_schema(
1068        metadata: Arc<ParquetMetaData>,
1069        supplied_schema: SchemaRef,
1070        virtual_columns: &[FieldRef],
1071    ) -> Result<Self> {
1072        let parquet_schema = metadata.file_metadata().schema_descr();
1073        let field_levels = parquet_to_arrow_field_levels_with_virtual(
1074            parquet_schema,
1075            ProjectionMask::all(),
1076            Some(supplied_schema.fields()),
1077            virtual_columns,
1078        )?;
1079        let fields = field_levels.fields;
1080        let inferred_len = fields.len();
1081        let supplied_len = supplied_schema.fields().len() + virtual_columns.len();
1082        // Ensure the supplied schema has the same number of columns as the parquet schema.
1083        // parquet_to_arrow_field_levels is expected to throw an error if the schemas have
1084        // different lengths, but we check here to be safe.
1085        if inferred_len != supplied_len {
1086            return Err(arrow_err!(format!(
1087                "Incompatible supplied Arrow schema: expected {} columns received {}",
1088                inferred_len, supplied_len
1089            )));
1090        }
1091
1092        let mut errors = Vec::new();
1093
1094        let field_iter = supplied_schema.fields().iter().zip(fields.iter());
1095
1096        for (field1, field2) in field_iter {
1097            if field1.data_type() != field2.data_type() {
1098                errors.push(format!(
1099                    "data type mismatch for field {}: requested {} but found {}",
1100                    field1.name(),
1101                    field1.data_type(),
1102                    field2.data_type()
1103                ));
1104            }
1105            if field1.is_nullable() != field2.is_nullable() {
1106                errors.push(format!(
1107                    "nullability mismatch for field {}: expected {:?} but found {:?}",
1108                    field1.name(),
1109                    field1.is_nullable(),
1110                    field2.is_nullable()
1111                ));
1112            }
1113            if field1.metadata() != field2.metadata() {
1114                errors.push(format!(
1115                    "metadata mismatch for field {}: expected {:?} but found {:?}",
1116                    field1.name(),
1117                    field1.metadata(),
1118                    field2.metadata()
1119                ));
1120            }
1121        }
1122
1123        if !errors.is_empty() {
1124            let message = errors.join(", ");
1125            return Err(ParquetError::ArrowError(format!(
1126                "Incompatible supplied Arrow schema: {message}",
1127            )));
1128        }
1129
1130        Ok(Self {
1131            metadata,
1132            schema: supplied_schema,
1133            fields: field_levels.levels.map(Arc::new),
1134        })
1135    }
1136
1137    /// Returns a reference to the [`ParquetMetaData`] for this parquet file
1138    pub fn metadata(&self) -> &Arc<ParquetMetaData> {
1139        &self.metadata
1140    }
1141
1142    /// Returns the parquet [`SchemaDescriptor`] for this parquet file
1143    pub fn parquet_schema(&self) -> &SchemaDescriptor {
1144        self.metadata.file_metadata().schema_descr()
1145    }
1146
1147    /// Returns the arrow [`SchemaRef`] for this parquet file
1148    pub fn schema(&self) -> &SchemaRef {
1149        &self.schema
1150    }
1151}
1152
1153#[doc(hidden)]
1154// A newtype used within `ReaderOptionsBuilder` to distinguish sync readers from async
1155pub struct SyncReader<T: ChunkReader>(T);
1156
1157impl<T: Debug + ChunkReader> Debug for SyncReader<T> {
1158    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1159        f.debug_tuple("SyncReader").field(&self.0).finish()
1160    }
1161}
1162
1163/// Creates [`ParquetRecordBatchReader`] for reading Parquet files into Arrow [`RecordBatch`]es
1164///
1165/// # See Also
1166/// * [`crate::arrow::async_reader::ParquetRecordBatchStreamBuilder`] for an async API
1167/// * [`crate::arrow::push_decoder::ParquetPushDecoderBuilder`] for a SansIO decoder API
1168/// * [`ArrowReaderBuilder`] for additional member functions
1169pub type ParquetRecordBatchReaderBuilder<T> = ArrowReaderBuilder<SyncReader<T>>;
1170
1171impl<T: ChunkReader + 'static> ParquetRecordBatchReaderBuilder<T> {
1172    /// Create a new [`ParquetRecordBatchReaderBuilder`]
1173    ///
1174    /// ```
1175    /// # use std::sync::Arc;
1176    /// # use bytes::Bytes;
1177    /// # use arrow_array::{Int32Array, RecordBatch};
1178    /// # use arrow_schema::{DataType, Field, Schema};
1179    /// # use parquet::arrow::arrow_reader::{ParquetRecordBatchReader, ParquetRecordBatchReaderBuilder};
1180    /// # use parquet::arrow::ArrowWriter;
1181    /// # let mut file: Vec<u8> = Vec::with_capacity(1024);
1182    /// # let schema = Arc::new(Schema::new(vec![Field::new("i32", DataType::Int32, false)]));
1183    /// # let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1184    /// # let batch = RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1, 2, 3]))]).unwrap();
1185    /// # writer.write(&batch).unwrap();
1186    /// # writer.close().unwrap();
1187    /// # let file = Bytes::from(file);
1188    /// // Build the reader from anything that implements `ChunkReader`
1189    /// // such as a `File`, or `Bytes`
1190    /// let mut builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
1191    /// // The builder has access to ParquetMetaData such
1192    /// // as the number and layout of row groups
1193    /// assert_eq!(builder.metadata().num_row_groups(), 1);
1194    /// // Call build to create the reader
1195    /// let mut reader: ParquetRecordBatchReader = builder.build().unwrap();
1196    /// // Read data
1197    /// while let Some(batch) = reader.next().transpose()? {
1198    ///     println!("Read {} rows", batch.num_rows());
1199    /// }
1200    /// # Ok::<(), parquet::errors::ParquetError>(())
1201    /// ```
1202    pub fn try_new(reader: T) -> Result<Self> {
1203        Self::try_new_with_options(reader, Default::default())
1204    }
1205
1206    /// Create a new [`ParquetRecordBatchReaderBuilder`] with [`ArrowReaderOptions`]
1207    ///
1208    /// Use this method if you want to control the options for reading the
1209    /// [`ParquetMetaData`]
1210    pub fn try_new_with_options(reader: T, options: ArrowReaderOptions) -> Result<Self> {
1211        let metadata = ArrowReaderMetadata::load(&reader, options)?;
1212        Ok(Self::new_with_metadata(reader, metadata))
1213    }
1214
1215    /// Create a [`ParquetRecordBatchReaderBuilder`] from the provided [`ArrowReaderMetadata`]
1216    ///
1217    /// Use this method if you already have [`ParquetMetaData`] for a file.
1218    /// This interface allows:
1219    ///
1220    /// 1. Loading metadata once and using it to create multiple builders with
1221    ///    potentially different settings or run on different threads
1222    ///
1223    /// 2. Using a cached copy of the metadata rather than re-reading it from the
1224    ///    file each time a reader is constructed.
1225    ///
1226    /// See the docs on [`ArrowReaderMetadata`] for more details
1227    ///
1228    /// # Example
1229    /// ```
1230    /// # use std::fs::metadata;
1231    /// # use std::sync::Arc;
1232    /// # use bytes::Bytes;
1233    /// # use arrow_array::{Int32Array, RecordBatch};
1234    /// # use arrow_schema::{DataType, Field, Schema};
1235    /// # use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ParquetRecordBatchReader, ParquetRecordBatchReaderBuilder};
1236    /// # use parquet::arrow::ArrowWriter;
1237    /// #
1238    /// # let mut file: Vec<u8> = Vec::with_capacity(1024);
1239    /// # let schema = Arc::new(Schema::new(vec![Field::new("i32", DataType::Int32, false)]));
1240    /// # let mut writer = ArrowWriter::try_new(&mut file, schema.clone(), None).unwrap();
1241    /// # let batch = RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1, 2, 3]))]).unwrap();
1242    /// # writer.write(&batch).unwrap();
1243    /// # writer.close().unwrap();
1244    /// # let file = Bytes::from(file);
1245    /// #
1246    /// let metadata = ArrowReaderMetadata::load(&file, Default::default()).unwrap();
1247    /// let mut a = ParquetRecordBatchReaderBuilder::new_with_metadata(file.clone(), metadata.clone()).build().unwrap();
1248    /// let mut b = ParquetRecordBatchReaderBuilder::new_with_metadata(file, metadata).build().unwrap();
1249    ///
1250    /// // Should be able to read from both in parallel
1251    /// assert_eq!(a.next().unwrap().unwrap(), b.next().unwrap().unwrap());
1252    /// ```
1253    pub fn new_with_metadata(input: T, metadata: ArrowReaderMetadata) -> Self {
1254        Self::new_builder(SyncReader(input), metadata)
1255    }
1256
1257    /// Read bloom filter for a column in a row group
1258    ///
1259    /// Returns `None` if the column does not have a bloom filter
1260    ///
1261    /// We should call this function after other forms pruning, such as projection and predicate pushdown.
1262    pub fn get_row_group_column_bloom_filter(
1263        &self,
1264        row_group_idx: usize,
1265        column_idx: usize,
1266    ) -> Result<Option<Sbbf>> {
1267        let metadata = self.metadata.row_group(row_group_idx);
1268        let column_metadata = metadata.column(column_idx);
1269
1270        let offset: u64 = if let Some(offset) = column_metadata.bloom_filter_offset() {
1271            offset
1272                .try_into()
1273                .map_err(|_| ParquetError::General("Bloom filter offset is invalid".to_string()))?
1274        } else {
1275            return Ok(None);
1276        };
1277
1278        let buffer = match column_metadata.bloom_filter_length() {
1279            Some(length) => self.input.0.get_bytes(offset, length as usize),
1280            None => self.input.0.get_bytes(offset, SBBF_HEADER_SIZE_ESTIMATE),
1281        }?;
1282
1283        let (header, bitset_offset) =
1284            chunk_read_bloom_filter_header_and_offset(offset, buffer.clone())?;
1285
1286        match header.algorithm {
1287            BloomFilterAlgorithm::BLOCK => {
1288                // this match exists to future proof the singleton algorithm enum
1289            }
1290        }
1291        match header.compression {
1292            BloomFilterCompression::UNCOMPRESSED => {
1293                // this match exists to future proof the singleton compression enum
1294            }
1295        }
1296        match header.hash {
1297            BloomFilterHash::XXHASH => {
1298                // this match exists to future proof the singleton hash enum
1299            }
1300        }
1301
1302        let bitset = match column_metadata.bloom_filter_length() {
1303            Some(_) => {
1304                let bitset_start = bitset_offset
1305                    .checked_sub(offset)
1306                    .and_then(|start| usize::try_from(start).ok())
1307                    .ok_or_else(|| {
1308                        ParquetError::General("Bloom filter offset is invalid".to_string())
1309                    })?;
1310                buffer.slice(bitset_start..)
1311            }
1312            None => {
1313                let bitset_length: usize = header.num_bytes.try_into().map_err(|_| {
1314                    ParquetError::General("Bloom filter length is invalid".to_string())
1315                })?;
1316                self.input.0.get_bytes(bitset_offset, bitset_length)?
1317            }
1318        };
1319        Ok(Some(Sbbf::new(&bitset)))
1320    }
1321
1322    /// Build a [`ParquetRecordBatchReader`]
1323    ///
1324    /// Note: this will eagerly evaluate any `RowFilter` before returning
1325    pub fn build(self) -> Result<ParquetRecordBatchReader> {
1326        let Self {
1327            input,
1328            metadata,
1329            schema: _,
1330            fields,
1331            batch_size,
1332            row_group_plan,
1333            projection,
1334            mut filter,
1335            row_selection_policy,
1336            limit,
1337            offset,
1338            metrics,
1339            // Not used for the sync reader, see https://github.com/apache/arrow-rs/issues/8000
1340            max_predicate_cache_size: _,
1341        } = self;
1342
1343        // Try to avoid allocate large buffer
1344        let batch_size = batch_size.min(metadata.file_metadata().num_rows() as usize);
1345
1346        let (row_groups, selection) = row_group_plan.into_global()?;
1347
1348        let row_groups = row_groups.unwrap_or_else(|| (0..metadata.num_row_groups()).collect());
1349
1350        let reader = ReaderRowGroups {
1351            reader: Arc::new(input.0),
1352            metadata,
1353            row_groups,
1354        };
1355
1356        let mut plan_builder = ReadPlanBuilder::new(batch_size)
1357            .with_selection(selection)
1358            .with_row_selection_policy(row_selection_policy);
1359
1360        // Update selection based on any filters
1361        if let Some(filter) = filter.as_mut() {
1362            for predicate in &mut filter.predicates {
1363                // break early if we have ruled out all rows
1364                if !plan_builder.selects_any() {
1365                    break;
1366                }
1367
1368                let mut cache_projection = predicate.projection().clone();
1369                cache_projection.intersect(&projection);
1370
1371                let array_reader = ArrayReaderBuilder::new(&reader, &metrics)
1372                    .with_batch_size(batch_size)
1373                    .with_parquet_metadata(&reader.metadata)
1374                    .build_array_reader(fields.as_deref(), predicate.projection())?;
1375
1376                plan_builder = plan_builder.with_predicate(array_reader, predicate.as_mut())?;
1377            }
1378        }
1379
1380        let array_reader = ArrayReaderBuilder::new(&reader, &metrics)
1381            .with_batch_size(batch_size)
1382            .with_parquet_metadata(&reader.metadata)
1383            .build_array_reader(fields.as_deref(), &projection)?;
1384
1385        let read_plan = plan_builder
1386            .limited(reader.num_rows())
1387            .with_offset(offset)
1388            .with_limit(limit)
1389            .build_limited()
1390            .build();
1391
1392        Ok(ParquetRecordBatchReader::new(array_reader, read_plan))
1393    }
1394}
1395
1396struct ReaderRowGroups<T: ChunkReader> {
1397    reader: Arc<T>,
1398
1399    metadata: Arc<ParquetMetaData>,
1400    /// Optional list of row group indices to scan
1401    row_groups: Vec<usize>,
1402}
1403
1404impl<T: ChunkReader + 'static> RowGroups for ReaderRowGroups<T> {
1405    fn num_rows(&self) -> usize {
1406        let meta = self.metadata.row_groups();
1407        self.row_groups
1408            .iter()
1409            .map(|x| meta[*x].num_rows() as usize)
1410            .sum()
1411    }
1412
1413    fn column_chunks(&self, i: usize) -> Result<Box<dyn PageIterator>> {
1414        Ok(Box::new(ReaderPageIterator {
1415            column_idx: i,
1416            reader: self.reader.clone(),
1417            metadata: self.metadata.clone(),
1418            row_groups: self.row_groups.clone().into_iter(),
1419        }))
1420    }
1421
1422    fn row_groups(&self) -> Box<dyn Iterator<Item = &RowGroupMetaData> + '_> {
1423        Box::new(
1424            self.row_groups
1425                .iter()
1426                .map(move |i| self.metadata.row_group(*i)),
1427        )
1428    }
1429
1430    fn metadata(&self) -> &ParquetMetaData {
1431        self.metadata.as_ref()
1432    }
1433}
1434
1435struct ReaderPageIterator<T: ChunkReader> {
1436    reader: Arc<T>,
1437    column_idx: usize,
1438    row_groups: std::vec::IntoIter<usize>,
1439    metadata: Arc<ParquetMetaData>,
1440}
1441
1442impl<T: ChunkReader + 'static> ReaderPageIterator<T> {
1443    /// Return the next SerializedPageReader
1444    fn next_page_reader(&self, rg_idx: usize) -> Result<SerializedPageReader<T>> {
1445        let rg = self.metadata.row_group(rg_idx);
1446        let column_chunk_metadata = rg.column(self.column_idx);
1447        let page_locations = self
1448            .metadata
1449            .page_index()
1450            .map(|i| i.page_locations(rg_idx, self.column_idx).cloned())
1451            .unwrap_or(None);
1452        let total_rows = rg.num_rows() as usize;
1453        let reader = self.reader.clone();
1454
1455        SerializedPageReader::new(reader, column_chunk_metadata, total_rows, page_locations)?
1456            .add_crypto_context(
1457                rg_idx,
1458                self.column_idx,
1459                self.metadata.as_ref(),
1460                column_chunk_metadata,
1461            )
1462    }
1463}
1464
1465impl<T: ChunkReader + 'static> Iterator for ReaderPageIterator<T> {
1466    type Item = Result<Box<dyn PageReader>>;
1467
1468    fn next(&mut self) -> Option<Self::Item> {
1469        let rg_idx = self.row_groups.next()?;
1470        let page_reader = self
1471            .next_page_reader(rg_idx)
1472            .map(|page_reader| Box::new(page_reader) as _);
1473        Some(page_reader)
1474    }
1475}
1476
1477impl<T: ChunkReader + 'static> PageIterator for ReaderPageIterator<T> {}
1478
1479/// Reads Parquet data as Arrow [`RecordBatch`]es
1480///
1481/// This struct implements the [`RecordBatchReader`] trait and is an
1482/// `Iterator<Item = ArrowResult<RecordBatch>>` that yields [`RecordBatch`]es.
1483///
1484/// Typically, either reads from a file or an in memory buffer [`Bytes`]
1485///
1486/// Created by [`ParquetRecordBatchReaderBuilder`]
1487///
1488/// [`Bytes`]: bytes::Bytes
1489pub struct ParquetRecordBatchReader {
1490    array_reader: Box<dyn ArrayReader>,
1491    schema: SchemaRef,
1492    read_plan: ReadPlan,
1493}
1494
1495/// Accumulates filter masks for decoded chunks in one logical output batch.
1496///
1497/// The first chunk keeps its [`BooleanBuffer`] without copying. A second chunk
1498/// promotes the accumulator to a [`BooleanBufferBuilder`], and later chunks are
1499/// appended to it. For example, chunks `1001` and `1` become `10011`:
1500///
1501/// ```text
1502///           append(1001)             append(1)
1503///   Empty ───────────────▶ Single ───────────────▶ Combined
1504///                          1001                    10011
1505///                          (zero copy)             (promoted to builder)
1506/// ```
1507///
1508/// The combined mask lines up with the rows that [`ArrayReader::read_records`]
1509/// buffered across the decoded chunks (see [`read_mask_batch`]). Consuming the
1510/// buffered batch and filtering it with the accumulated mask yields the output.
1511///
1512/// ```text
1513///   decoded rows:   0 1 2 3   11   <-- buffered by the array reader
1514///   chunk masks:   [1 0 0 1] [1]
1515///   finish():       1 0 0 1   1    <-- filters the whole batch in one pass
1516/// ```
1517#[derive(Default)]
1518enum FilterMaskAccumulator {
1519    #[default]
1520    Empty,
1521    Single(BooleanBuffer),
1522    Combined(BooleanBufferBuilder),
1523}
1524
1525impl FilterMaskAccumulator {
1526    fn append(&mut self, mask: BooleanBuffer) {
1527        *self = match std::mem::take(self) {
1528            Self::Empty => Self::Single(mask),
1529            Self::Single(first) => {
1530                let mut combined = BooleanBufferBuilder::new(first.len() + mask.len());
1531                combined.append_buffer(&first);
1532                combined.append_buffer(&mask);
1533                Self::Combined(combined)
1534            }
1535            Self::Combined(mut combined) => {
1536                combined.append_buffer(&mask);
1537                Self::Combined(combined)
1538            }
1539        };
1540    }
1541
1542    fn finish(self) -> Option<BooleanBuffer> {
1543        match self {
1544            Self::Empty => None,
1545            Self::Single(mask) => Some(mask),
1546            Self::Combined(combined) => Some(combined.build()),
1547        }
1548    }
1549}
1550
1551/// Converts the projection buffered by `array_reader` into a record batch.
1552fn consume_record_batch(array_reader: &mut dyn ArrayReader) -> Result<RecordBatch> {
1553    let array = array_reader.consume_batch()?;
1554    let struct_array = array.as_struct_opt().ok_or_else(|| {
1555        ArrowError::ParquetError("Struct array reader should return struct array".to_string())
1556    })?;
1557    Ok(RecordBatch::from(struct_array))
1558}
1559
1560/// Reads one logical Mask batch, potentially spanning multiple loaded ranges.
1561///
1562/// Each [`MaskCursor`] chunk is safe to decode because it stays within loaded
1563/// pages. Gaps are crossed with [`ArrayReader::skip_records`], while decoded
1564/// arrays and their mask fragments remain buffered. Once `batch_size` selected
1565/// rows have accumulated, this consumes the underlying batch and filters it
1566/// once with the combined mask.
1567fn read_mask_batch(
1568    array_reader: &mut dyn ArrayReader,
1569    mask_cursor: &mut MaskCursor,
1570    batch_size: usize,
1571) -> Result<Option<RecordBatch>> {
1572    let mut selected_rows = 0;
1573    let mut filter_mask = FilterMaskAccumulator::default();
1574
1575    while selected_rows < batch_size && !mask_cursor.is_empty() {
1576        let mask_chunk = mask_cursor.next_chunk(batch_size - selected_rows)?;
1577
1578        if mask_chunk.initial_skip > 0 {
1579            let skipped = array_reader.skip_records(mask_chunk.initial_skip)?;
1580            if skipped != mask_chunk.initial_skip {
1581                return Err(general_err!(
1582                    "failed to skip rows, expected {}, got {}",
1583                    mask_chunk.initial_skip,
1584                    skipped
1585                ));
1586            }
1587        }
1588
1589        let mask = mask_cursor.mask_values_for(&mask_chunk)?;
1590        let read = array_reader.read_records(mask_chunk.chunk_rows)?;
1591        if read == 0 {
1592            return Err(general_err!(
1593                "reached end of column while expecting {} rows",
1594                mask_chunk.chunk_rows
1595            ));
1596        }
1597        if read != mask_chunk.chunk_rows {
1598            return Err(general_err!(
1599                "insufficient rows read from array reader - expected {}, got {}",
1600                mask_chunk.chunk_rows,
1601                read
1602            ));
1603        }
1604
1605        filter_mask.append(mask.values().clone());
1606        selected_rows += mask_chunk.selected_rows;
1607    }
1608
1609    if selected_rows == 0 {
1610        return Ok(None);
1611    }
1612
1613    let filter_mask = filter_mask
1614        .finish()
1615        .ok_or_else(|| general_err!("Internal Error: decoded Mask batch has no filter values"))?;
1616    let batch = consume_record_batch(array_reader)?;
1617    let filtered_batch = filter_record_batch(&batch, &BooleanArray::from(filter_mask))?;
1618    if filtered_batch.num_rows() != selected_rows {
1619        return Err(general_err!(
1620            "filtered rows mismatch selection - expected {}, got {}",
1621            selected_rows,
1622            filtered_batch.num_rows()
1623        ));
1624    }
1625
1626    Ok(Some(filtered_batch))
1627}
1628
1629impl Debug for ParquetRecordBatchReader {
1630    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1631        f.debug_struct("ParquetRecordBatchReader")
1632            .field("array_reader", &"...")
1633            .field("schema", &self.schema)
1634            .field("read_plan", &self.read_plan)
1635            .finish()
1636    }
1637}
1638
1639impl Iterator for ParquetRecordBatchReader {
1640    type Item = Result<RecordBatch, ArrowError>;
1641
1642    fn next(&mut self) -> Option<Self::Item> {
1643        self.next_inner()
1644            .map_err(|arrow_err| arrow_err.into())
1645            .transpose()
1646    }
1647}
1648
1649impl ParquetRecordBatchReader {
1650    /// Returns the next `RecordBatch` from the reader, or `None` if the reader
1651    /// has reached the end of the file.
1652    ///
1653    /// Returns `Result<Option<..>>` rather than `Option<Result<..>>` to
1654    /// simplify error handling with `?`
1655    fn next_inner(&mut self) -> Result<Option<RecordBatch>> {
1656        let mut read_records = 0;
1657        let batch_size = self.batch_size();
1658        if batch_size == 0 {
1659            return Ok(None);
1660        }
1661        match self.read_plan.row_selection_cursor_mut() {
1662            RowSelectionCursor::Mask(mask_cursor) => {
1663                return read_mask_batch(self.array_reader.as_mut(), mask_cursor, batch_size);
1664            }
1665            RowSelectionCursor::Selectors(selectors_cursor) => {
1666                while read_records < batch_size && !selectors_cursor.is_empty() {
1667                    let front = selectors_cursor.next_selector();
1668                    if front.skip {
1669                        let skipped = self.array_reader.skip_records(front.row_count)?;
1670
1671                        if skipped != front.row_count {
1672                            return Err(general_err!(
1673                                "failed to skip rows, expected {}, got {}",
1674                                front.row_count,
1675                                skipped
1676                            ));
1677                        }
1678                        continue;
1679                    }
1680
1681                    //Currently, when RowSelectors with row_count = 0 are included then its interpreted as end of reader.
1682                    //Fix is to skip such entries. See https://github.com/apache/arrow-rs/issues/2669
1683                    if front.row_count == 0 {
1684                        continue;
1685                    }
1686
1687                    // try to read record
1688                    let need_read = batch_size - read_records;
1689                    let to_read = match front.row_count.checked_sub(need_read) {
1690                        Some(remaining) if remaining != 0 => {
1691                            // if page row count less than batch_size we must set batch size to page row count.
1692                            // add check avoid dead loop
1693                            selectors_cursor.return_selector(RowSelector::select(remaining));
1694                            need_read
1695                        }
1696                        _ => front.row_count,
1697                    };
1698                    match self.array_reader.read_records(to_read)? {
1699                        0 => break,
1700                        rec => read_records += rec,
1701                    }
1702                }
1703            }
1704            RowSelectionCursor::All => {
1705                self.array_reader.read_records(batch_size)?;
1706            }
1707        }
1708
1709        let batch = consume_record_batch(self.array_reader.as_mut())?;
1710        Ok(if batch.num_rows() > 0 {
1711            Some(batch)
1712        } else {
1713            None
1714        })
1715    }
1716}
1717
1718impl RecordBatchReader for ParquetRecordBatchReader {
1719    /// Returns the projected [`SchemaRef`] for reading the parquet file.
1720    ///
1721    /// Note that the schema metadata will be stripped here. See
1722    /// [`ParquetRecordBatchReaderBuilder::schema`] if the metadata is desired.
1723    fn schema(&self) -> SchemaRef {
1724        self.schema.clone()
1725    }
1726}
1727
1728impl ParquetRecordBatchReader {
1729    /// Create a new [`ParquetRecordBatchReader`] from the provided chunk reader
1730    ///
1731    /// See [`ParquetRecordBatchReaderBuilder`] for more options
1732    pub fn try_new<T: ChunkReader + 'static>(reader: T, batch_size: usize) -> Result<Self> {
1733        ParquetRecordBatchReaderBuilder::try_new(reader)?
1734            .with_batch_size(batch_size)
1735            .build()
1736    }
1737
1738    /// Create a new [`ParquetRecordBatchReader`] from the provided [`RowGroups`]
1739    ///
1740    /// Note: this is a low-level interface see [`ParquetRecordBatchReader::try_new`] for a
1741    /// higher-level interface for reading parquet data from a file
1742    pub fn try_new_with_row_groups(
1743        levels: &FieldLevels,
1744        row_groups: &dyn RowGroups,
1745        batch_size: usize,
1746        selection: Option<RowSelection>,
1747    ) -> Result<Self> {
1748        // note metrics are not supported in this API
1749        let metrics = ArrowReaderMetrics::disabled();
1750        let array_reader = ArrayReaderBuilder::new(row_groups, &metrics)
1751            .with_batch_size(batch_size)
1752            .with_parquet_metadata(row_groups.metadata())
1753            .build_array_reader(levels.levels.as_ref(), &ProjectionMask::all())?;
1754
1755        let read_plan = ReadPlanBuilder::new(batch_size)
1756            .with_selection(selection)
1757            .build();
1758
1759        Ok(Self {
1760            array_reader,
1761            schema: Arc::new(Schema::new(levels.fields.clone())),
1762            read_plan,
1763        })
1764    }
1765
1766    /// Create a new [`ParquetRecordBatchReader`] that will read at most `batch_size` rows at
1767    /// a time from [`ArrayReader`] based on the configured `selection`. If `selection` is `None`
1768    /// all rows will be returned
1769    pub(crate) fn new(array_reader: Box<dyn ArrayReader>, read_plan: ReadPlan) -> Self {
1770        let schema = match array_reader.get_data_type() {
1771            ArrowType::Struct(fields) => Schema::new(fields.clone()),
1772            _ => unreachable!("Struct array reader's data type is not struct!"),
1773        };
1774
1775        Self {
1776            array_reader,
1777            schema: Arc::new(schema),
1778            read_plan,
1779        }
1780    }
1781
1782    #[inline(always)]
1783    pub(crate) fn batch_size(&self) -> usize {
1784        self.read_plan.batch_size()
1785    }
1786}
1787
1788#[cfg(test)]
1789pub(crate) mod tests {
1790    use std::cmp::min;
1791    use std::collections::{HashMap, VecDeque};
1792    use std::fmt::Formatter;
1793    use std::fs::File;
1794    use std::io::Seek;
1795    use std::path::PathBuf;
1796    use std::sync::Arc;
1797
1798    use rand::rngs::StdRng;
1799    use rand::{Rng, RngExt, SeedableRng, random, rng};
1800    use tempfile::tempfile;
1801
1802    use crate::arrow::arrow_reader::{
1803        ArrowPredicateFn, ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReader,
1804        ParquetRecordBatchReaderBuilder, RowFilter, RowGroupPlan, RowGroupSelection, RowSelection,
1805        RowSelector,
1806    };
1807    use crate::arrow::schema::{
1808        add_encoded_arrow_schema_to_metadata,
1809        virtual_type::{RowGroupIndex, RowNumber},
1810    };
1811    use crate::arrow::{ArrowWriter, ProjectionMask};
1812    use crate::basic::{ConvertedType, Encoding, Repetition, Type as PhysicalType};
1813    use crate::column::reader::decoder::REPETITION_LEVELS_BATCH_SIZE;
1814    use crate::data_type::{
1815        BoolType, ByteArray, ByteArrayType, DataType, DoubleType, FixedLenByteArray,
1816        FixedLenByteArrayType, FloatType, Int32Type, Int64Type, Int96, Int96Type,
1817    };
1818    use crate::errors::Result;
1819    use crate::file::metadata::{PageIndexPolicy, ParquetMetaData, ParquetStatisticsPolicy};
1820    use crate::file::properties::{EnabledStatistics, WriterProperties, WriterVersion};
1821    use crate::file::writer::{SerializedFileWriter, SerializedRowGroupWriter};
1822    use crate::schema::parser::parse_message_type;
1823    use crate::schema::types::{Type, TypePtr};
1824    use crate::util::test_common::rand_gen::RandGen;
1825    use arrow_array::builder::*;
1826    use arrow_array::cast::AsArray;
1827    use arrow_array::types::{
1828        Date32Type, Date64Type, Decimal32Type, Decimal64Type, Decimal128Type, Decimal256Type,
1829        DecimalType, Float16Type, Float32Type, Float64Type, Time32MillisecondType,
1830        Time64MicrosecondType,
1831    };
1832    use arrow_array::*;
1833    use arrow_buffer::{ArrowNativeType, BooleanBuffer, Buffer, IntervalDayTime, NullBuffer, i256};
1834    use arrow_data::{ArrayData, ArrayDataBuilder};
1835    use arrow_schema::{DataType as ArrowDataType, Field, Fields, Schema, SchemaRef, TimeUnit};
1836    use arrow_select::concat::concat_batches;
1837    use bytes::Bytes;
1838    use half::f16;
1839    use num_traits::PrimInt;
1840
1841    fn row_selection(rows: usize) -> RowSelection {
1842        RowSelection::from(vec![RowSelector::select(rows)])
1843    }
1844
1845    #[test]
1846    fn row_group_selection_accessors() {
1847        let row_group = RowGroupSelection::new(3, Some(row_selection(5)));
1848        assert_eq!(row_group.row_group_index(), 3);
1849        assert_eq!(row_group.selection().unwrap().row_count(), 5);
1850
1851        let row_group = RowGroupSelection::new(4, None);
1852        assert_eq!(row_group.row_group_index(), 4);
1853        assert!(row_group.selection().is_none());
1854    }
1855
1856    #[test]
1857    fn row_group_plan_tracks_global_configuration() {
1858        let mut plan = RowGroupPlan::Global {
1859            row_groups: None,
1860            selection: None,
1861        };
1862        plan.set_row_groups(vec![0]);
1863        plan.set_row_groups(vec![1, 2]);
1864        plan.set_row_selection(row_selection(3));
1865        plan.set_row_selection(row_selection(4));
1866
1867        let (row_groups, selection) = plan.into_global().unwrap();
1868        assert_eq!(row_groups, Some(vec![1, 2]));
1869        assert_eq!(selection.unwrap().row_count(), 4);
1870    }
1871
1872    #[test]
1873    fn row_group_plan_replaces_local_configuration() {
1874        let mut plan = RowGroupPlan::Global {
1875            row_groups: None,
1876            selection: None,
1877        };
1878        plan.set_row_group_selections(vec![RowGroupSelection::new(0, None)]);
1879        plan.set_row_group_selections(vec![RowGroupSelection::new(1, None)]);
1880
1881        let RowGroupPlan::PerRowGroup(row_groups) = plan else {
1882            panic!("expected per-row-group plan");
1883        };
1884        assert_eq!(row_groups, vec![RowGroupSelection::new(1, None)]);
1885    }
1886
1887    #[test]
1888    fn row_group_plan_rejects_mixed_configuration() {
1889        let mut row_groups_then_local = RowGroupPlan::Global {
1890            row_groups: None,
1891            selection: None,
1892        };
1893        row_groups_then_local.set_row_groups(vec![0]);
1894        row_groups_then_local.set_row_group_selections(vec![RowGroupSelection::new(0, None)]);
1895        assert!(matches!(row_groups_then_local, RowGroupPlan::Conflicting));
1896        row_groups_then_local.set_row_groups(vec![1]);
1897        row_groups_then_local.set_row_selection(row_selection(1));
1898        row_groups_then_local.set_row_group_selections(vec![RowGroupSelection::new(1, None)]);
1899        assert!(row_groups_then_local.into_global().is_err());
1900
1901        let mut local_then_row_groups =
1902            RowGroupPlan::PerRowGroup(vec![RowGroupSelection::new(0, None)]);
1903        local_then_row_groups.set_row_groups(vec![0]);
1904        assert!(matches!(local_then_row_groups, RowGroupPlan::Conflicting));
1905
1906        let mut local_then_selection =
1907            RowGroupPlan::PerRowGroup(vec![RowGroupSelection::new(0, None)]);
1908        local_then_selection.set_row_selection(row_selection(1));
1909        assert!(matches!(local_then_selection, RowGroupPlan::Conflicting));
1910
1911        let local = RowGroupPlan::PerRowGroup(vec![RowGroupSelection::new(0, None)]);
1912        assert!(local.into_global().is_err());
1913    }
1914
1915    #[test]
1916    fn filter_mask_accumulator_handles_empty_single_and_multiple_chunks() {
1917        let first = BooleanBuffer::from(vec![true, false, false, false]);
1918        let second = BooleanBuffer::from(vec![true]);
1919        let third = BooleanBuffer::from(vec![false, true]);
1920
1921        assert!(super::FilterMaskAccumulator::default().finish().is_none());
1922
1923        let mut single = super::FilterMaskAccumulator::default();
1924        single.append(first.clone());
1925        assert_eq!(single.finish().unwrap(), first);
1926
1927        let mut combined = super::FilterMaskAccumulator::default();
1928        combined.append(BooleanBuffer::from(vec![true, false, false, false]));
1929        combined.append(second);
1930        combined.append(third);
1931        assert_eq!(
1932            combined.finish().unwrap(),
1933            BooleanBuffer::from(vec![true, false, false, false, true, false, true])
1934        );
1935    }
1936
1937    #[test]
1938    fn test_arrow_reader_all_columns() {
1939        let file = get_test_file("parquet/generated_simple_numerics/blogs.parquet");
1940
1941        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
1942        let original_schema = Arc::clone(builder.schema());
1943        let reader = builder.build().unwrap();
1944
1945        // Verify that the schema was correctly parsed
1946        assert_eq!(original_schema.fields(), reader.schema().fields());
1947    }
1948
1949    #[test]
1950    fn test_reuse_schema() {
1951        let file = get_test_file("parquet/alltypes-java.parquet");
1952
1953        let builder = ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1954        let expected = builder.metadata;
1955        let schema = expected.file_metadata().schema_descr_ptr();
1956
1957        let arrow_options = ArrowReaderOptions::new().with_parquet_schema(schema.clone());
1958        let builder =
1959            ParquetRecordBatchReaderBuilder::try_new_with_options(file, arrow_options).unwrap();
1960
1961        // Verify that the metadata matches
1962        assert_eq!(expected.as_ref(), builder.metadata.as_ref());
1963    }
1964
1965    #[test]
1966    fn test_page_encoding_stats_mask() {
1967        let testdata = arrow::util::test_util::parquet_test_data();
1968        let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1969        let file = File::open(path).unwrap();
1970
1971        let arrow_options = ArrowReaderOptions::new().with_encoding_stats_as_mask(true);
1972        let builder =
1973            ParquetRecordBatchReaderBuilder::try_new_with_options(file, arrow_options).unwrap();
1974
1975        let row_group_metadata = builder.metadata.row_group(0);
1976
1977        // test page encoding stats
1978        let page_encoding_stats = row_group_metadata
1979            .column(0)
1980            .page_encoding_stats_mask()
1981            .unwrap();
1982        assert!(page_encoding_stats.is_only(Encoding::PLAIN));
1983        let page_encoding_stats = row_group_metadata
1984            .column(2)
1985            .page_encoding_stats_mask()
1986            .unwrap();
1987        assert!(page_encoding_stats.is_only(Encoding::PLAIN_DICTIONARY));
1988    }
1989
1990    #[test]
1991    fn test_stats_stats_skipped() {
1992        let testdata = arrow::util::test_util::parquet_test_data();
1993        let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1994        let file = File::open(path).unwrap();
1995
1996        // test skipping all
1997        let arrow_options = ArrowReaderOptions::new()
1998            .with_encoding_stats_policy(ParquetStatisticsPolicy::SkipAll)
1999            .with_column_stats_policy(ParquetStatisticsPolicy::SkipAll);
2000        let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2001            file.try_clone().unwrap(),
2002            arrow_options,
2003        )
2004        .unwrap();
2005
2006        let row_group_metadata = builder.metadata.row_group(0);
2007        for column in row_group_metadata.columns() {
2008            assert!(column.page_encoding_stats().is_none());
2009            assert!(column.page_encoding_stats_mask().is_none());
2010            assert!(column.statistics().is_none());
2011        }
2012
2013        // test skipping all but one column and converting to mask
2014        let arrow_options = ArrowReaderOptions::new()
2015            .with_encoding_stats_as_mask(true)
2016            .with_encoding_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
2017            .with_column_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]));
2018        let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2019            file.try_clone().unwrap(),
2020            arrow_options,
2021        )
2022        .unwrap();
2023
2024        let row_group_metadata = builder.metadata.row_group(0);
2025        for (idx, column) in row_group_metadata.columns().iter().enumerate() {
2026            assert!(column.page_encoding_stats().is_none());
2027            assert_eq!(column.page_encoding_stats_mask().is_some(), idx == 0);
2028            assert_eq!(column.statistics().is_some(), idx == 0);
2029        }
2030    }
2031
2032    #[test]
2033    fn test_size_stats_stats_skipped() {
2034        let testdata = arrow::util::test_util::parquet_test_data();
2035        let path = format!("{testdata}/repeated_primitive_no_list.parquet");
2036        let file = File::open(path).unwrap();
2037
2038        // test skipping all
2039        let arrow_options =
2040            ArrowReaderOptions::new().with_size_stats_policy(ParquetStatisticsPolicy::SkipAll);
2041        let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2042            file.try_clone().unwrap(),
2043            arrow_options,
2044        )
2045        .unwrap();
2046
2047        let row_group_metadata = builder.metadata.row_group(0);
2048        for column in row_group_metadata.columns() {
2049            assert!(column.repetition_level_histogram().is_none());
2050            assert!(column.definition_level_histogram().is_none());
2051            assert!(column.unencoded_byte_array_data_bytes().is_none());
2052        }
2053
2054        // test skipping all but one column and converting to mask
2055        let arrow_options = ArrowReaderOptions::new()
2056            .with_encoding_stats_as_mask(true)
2057            .with_size_stats_policy(ParquetStatisticsPolicy::skip_except(&[1]));
2058        let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2059            file.try_clone().unwrap(),
2060            arrow_options,
2061        )
2062        .unwrap();
2063
2064        let row_group_metadata = builder.metadata.row_group(0);
2065        for (idx, column) in row_group_metadata.columns().iter().enumerate() {
2066            assert_eq!(column.repetition_level_histogram().is_some(), idx == 1);
2067            assert_eq!(column.definition_level_histogram().is_some(), idx == 1);
2068            assert_eq!(column.unencoded_byte_array_data_bytes().is_some(), idx == 1);
2069        }
2070    }
2071
2072    #[test]
2073    fn test_arrow_reader_single_column() {
2074        let file = get_test_file("parquet/generated_simple_numerics/blogs.parquet");
2075
2076        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2077        let original_schema = Arc::clone(builder.schema());
2078
2079        let mask = ProjectionMask::leaves(builder.parquet_schema(), [2]);
2080        let reader = builder.with_projection(mask).build().unwrap();
2081
2082        // Verify that the schema was correctly parsed
2083        assert_eq!(1, reader.schema().fields().len());
2084        assert_eq!(original_schema.fields()[1], reader.schema().fields()[0]);
2085    }
2086
2087    #[test]
2088    fn test_arrow_reader_single_column_by_name() {
2089        let file = get_test_file("parquet/generated_simple_numerics/blogs.parquet");
2090
2091        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2092        let original_schema = Arc::clone(builder.schema());
2093
2094        let mask = ProjectionMask::columns(builder.parquet_schema(), ["blog_id"]);
2095        let reader = builder.with_projection(mask).build().unwrap();
2096
2097        // Verify that the schema was correctly parsed
2098        assert_eq!(1, reader.schema().fields().len());
2099        assert_eq!(original_schema.fields()[1], reader.schema().fields()[0]);
2100    }
2101
2102    #[test]
2103    fn test_null_column_reader_test() {
2104        let mut file = tempfile::tempfile().unwrap();
2105
2106        let schema = "
2107            message message {
2108                OPTIONAL INT32 int32;
2109            }
2110        ";
2111        let schema = Arc::new(parse_message_type(schema).unwrap());
2112
2113        let def_levels = vec![vec![0, 0, 0], vec![0, 0, 0, 0]];
2114        generate_single_column_file_with_data::<Int32Type>(
2115            &[vec![], vec![]],
2116            Some(&def_levels),
2117            file.try_clone().unwrap(), // Cannot use &mut File (#1163)
2118            schema,
2119            Some(Field::new("int32", ArrowDataType::Null, true)),
2120            &Default::default(),
2121        )
2122        .unwrap();
2123
2124        file.rewind().unwrap();
2125
2126        let record_reader = ParquetRecordBatchReader::try_new(file, 2).unwrap();
2127        let batches = record_reader.collect::<Result<Vec<_>, _>>().unwrap();
2128
2129        assert_eq!(batches.len(), 4);
2130        for batch in &batches[0..3] {
2131            assert_eq!(batch.num_rows(), 2);
2132            assert_eq!(batch.num_columns(), 1);
2133            assert_eq!(batch.column(0).null_count(), 2);
2134        }
2135
2136        assert_eq!(batches[3].num_rows(), 1);
2137        assert_eq!(batches[3].num_columns(), 1);
2138        assert_eq!(batches[3].column(0).null_count(), 1);
2139    }
2140
2141    #[test]
2142    #[cfg_attr(miri, ignore)] // Takes too long
2143    fn test_primitive_single_column_reader_test() {
2144        run_single_column_reader_tests::<BoolType, _, BoolType>(
2145            2,
2146            ConvertedType::NONE,
2147            None,
2148            |vals| Arc::new(BooleanArray::from_iter(vals.iter().copied())),
2149            &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2150        );
2151        run_single_column_reader_tests::<Int32Type, _, Int32Type>(
2152            2,
2153            ConvertedType::NONE,
2154            None,
2155            |vals| Arc::new(Int32Array::from_iter(vals.iter().copied())),
2156            &[
2157                Encoding::PLAIN,
2158                Encoding::RLE_DICTIONARY,
2159                Encoding::DELTA_BINARY_PACKED,
2160                Encoding::BYTE_STREAM_SPLIT,
2161            ],
2162        );
2163        run_single_column_reader_tests::<Int64Type, _, Int64Type>(
2164            2,
2165            ConvertedType::NONE,
2166            None,
2167            |vals| Arc::new(Int64Array::from_iter(vals.iter().copied())),
2168            &[
2169                Encoding::PLAIN,
2170                Encoding::RLE_DICTIONARY,
2171                Encoding::DELTA_BINARY_PACKED,
2172                Encoding::BYTE_STREAM_SPLIT,
2173            ],
2174        );
2175        run_single_column_reader_tests::<FloatType, _, FloatType>(
2176            2,
2177            ConvertedType::NONE,
2178            None,
2179            |vals| Arc::new(Float32Array::from_iter(vals.iter().copied())),
2180            &[Encoding::PLAIN, Encoding::BYTE_STREAM_SPLIT],
2181        );
2182    }
2183
2184    #[test]
2185    #[cfg_attr(miri, ignore)] // Takes too long
2186    fn test_unsigned_primitive_single_column_reader_test() {
2187        run_single_column_reader_tests::<Int32Type, _, Int32Type>(
2188            2,
2189            ConvertedType::UINT_32,
2190            Some(ArrowDataType::UInt32),
2191            |vals| {
2192                Arc::new(UInt32Array::from_iter(
2193                    vals.iter().map(|x| x.map(|x| x as u32)),
2194                ))
2195            },
2196            &[
2197                Encoding::PLAIN,
2198                Encoding::RLE_DICTIONARY,
2199                Encoding::DELTA_BINARY_PACKED,
2200            ],
2201        );
2202        run_single_column_reader_tests::<Int64Type, _, Int64Type>(
2203            2,
2204            ConvertedType::UINT_64,
2205            Some(ArrowDataType::UInt64),
2206            |vals| {
2207                Arc::new(UInt64Array::from_iter(
2208                    vals.iter().map(|x| x.map(|x| x as u64)),
2209                ))
2210            },
2211            &[
2212                Encoding::PLAIN,
2213                Encoding::RLE_DICTIONARY,
2214                Encoding::DELTA_BINARY_PACKED,
2215            ],
2216        );
2217    }
2218
2219    #[test]
2220    fn test_unsigned_roundtrip() {
2221        let schema = Arc::new(Schema::new(vec![
2222            Field::new("uint32", ArrowDataType::UInt32, true),
2223            Field::new("uint64", ArrowDataType::UInt64, true),
2224        ]));
2225
2226        let mut buf = Vec::with_capacity(1024);
2227        let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None).unwrap();
2228
2229        let original = RecordBatch::try_new(
2230            schema,
2231            vec![
2232                Arc::new(UInt32Array::from_iter_values([
2233                    0,
2234                    i32::MAX as u32,
2235                    u32::MAX,
2236                ])),
2237                Arc::new(UInt64Array::from_iter_values([
2238                    0,
2239                    i64::MAX as u64,
2240                    u64::MAX,
2241                ])),
2242            ],
2243        )
2244        .unwrap();
2245
2246        writer.write(&original).unwrap();
2247        writer.close().unwrap();
2248
2249        let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024).unwrap();
2250        let ret = reader.next().unwrap().unwrap();
2251        assert_eq!(ret, original);
2252
2253        // Check they can be downcast to the correct type
2254        ret.column(0)
2255            .as_any()
2256            .downcast_ref::<UInt32Array>()
2257            .unwrap();
2258
2259        ret.column(1)
2260            .as_any()
2261            .downcast_ref::<UInt64Array>()
2262            .unwrap();
2263    }
2264
2265    #[test]
2266    fn test_float16_roundtrip() -> Result<()> {
2267        let schema = Arc::new(Schema::new(vec![
2268            Field::new("float16", ArrowDataType::Float16, false),
2269            Field::new("float16-nullable", ArrowDataType::Float16, true),
2270        ]));
2271
2272        let mut buf = Vec::with_capacity(1024);
2273        let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
2274
2275        let original = RecordBatch::try_new(
2276            schema,
2277            vec![
2278                Arc::new(Float16Array::from_iter_values([
2279                    f16::EPSILON,
2280                    f16::MIN,
2281                    f16::MAX,
2282                    f16::NAN,
2283                    f16::INFINITY,
2284                    f16::NEG_INFINITY,
2285                    f16::ONE,
2286                    f16::NEG_ONE,
2287                    f16::ZERO,
2288                    f16::NEG_ZERO,
2289                    f16::E,
2290                    f16::PI,
2291                    f16::FRAC_1_PI,
2292                ])),
2293                Arc::new(Float16Array::from(vec![
2294                    None,
2295                    None,
2296                    None,
2297                    Some(f16::NAN),
2298                    Some(f16::INFINITY),
2299                    Some(f16::NEG_INFINITY),
2300                    None,
2301                    None,
2302                    None,
2303                    None,
2304                    None,
2305                    None,
2306                    Some(f16::FRAC_1_PI),
2307                ])),
2308            ],
2309        )?;
2310
2311        writer.write(&original)?;
2312        writer.close()?;
2313
2314        let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024)?;
2315        let ret = reader.next().unwrap()?;
2316        assert_eq!(ret, original);
2317
2318        // Ensure can be downcast to the correct type
2319        ret.column(0).as_primitive::<Float16Type>();
2320        ret.column(1).as_primitive::<Float16Type>();
2321
2322        Ok(())
2323    }
2324
2325    #[test]
2326    fn test_time_utc_roundtrip() -> Result<()> {
2327        let schema = Arc::new(Schema::new(vec![
2328            Field::new(
2329                "time_millis",
2330                ArrowDataType::Time32(TimeUnit::Millisecond),
2331                true,
2332            )
2333            .with_metadata(HashMap::from_iter(vec![(
2334                "adjusted_to_utc".to_string(),
2335                String::new(),
2336            )])),
2337            Field::new(
2338                "time_micros",
2339                ArrowDataType::Time64(TimeUnit::Microsecond),
2340                true,
2341            )
2342            .with_metadata(HashMap::from_iter(vec![(
2343                "adjusted_to_utc".to_string(),
2344                String::new(),
2345            )])),
2346        ]));
2347
2348        let mut buf = Vec::with_capacity(1024);
2349        let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
2350
2351        let original = RecordBatch::try_new(
2352            schema,
2353            vec![
2354                Arc::new(Time32MillisecondArray::from(vec![
2355                    Some(-1),
2356                    Some(0),
2357                    Some(86_399_000),
2358                    Some(86_400_000),
2359                    Some(86_401_000),
2360                    None,
2361                ])),
2362                Arc::new(Time64MicrosecondArray::from(vec![
2363                    Some(-1),
2364                    Some(0),
2365                    Some(86_399 * 1_000_000),
2366                    Some(86_400 * 1_000_000),
2367                    Some(86_401 * 1_000_000),
2368                    None,
2369                ])),
2370            ],
2371        )?;
2372
2373        writer.write(&original)?;
2374        writer.close()?;
2375
2376        let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024)?;
2377        let ret = reader.next().unwrap()?;
2378        assert_eq!(ret, original);
2379
2380        // Ensure can be downcast to the correct type
2381        ret.column(0).as_primitive::<Time32MillisecondType>();
2382        ret.column(1).as_primitive::<Time64MicrosecondType>();
2383
2384        Ok(())
2385    }
2386
2387    #[test]
2388    fn test_date32_roundtrip() -> Result<()> {
2389        use arrow_array::Date32Array;
2390
2391        let schema = Arc::new(Schema::new(vec![Field::new(
2392            "date32",
2393            ArrowDataType::Date32,
2394            false,
2395        )]));
2396
2397        let mut buf = Vec::with_capacity(1024);
2398
2399        let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
2400
2401        let original = RecordBatch::try_new(
2402            schema,
2403            vec![Arc::new(Date32Array::from(vec![
2404                -1_000_000, -100_000, -10_000, -1_000, 0, 1_000, 10_000, 100_000, 1_000_000,
2405            ]))],
2406        )?;
2407
2408        writer.write(&original)?;
2409        writer.close()?;
2410
2411        let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024)?;
2412        let ret = reader.next().unwrap()?;
2413        assert_eq!(ret, original);
2414
2415        // Ensure can be downcast to the correct type
2416        ret.column(0).as_primitive::<Date32Type>();
2417
2418        Ok(())
2419    }
2420
2421    #[test]
2422    fn test_date64_roundtrip() -> Result<()> {
2423        use arrow_array::Date64Array;
2424
2425        let schema = Arc::new(Schema::new(vec![
2426            Field::new("small-date64", ArrowDataType::Date64, false),
2427            Field::new("big-date64", ArrowDataType::Date64, false),
2428            Field::new("invalid-date64", ArrowDataType::Date64, false),
2429        ]));
2430
2431        let mut default_buf = Vec::with_capacity(1024);
2432        let mut coerce_buf = Vec::with_capacity(1024);
2433
2434        let coerce_props = WriterProperties::builder().set_coerce_types(true).build();
2435
2436        let mut default_writer = ArrowWriter::try_new(&mut default_buf, schema.clone(), None)?;
2437        let mut coerce_writer =
2438            ArrowWriter::try_new(&mut coerce_buf, schema.clone(), Some(coerce_props))?;
2439
2440        static NUM_MILLISECONDS_IN_DAY: i64 = 1000 * 60 * 60 * 24;
2441
2442        let original = RecordBatch::try_new(
2443            schema,
2444            vec![
2445                // small-date64
2446                Arc::new(Date64Array::from(vec![
2447                    -1_000_000 * NUM_MILLISECONDS_IN_DAY,
2448                    -1_000 * NUM_MILLISECONDS_IN_DAY,
2449                    0,
2450                    1_000 * NUM_MILLISECONDS_IN_DAY,
2451                    1_000_000 * NUM_MILLISECONDS_IN_DAY,
2452                ])),
2453                // big-date64
2454                Arc::new(Date64Array::from(vec![
2455                    -10_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2456                    -1_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2457                    0,
2458                    1_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2459                    10_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2460                ])),
2461                // invalid-date64
2462                Arc::new(Date64Array::from(vec![
2463                    -1_000_000 * NUM_MILLISECONDS_IN_DAY + 1,
2464                    -1_000 * NUM_MILLISECONDS_IN_DAY + 1,
2465                    1,
2466                    1_000 * NUM_MILLISECONDS_IN_DAY + 1,
2467                    1_000_000 * NUM_MILLISECONDS_IN_DAY + 1,
2468                ])),
2469            ],
2470        )?;
2471
2472        default_writer.write(&original)?;
2473        coerce_writer.write(&original)?;
2474
2475        default_writer.close()?;
2476        coerce_writer.close()?;
2477
2478        let mut default_reader = ParquetRecordBatchReader::try_new(Bytes::from(default_buf), 1024)?;
2479        let mut coerce_reader = ParquetRecordBatchReader::try_new(Bytes::from(coerce_buf), 1024)?;
2480
2481        let default_ret = default_reader.next().unwrap()?;
2482        let coerce_ret = coerce_reader.next().unwrap()?;
2483
2484        // Roundtrip should be successful when default writer used
2485        assert_eq!(default_ret, original);
2486
2487        // Only small-date64 should roundtrip successfully when coerce_types writer is used
2488        assert_eq!(coerce_ret.column(0), original.column(0));
2489        assert_ne!(coerce_ret.column(1), original.column(1));
2490        assert_ne!(coerce_ret.column(2), original.column(2));
2491
2492        // Ensure both can be downcast to the correct type
2493        default_ret.column(0).as_primitive::<Date64Type>();
2494        coerce_ret.column(0).as_primitive::<Date64Type>();
2495
2496        Ok(())
2497    }
2498    struct RandFixedLenGen {}
2499
2500    impl RandGen<FixedLenByteArrayType> for RandFixedLenGen {
2501        fn r#gen(len: i32) -> FixedLenByteArray {
2502            let mut v = vec![0u8; len as usize];
2503            rng().fill_bytes(&mut v);
2504            ByteArray::from(v).into()
2505        }
2506    }
2507
2508    #[test]
2509    #[cfg_attr(miri, ignore)] // Takes too long
2510    fn test_fixed_length_binary_column_reader() {
2511        run_single_column_reader_tests::<FixedLenByteArrayType, _, RandFixedLenGen>(
2512            20,
2513            ConvertedType::NONE,
2514            None,
2515            |vals| {
2516                let mut builder = FixedSizeBinaryBuilder::with_capacity(vals.len(), 20);
2517                for val in vals {
2518                    match val {
2519                        Some(b) => builder.append_value(b).unwrap(),
2520                        None => builder.append_null(),
2521                    }
2522                }
2523                Arc::new(builder.finish())
2524            },
2525            &[Encoding::PLAIN, Encoding::RLE_DICTIONARY],
2526        );
2527    }
2528
2529    #[test]
2530    #[cfg_attr(miri, ignore)] // Takes too long
2531    fn test_interval_day_time_column_reader() {
2532        run_single_column_reader_tests::<FixedLenByteArrayType, _, RandFixedLenGen>(
2533            12,
2534            ConvertedType::INTERVAL,
2535            None,
2536            |vals| {
2537                Arc::new(
2538                    vals.iter()
2539                        .map(|x| {
2540                            x.as_ref().map(|b| IntervalDayTime {
2541                                days: i32::from_le_bytes(b.as_ref()[4..8].try_into().unwrap()),
2542                                milliseconds: i32::from_le_bytes(
2543                                    b.as_ref()[8..12].try_into().unwrap(),
2544                                ),
2545                            })
2546                        })
2547                        .collect::<IntervalDayTimeArray>(),
2548                )
2549            },
2550            &[Encoding::PLAIN, Encoding::RLE_DICTIONARY],
2551        );
2552    }
2553
2554    #[test]
2555    #[cfg_attr(miri, ignore)] // Takes too long
2556    fn test_int96_single_column_reader_test() {
2557        let encodings = &[Encoding::PLAIN, Encoding::RLE_DICTIONARY];
2558
2559        type TypeHintAndConversionFunction =
2560            (Option<ArrowDataType>, fn(&[Option<Int96>]) -> ArrayRef);
2561
2562        let resolutions: Vec<TypeHintAndConversionFunction> = vec![
2563            // Test without a specified ArrowType hint.
2564            (None, |vals: &[Option<Int96>]| {
2565                Arc::new(TimestampNanosecondArray::from_iter(
2566                    vals.iter().map(|x| x.map(|x| x.to_nanos())),
2567                )) as ArrayRef
2568            }),
2569            // Test other TimeUnits as ArrowType hints.
2570            (
2571                Some(ArrowDataType::Timestamp(TimeUnit::Second, None)),
2572                |vals: &[Option<Int96>]| {
2573                    Arc::new(TimestampSecondArray::from_iter(
2574                        vals.iter().map(|x| x.map(|x| x.to_seconds())),
2575                    )) as ArrayRef
2576                },
2577            ),
2578            (
2579                Some(ArrowDataType::Timestamp(TimeUnit::Millisecond, None)),
2580                |vals: &[Option<Int96>]| {
2581                    Arc::new(TimestampMillisecondArray::from_iter(
2582                        vals.iter().map(|x| x.map(|x| x.to_millis())),
2583                    )) as ArrayRef
2584                },
2585            ),
2586            (
2587                Some(ArrowDataType::Timestamp(TimeUnit::Microsecond, None)),
2588                |vals: &[Option<Int96>]| {
2589                    Arc::new(TimestampMicrosecondArray::from_iter(
2590                        vals.iter().map(|x| x.map(|x| x.to_micros())),
2591                    )) as ArrayRef
2592                },
2593            ),
2594            (
2595                Some(ArrowDataType::Timestamp(TimeUnit::Nanosecond, None)),
2596                |vals: &[Option<Int96>]| {
2597                    Arc::new(TimestampNanosecondArray::from_iter(
2598                        vals.iter().map(|x| x.map(|x| x.to_nanos())),
2599                    )) as ArrayRef
2600                },
2601            ),
2602            // Test another timezone with TimeUnit as ArrowType hints.
2603            (
2604                Some(ArrowDataType::Timestamp(
2605                    TimeUnit::Second,
2606                    Some(Arc::from("-05:00")),
2607                )),
2608                |vals: &[Option<Int96>]| {
2609                    Arc::new(
2610                        TimestampSecondArray::from_iter(
2611                            vals.iter().map(|x| x.map(|x| x.to_seconds())),
2612                        )
2613                        .with_timezone("-05:00"),
2614                    ) as ArrayRef
2615                },
2616            ),
2617        ];
2618
2619        resolutions.iter().for_each(|(arrow_type, converter)| {
2620            run_single_column_reader_tests::<Int96Type, _, Int96Type>(
2621                2,
2622                ConvertedType::NONE,
2623                arrow_type.clone(),
2624                converter,
2625                encodings,
2626            );
2627        })
2628    }
2629
2630    struct RandUtf8Gen {}
2631
2632    impl RandGen<ByteArrayType> for RandUtf8Gen {
2633        fn r#gen(len: i32) -> ByteArray {
2634            Int32Type::r#gen(len).to_string().as_str().into()
2635        }
2636    }
2637
2638    #[test]
2639    #[cfg_attr(miri, ignore)] // Takes too long
2640    fn test_utf8_single_column_reader_test() {
2641        fn string_converter<O: OffsetSizeTrait>(vals: &[Option<ByteArray>]) -> ArrayRef {
2642            Arc::new(GenericStringArray::<O>::from_iter(vals.iter().map(|x| {
2643                x.as_ref().map(|b| std::str::from_utf8(b.data()).unwrap())
2644            })))
2645        }
2646
2647        let encodings = &[
2648            Encoding::PLAIN,
2649            Encoding::RLE_DICTIONARY,
2650            Encoding::DELTA_LENGTH_BYTE_ARRAY,
2651            Encoding::DELTA_BYTE_ARRAY,
2652        ];
2653
2654        run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2655            2,
2656            ConvertedType::NONE,
2657            None,
2658            |vals| {
2659                Arc::new(BinaryArray::from_iter(
2660                    vals.iter().map(|x| x.as_ref().map(|x| x.data())),
2661                ))
2662            },
2663            encodings,
2664        );
2665
2666        run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2667            2,
2668            ConvertedType::UTF8,
2669            None,
2670            string_converter::<i32>,
2671            encodings,
2672        );
2673
2674        run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2675            2,
2676            ConvertedType::UTF8,
2677            Some(ArrowDataType::Utf8),
2678            string_converter::<i32>,
2679            encodings,
2680        );
2681
2682        run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2683            2,
2684            ConvertedType::UTF8,
2685            Some(ArrowDataType::LargeUtf8),
2686            string_converter::<i64>,
2687            encodings,
2688        );
2689
2690        let small_key_types = [ArrowDataType::Int8, ArrowDataType::UInt8];
2691        for key in &small_key_types {
2692            for encoding in encodings {
2693                let mut opts = TestOptions::new(2, 20, 15).with_null_percent(50);
2694                opts.encoding = *encoding;
2695
2696                let data_type =
2697                    ArrowDataType::Dictionary(Box::new(key.clone()), Box::new(ArrowDataType::Utf8));
2698
2699                // Cannot run full test suite as keys overflow, run small test instead
2700                single_column_reader_test::<ByteArrayType, _, RandUtf8Gen>(
2701                    opts,
2702                    2,
2703                    ConvertedType::UTF8,
2704                    Some(data_type.clone()),
2705                    move |vals| {
2706                        let vals = string_converter::<i32>(vals);
2707                        arrow::compute::cast(&vals, &data_type).unwrap()
2708                    },
2709                );
2710            }
2711        }
2712
2713        let key_types = [
2714            ArrowDataType::Int16,
2715            ArrowDataType::UInt16,
2716            ArrowDataType::Int32,
2717            ArrowDataType::UInt32,
2718            ArrowDataType::Int64,
2719            ArrowDataType::UInt64,
2720        ];
2721
2722        for key in &key_types {
2723            let data_type =
2724                ArrowDataType::Dictionary(Box::new(key.clone()), Box::new(ArrowDataType::Utf8));
2725
2726            run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2727                2,
2728                ConvertedType::UTF8,
2729                Some(data_type.clone()),
2730                move |vals| {
2731                    let vals = string_converter::<i32>(vals);
2732                    arrow::compute::cast(&vals, &data_type).unwrap()
2733                },
2734                encodings,
2735            );
2736
2737            let data_type = ArrowDataType::Dictionary(
2738                Box::new(key.clone()),
2739                Box::new(ArrowDataType::LargeUtf8),
2740            );
2741
2742            run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2743                2,
2744                ConvertedType::UTF8,
2745                Some(data_type.clone()),
2746                move |vals| {
2747                    let vals = string_converter::<i64>(vals);
2748                    arrow::compute::cast(&vals, &data_type).unwrap()
2749                },
2750                encodings,
2751            );
2752        }
2753    }
2754
2755    #[test]
2756    fn test_decimal_nullable_struct() {
2757        let decimals = Decimal256Array::from_iter_values(
2758            [1, 2, 3, 4, 5, 6, 7, 8].into_iter().map(i256::from_i128),
2759        );
2760
2761        let data = ArrayDataBuilder::new(ArrowDataType::Struct(Fields::from(vec![Field::new(
2762            "decimals",
2763            decimals.data_type().clone(),
2764            false,
2765        )])))
2766        .len(8)
2767        .null_bit_buffer(Some(Buffer::from(&[0b11101111])))
2768        .child_data(vec![decimals.into_data()])
2769        .build()
2770        .unwrap();
2771
2772        let written =
2773            RecordBatch::try_from_iter([("struct", Arc::new(StructArray::from(data)) as ArrayRef)])
2774                .unwrap();
2775
2776        let mut buffer = Vec::with_capacity(1024);
2777        let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
2778        writer.write(&written).unwrap();
2779        writer.close().unwrap();
2780
2781        let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
2782            .unwrap()
2783            .collect::<Result<Vec<_>, _>>()
2784            .unwrap();
2785
2786        assert_eq!(&written.slice(0, 3), &read[0]);
2787        assert_eq!(&written.slice(3, 3), &read[1]);
2788        assert_eq!(&written.slice(6, 2), &read[2]);
2789    }
2790
2791    #[test]
2792    fn test_int32_nullable_struct() {
2793        let int32 = Int32Array::from_iter_values([1, 2, 3, 4, 5, 6, 7, 8]);
2794        let data = ArrayDataBuilder::new(ArrowDataType::Struct(Fields::from(vec![Field::new(
2795            "int32",
2796            int32.data_type().clone(),
2797            false,
2798        )])))
2799        .len(8)
2800        .null_bit_buffer(Some(Buffer::from(&[0b11101111])))
2801        .child_data(vec![int32.into_data()])
2802        .build()
2803        .unwrap();
2804
2805        let written =
2806            RecordBatch::try_from_iter([("struct", Arc::new(StructArray::from(data)) as ArrayRef)])
2807                .unwrap();
2808
2809        let mut buffer = Vec::with_capacity(1024);
2810        let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
2811        writer.write(&written).unwrap();
2812        writer.close().unwrap();
2813
2814        let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
2815            .unwrap()
2816            .collect::<Result<Vec<_>, _>>()
2817            .unwrap();
2818
2819        assert_eq!(&written.slice(0, 3), &read[0]);
2820        assert_eq!(&written.slice(3, 3), &read[1]);
2821        assert_eq!(&written.slice(6, 2), &read[2]);
2822    }
2823
2824    #[test]
2825    fn test_decimal_list() {
2826        let decimals = Decimal128Array::from_iter_values([1, 2, 3, 4, 5, 6, 7, 8]);
2827
2828        // [[], [1], [2, 3], null, [4], null, [6, 7, 8]]
2829        let data = ArrayDataBuilder::new(ArrowDataType::List(Arc::new(Field::new_list_field(
2830            decimals.data_type().clone(),
2831            false,
2832        ))))
2833        .len(7)
2834        .add_buffer(Buffer::from_iter([0_i32, 0, 1, 3, 3, 4, 5, 8]))
2835        .null_bit_buffer(Some(Buffer::from(&[0b01010111])))
2836        .child_data(vec![decimals.into_data()])
2837        .build()
2838        .unwrap();
2839
2840        let written =
2841            RecordBatch::try_from_iter([("list", Arc::new(ListArray::from(data)) as ArrayRef)])
2842                .unwrap();
2843
2844        let mut buffer = Vec::with_capacity(1024);
2845        let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
2846        writer.write(&written).unwrap();
2847        writer.close().unwrap();
2848
2849        let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
2850            .unwrap()
2851            .collect::<Result<Vec<_>, _>>()
2852            .unwrap();
2853
2854        assert_eq!(&written.slice(0, 3), &read[0]);
2855        assert_eq!(&written.slice(3, 3), &read[1]);
2856        assert_eq!(&written.slice(6, 1), &read[2]);
2857    }
2858
2859    #[test]
2860    fn test_read_decimal_file() {
2861        use arrow_array::Decimal128Array;
2862        let testdata = arrow::util::test_util::parquet_test_data();
2863        let file_variants = vec![
2864            ("byte_array", 4),
2865            ("fixed_length", 25),
2866            ("int32", 4),
2867            ("int64", 10),
2868        ];
2869        for (prefix, target_precision) in file_variants {
2870            let path = format!("{testdata}/{prefix}_decimal.parquet");
2871            let file = File::open(path).unwrap();
2872            let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
2873
2874            let batch = record_reader.next().unwrap().unwrap();
2875            assert_eq!(batch.num_rows(), 24);
2876            let col = batch
2877                .column(0)
2878                .as_any()
2879                .downcast_ref::<Decimal128Array>()
2880                .unwrap();
2881
2882            let expected = 1..25;
2883
2884            assert_eq!(col.precision(), target_precision);
2885            assert_eq!(col.scale(), 2);
2886
2887            for (i, v) in expected.enumerate() {
2888                assert_eq!(col.value(i), v * 100_i128);
2889            }
2890        }
2891    }
2892
2893    #[test]
2894    #[cfg_attr(miri, ignore)] // inline assembly is not supported
2895    fn test_read_float16_nonzeros_file() {
2896        use arrow_array::Float16Array;
2897        let testdata = arrow::util::test_util::parquet_test_data();
2898        // see https://github.com/apache/parquet-testing/pull/40
2899        let path = format!("{testdata}/float16_nonzeros_and_nans.parquet");
2900        let file = File::open(path).unwrap();
2901        let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
2902
2903        let batch = record_reader.next().unwrap().unwrap();
2904        assert_eq!(batch.num_rows(), 8);
2905        let col = batch
2906            .column(0)
2907            .as_any()
2908            .downcast_ref::<Float16Array>()
2909            .unwrap();
2910
2911        let f16_two = f16::ONE + f16::ONE;
2912
2913        assert_eq!(col.null_count(), 1);
2914        assert!(col.is_null(0));
2915        assert_eq!(col.value(1), f16::ONE);
2916        assert_eq!(col.value(2), -f16_two);
2917        assert!(col.value(3).is_nan());
2918        assert_eq!(col.value(4), f16::ZERO);
2919        assert!(col.value(4).is_sign_positive());
2920        assert_eq!(col.value(5), f16::NEG_ONE);
2921        assert_eq!(col.value(6), f16::NEG_ZERO);
2922        assert!(col.value(6).is_sign_negative());
2923        assert_eq!(col.value(7), f16_two);
2924    }
2925
2926    #[test]
2927    fn test_read_float16_zeros_file() {
2928        use arrow_array::Float16Array;
2929        let testdata = arrow::util::test_util::parquet_test_data();
2930        // see https://github.com/apache/parquet-testing/pull/40
2931        let path = format!("{testdata}/float16_zeros_and_nans.parquet");
2932        let file = File::open(path).unwrap();
2933        let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
2934
2935        let batch = record_reader.next().unwrap().unwrap();
2936        assert_eq!(batch.num_rows(), 3);
2937        let col = batch
2938            .column(0)
2939            .as_any()
2940            .downcast_ref::<Float16Array>()
2941            .unwrap();
2942
2943        assert_eq!(col.null_count(), 1);
2944        assert!(col.is_null(0));
2945        assert_eq!(col.value(1), f16::ZERO);
2946        assert!(col.value(1).is_sign_positive());
2947        assert!(col.value(2).is_nan());
2948    }
2949
2950    #[test]
2951    #[cfg_attr(miri, ignore)] // Zstd calls native C functions unsupported by Miri
2952    fn test_read_float32_float64_byte_stream_split() {
2953        let path = format!(
2954            "{}/byte_stream_split.zstd.parquet",
2955            arrow::util::test_util::parquet_test_data(),
2956        );
2957        let file = File::open(path).unwrap();
2958        let record_reader = ParquetRecordBatchReader::try_new(file, 128).unwrap();
2959
2960        let mut row_count = 0;
2961        for batch in record_reader {
2962            let batch = batch.unwrap();
2963            row_count += batch.num_rows();
2964            let f32_col = batch.column(0).as_primitive::<Float32Type>();
2965            let f64_col = batch.column(1).as_primitive::<Float64Type>();
2966
2967            // This file contains floats from a standard normal distribution
2968            for &x in f32_col.values() {
2969                assert!(x > -10.0);
2970                assert!(x < 10.0);
2971            }
2972            for &x in f64_col.values() {
2973                assert!(x > -10.0);
2974                assert!(x < 10.0);
2975            }
2976        }
2977        assert_eq!(row_count, 300);
2978    }
2979
2980    #[test]
2981    #[cfg_attr(miri, ignore)] // Takes too long
2982    fn test_read_extended_byte_stream_split() {
2983        let path = format!(
2984            "{}/byte_stream_split_extended.gzip.parquet",
2985            arrow::util::test_util::parquet_test_data(),
2986        );
2987        let file = File::open(path).unwrap();
2988        let record_reader = ParquetRecordBatchReader::try_new(file, 128).unwrap();
2989
2990        let mut row_count = 0;
2991        for batch in record_reader {
2992            let batch = batch.unwrap();
2993            row_count += batch.num_rows();
2994
2995            // 0,1 are f16
2996            let f16_col = batch.column(0).as_primitive::<Float16Type>();
2997            let f16_bss = batch.column(1).as_primitive::<Float16Type>();
2998            assert_eq!(f16_col.len(), f16_bss.len());
2999            f16_col
3000                .iter()
3001                .zip(f16_bss.iter())
3002                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3003
3004            // 2,3 are f32
3005            let f32_col = batch.column(2).as_primitive::<Float32Type>();
3006            let f32_bss = batch.column(3).as_primitive::<Float32Type>();
3007            assert_eq!(f32_col.len(), f32_bss.len());
3008            f32_col
3009                .iter()
3010                .zip(f32_bss.iter())
3011                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3012
3013            // 4,5 are f64
3014            let f64_col = batch.column(4).as_primitive::<Float64Type>();
3015            let f64_bss = batch.column(5).as_primitive::<Float64Type>();
3016            assert_eq!(f64_col.len(), f64_bss.len());
3017            f64_col
3018                .iter()
3019                .zip(f64_bss.iter())
3020                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3021
3022            // 6,7 are i32
3023            let i32_col = batch.column(6).as_primitive::<types::Int32Type>();
3024            let i32_bss = batch.column(7).as_primitive::<types::Int32Type>();
3025            assert_eq!(i32_col.len(), i32_bss.len());
3026            i32_col
3027                .iter()
3028                .zip(i32_bss.iter())
3029                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3030
3031            // 8,9 are i64
3032            let i64_col = batch.column(8).as_primitive::<types::Int64Type>();
3033            let i64_bss = batch.column(9).as_primitive::<types::Int64Type>();
3034            assert_eq!(i64_col.len(), i64_bss.len());
3035            i64_col
3036                .iter()
3037                .zip(i64_bss.iter())
3038                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3039
3040            // 10,11 are FLBA(5)
3041            let flba_col = batch.column(10).as_fixed_size_binary();
3042            let flba_bss = batch.column(11).as_fixed_size_binary();
3043            assert_eq!(flba_col.len(), flba_bss.len());
3044            flba_col
3045                .iter()
3046                .zip(flba_bss.iter())
3047                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3048
3049            // 12,13 are FLBA(4) (decimal(7,3))
3050            let dec_col = batch.column(12).as_primitive::<Decimal128Type>();
3051            let dec_bss = batch.column(13).as_primitive::<Decimal128Type>();
3052            assert_eq!(dec_col.len(), dec_bss.len());
3053            dec_col
3054                .iter()
3055                .zip(dec_bss.iter())
3056                .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3057        }
3058        assert_eq!(row_count, 200);
3059    }
3060
3061    #[test]
3062    fn test_read_incorrect_map_schema_file() {
3063        let testdata = arrow::util::test_util::parquet_test_data();
3064        // see https://github.com/apache/parquet-testing/pull/47
3065        let path = format!("{testdata}/incorrect_map_schema.parquet");
3066        let file = File::open(path).unwrap();
3067        let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
3068
3069        let batch = record_reader.next().unwrap().unwrap();
3070        assert_eq!(batch.num_rows(), 1);
3071
3072        let expected_schema = Schema::new(vec![Field::new(
3073            "my_map",
3074            ArrowDataType::Map(
3075                Arc::new(Field::new(
3076                    "key_value",
3077                    ArrowDataType::Struct(Fields::from(vec![
3078                        Field::new("key", ArrowDataType::Utf8, false),
3079                        Field::new("value", ArrowDataType::Utf8, true),
3080                    ])),
3081                    false,
3082                )),
3083                false,
3084            ),
3085            true,
3086        )]);
3087        assert_eq!(batch.schema().as_ref(), &expected_schema);
3088
3089        assert_eq!(batch.num_rows(), 1);
3090        assert_eq!(batch.column(0).null_count(), 0);
3091        assert_eq!(
3092            batch.column(0).as_map().keys().as_ref(),
3093            &StringArray::from(vec!["parent", "name"])
3094        );
3095        assert_eq!(
3096            batch.column(0).as_map().values().as_ref(),
3097            &StringArray::from(vec!["another", "report"])
3098        );
3099    }
3100
3101    #[test]
3102    fn test_read_dict_fixed_size_binary() {
3103        let schema = Arc::new(Schema::new(vec![Field::new(
3104            "a",
3105            ArrowDataType::Dictionary(
3106                Box::new(ArrowDataType::UInt8),
3107                Box::new(ArrowDataType::FixedSizeBinary(8)),
3108            ),
3109            true,
3110        )]));
3111        let keys = UInt8Array::from_iter_values(vec![0, 0, 1]);
3112        let values = FixedSizeBinaryArray::try_from_iter(
3113            vec![
3114                (0u8..8u8).collect::<Vec<u8>>(),
3115                (24u8..32u8).collect::<Vec<u8>>(),
3116            ]
3117            .into_iter(),
3118        )
3119        .unwrap();
3120        let arr = UInt8DictionaryArray::new(keys, Arc::new(values));
3121        let batch = RecordBatch::try_new(schema, vec![Arc::new(arr)]).unwrap();
3122
3123        let mut buffer = Vec::with_capacity(1024);
3124        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
3125        writer.write(&batch).unwrap();
3126        writer.close().unwrap();
3127        let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
3128            .unwrap()
3129            .collect::<Result<Vec<_>, _>>()
3130            .unwrap();
3131
3132        assert_eq!(read.len(), 1);
3133        assert_eq!(&batch, &read[0])
3134    }
3135
3136    #[test]
3137    fn test_read_nullable_structs_with_binary_dict_as_first_child_column() {
3138        // the `StructArrayReader` will check the definition and repetition levels of the first
3139        // child column in the struct to determine nullability for the struct. If the first
3140        // column's is being read by `ByteArrayDictionaryReader` we need to ensure that the
3141        // nullability is interpreted  correctly from the rep/def level buffers managed by the
3142        // buffers managed by this array reader.
3143
3144        let struct_fields = Fields::from(vec![
3145            Field::new(
3146                "city",
3147                ArrowDataType::Dictionary(
3148                    Box::new(ArrowDataType::UInt8),
3149                    Box::new(ArrowDataType::Utf8),
3150                ),
3151                true,
3152            ),
3153            Field::new("name", ArrowDataType::Utf8, true),
3154        ]);
3155        let schema = Arc::new(Schema::new(vec![Field::new(
3156            "items",
3157            ArrowDataType::Struct(struct_fields.clone()),
3158            true,
3159        )]));
3160
3161        let items_arr = StructArray::new(
3162            struct_fields,
3163            vec![
3164                Arc::new(DictionaryArray::new(
3165                    UInt8Array::from_iter_values(vec![0, 1, 1, 0, 2]),
3166                    Arc::new(StringArray::from_iter_values(vec![
3167                        "quebec",
3168                        "fredericton",
3169                        "halifax",
3170                    ])),
3171                )),
3172                Arc::new(StringArray::from_iter_values(vec![
3173                    "albert", "terry", "lance", "", "tim",
3174                ])),
3175            ],
3176            Some(NullBuffer::from_iter(vec![true, true, true, false, true])),
3177        );
3178
3179        let batch = RecordBatch::try_new(schema, vec![Arc::new(items_arr)]).unwrap();
3180        let mut buffer = Vec::with_capacity(1024);
3181        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
3182        writer.write(&batch).unwrap();
3183        writer.close().unwrap();
3184        let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 8)
3185            .unwrap()
3186            .collect::<Result<Vec<_>, _>>()
3187            .unwrap();
3188
3189        assert_eq!(read.len(), 1);
3190        assert_eq!(&batch, &read[0])
3191    }
3192
3193    /// Parameters for single_column_reader_test
3194    #[derive(Clone)]
3195    struct TestOptions {
3196        /// Number of row group to write to parquet (row group size =
3197        /// num_row_groups / num_rows)
3198        num_row_groups: usize,
3199        /// Total number of rows per row group
3200        num_rows: usize,
3201        /// Size of batches to read back
3202        record_batch_size: usize,
3203        /// Percentage of nulls in column or None if required
3204        null_percent: Option<usize>,
3205        /// Set write batch size
3206        ///
3207        /// This is the number of rows that are written at once to a page and
3208        /// therefore acts as a bound on the page granularity of a row group
3209        write_batch_size: usize,
3210        /// Maximum size of page in bytes
3211        max_data_page_size: usize,
3212        /// Maximum size of dictionary page in bytes
3213        max_dict_page_size: usize,
3214        /// Writer version
3215        writer_version: WriterVersion,
3216        /// Enabled statistics
3217        enabled_statistics: EnabledStatistics,
3218        /// Encoding
3219        encoding: Encoding,
3220        /// row selections and total selected row count
3221        row_selections: Option<(RowSelection, usize)>,
3222        /// row filter
3223        row_filter: Option<Vec<bool>>,
3224        /// limit
3225        limit: Option<usize>,
3226        /// offset
3227        offset: Option<usize>,
3228    }
3229
3230    /// Manually implement this to avoid printing entire contents of row_selections and row_filter
3231    impl std::fmt::Debug for TestOptions {
3232        fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
3233            f.debug_struct("TestOptions")
3234                .field("num_row_groups", &self.num_row_groups)
3235                .field("num_rows", &self.num_rows)
3236                .field("record_batch_size", &self.record_batch_size)
3237                .field("null_percent", &self.null_percent)
3238                .field("write_batch_size", &self.write_batch_size)
3239                .field("max_data_page_size", &self.max_data_page_size)
3240                .field("max_dict_page_size", &self.max_dict_page_size)
3241                .field("writer_version", &self.writer_version)
3242                .field("enabled_statistics", &self.enabled_statistics)
3243                .field("encoding", &self.encoding)
3244                .field("row_selections", &self.row_selections.is_some())
3245                .field("row_filter", &self.row_filter.is_some())
3246                .field("limit", &self.limit)
3247                .field("offset", &self.offset)
3248                .finish()
3249        }
3250    }
3251
3252    impl Default for TestOptions {
3253        fn default() -> Self {
3254            Self {
3255                num_row_groups: 2,
3256                num_rows: 100,
3257                record_batch_size: 15,
3258                null_percent: None,
3259                write_batch_size: 64,
3260                max_data_page_size: 1024 * 1024,
3261                max_dict_page_size: 1024 * 1024,
3262                writer_version: WriterVersion::PARQUET_1_0,
3263                enabled_statistics: EnabledStatistics::Page,
3264                encoding: Encoding::PLAIN,
3265                row_selections: None,
3266                row_filter: None,
3267                limit: None,
3268                offset: None,
3269            }
3270        }
3271    }
3272
3273    impl TestOptions {
3274        fn new(num_row_groups: usize, num_rows: usize, record_batch_size: usize) -> Self {
3275            Self {
3276                num_row_groups,
3277                num_rows,
3278                record_batch_size,
3279                ..Default::default()
3280            }
3281        }
3282
3283        fn with_null_percent(self, null_percent: usize) -> Self {
3284            Self {
3285                null_percent: Some(null_percent),
3286                ..self
3287            }
3288        }
3289
3290        fn with_max_data_page_size(self, max_data_page_size: usize) -> Self {
3291            Self {
3292                max_data_page_size,
3293                ..self
3294            }
3295        }
3296
3297        fn with_max_dict_page_size(self, max_dict_page_size: usize) -> Self {
3298            Self {
3299                max_dict_page_size,
3300                ..self
3301            }
3302        }
3303
3304        fn with_enabled_statistics(self, enabled_statistics: EnabledStatistics) -> Self {
3305            Self {
3306                enabled_statistics,
3307                ..self
3308            }
3309        }
3310
3311        fn with_row_selections(self) -> Self {
3312            assert!(self.row_filter.is_none(), "Must set row selection first");
3313
3314            let mut rng = rng();
3315            let step = rng.random_range(self.record_batch_size..self.num_rows);
3316            let row_selections = create_test_selection(
3317                step,
3318                self.num_row_groups * self.num_rows,
3319                rng.random::<bool>(),
3320            );
3321            Self {
3322                row_selections: Some(row_selections),
3323                ..self
3324            }
3325        }
3326
3327        fn with_row_filter(self) -> Self {
3328            let row_count = match &self.row_selections {
3329                Some((_, count)) => *count,
3330                None => self.num_row_groups * self.num_rows,
3331            };
3332
3333            let mut rng = rng();
3334            Self {
3335                row_filter: Some((0..row_count).map(|_| rng.random_bool(0.9)).collect()),
3336                ..self
3337            }
3338        }
3339
3340        fn with_limit(self, limit: usize) -> Self {
3341            Self {
3342                limit: Some(limit),
3343                ..self
3344            }
3345        }
3346
3347        fn with_offset(self, offset: usize) -> Self {
3348            Self {
3349                offset: Some(offset),
3350                ..self
3351            }
3352        }
3353
3354        fn writer_props(&self) -> WriterProperties {
3355            let builder = WriterProperties::builder()
3356                .set_data_page_size_limit(self.max_data_page_size)
3357                .set_write_batch_size(self.write_batch_size)
3358                .set_writer_version(self.writer_version)
3359                .set_statistics_enabled(self.enabled_statistics);
3360
3361            let builder = match self.encoding {
3362                Encoding::RLE_DICTIONARY | Encoding::PLAIN_DICTIONARY => builder
3363                    .set_dictionary_enabled(true)
3364                    .set_dictionary_page_size_limit(self.max_dict_page_size),
3365                _ => builder
3366                    .set_dictionary_enabled(false)
3367                    .set_encoding(self.encoding),
3368            };
3369
3370            builder.build()
3371        }
3372    }
3373
3374    /// Create a parquet file and then read it using
3375    /// `ParquetFileArrowReader` using a standard set of parameters
3376    /// `opts`.
3377    ///
3378    /// `rand_max` represents the maximum size of value to pass to to
3379    /// value generator
3380    fn run_single_column_reader_tests<T, F, G>(
3381        rand_max: i32,
3382        converted_type: ConvertedType,
3383        arrow_type: Option<ArrowDataType>,
3384        converter: F,
3385        encodings: &[Encoding],
3386    ) where
3387        T: DataType,
3388        G: RandGen<T>,
3389        F: Fn(&[Option<T::T>]) -> ArrayRef,
3390    {
3391        let all_options = vec![
3392            // choose record_batch_batch (15) so batches cross row
3393            // group boundaries (50 rows in 2 row groups) cases.
3394            TestOptions::new(2, 100, 15),
3395            // choose record_batch_batch (5) so batches sometime fall
3396            // on row group boundaries and (25 rows in 3 row groups
3397            // --> row groups of 10, 10, and 5). Tests buffer
3398            // refilling edge cases.
3399            TestOptions::new(3, 25, 5),
3400            // Choose record_batch_size (25) so all batches fall
3401            // exactly on row group boundary (25). Tests buffer
3402            // refilling edge cases.
3403            TestOptions::new(4, 100, 25),
3404            // Set maximum page size so row groups have multiple pages
3405            TestOptions::new(3, 256, 73).with_max_data_page_size(128),
3406            // Set small dictionary page size to test dictionary fallback
3407            TestOptions::new(3, 256, 57).with_max_dict_page_size(128),
3408            // Test optional but with no nulls
3409            TestOptions::new(2, 256, 127).with_null_percent(0),
3410            // Test optional with nulls
3411            TestOptions::new(2, 256, 93).with_null_percent(25),
3412            // Test with limit of 0
3413            TestOptions::new(4, 100, 25).with_limit(0),
3414            // Test with limit of 50
3415            TestOptions::new(4, 100, 25).with_limit(50),
3416            // Test with limit equal to number of rows
3417            TestOptions::new(4, 100, 25).with_limit(10),
3418            // Test with limit larger than number of rows
3419            TestOptions::new(4, 100, 25).with_limit(101),
3420            // Test with limit + offset equal to number of rows
3421            TestOptions::new(4, 100, 25).with_offset(30).with_limit(20),
3422            // Test with limit + offset equal to number of rows
3423            TestOptions::new(4, 100, 25).with_offset(20).with_limit(80),
3424            // Test with limit + offset larger than number of rows
3425            TestOptions::new(4, 100, 25).with_offset(20).with_limit(81),
3426            // Test with no page-level statistics
3427            TestOptions::new(2, 256, 91)
3428                .with_null_percent(25)
3429                .with_enabled_statistics(EnabledStatistics::Chunk),
3430            // Test with no statistics
3431            TestOptions::new(2, 256, 91)
3432                .with_null_percent(25)
3433                .with_enabled_statistics(EnabledStatistics::None),
3434            // Test with all null
3435            TestOptions::new(2, 128, 91)
3436                .with_null_percent(100)
3437                .with_enabled_statistics(EnabledStatistics::None),
3438            // Test skip
3439
3440            // choose record_batch_batch (15) so batches cross row
3441            // group boundaries (50 rows in 2 row groups) cases.
3442            TestOptions::new(2, 100, 15).with_row_selections(),
3443            // choose record_batch_batch (5) so batches sometime fall
3444            // on row group boundaries and (25 rows in 3 row groups
3445            // --> row groups of 10, 10, and 5). Tests buffer
3446            // refilling edge cases.
3447            TestOptions::new(3, 25, 5).with_row_selections(),
3448            // Choose record_batch_size (25) so all batches fall
3449            // exactly on row group boundary (25). Tests buffer
3450            // refilling edge cases.
3451            TestOptions::new(4, 100, 25).with_row_selections(),
3452            // Set maximum page size so row groups have multiple pages
3453            TestOptions::new(3, 256, 73)
3454                .with_max_data_page_size(128)
3455                .with_row_selections(),
3456            // Set small dictionary page size to test dictionary fallback
3457            TestOptions::new(3, 256, 57)
3458                .with_max_dict_page_size(128)
3459                .with_row_selections(),
3460            // Test optional but with no nulls
3461            TestOptions::new(2, 256, 127)
3462                .with_null_percent(0)
3463                .with_row_selections(),
3464            // Test optional with nulls
3465            TestOptions::new(2, 256, 93)
3466                .with_null_percent(25)
3467                .with_row_selections(),
3468            // Test optional with nulls
3469            TestOptions::new(2, 256, 93)
3470                .with_null_percent(25)
3471                .with_row_selections()
3472                .with_limit(10),
3473            // Test optional with nulls
3474            TestOptions::new(2, 256, 93)
3475                .with_null_percent(25)
3476                .with_row_selections()
3477                .with_offset(20)
3478                .with_limit(10),
3479            // Test filter
3480
3481            // Test with row filter
3482            TestOptions::new(4, 100, 25).with_row_filter(),
3483            // Test with row selection and row filter
3484            TestOptions::new(4, 100, 25)
3485                .with_row_selections()
3486                .with_row_filter(),
3487            // Test with nulls and row filter
3488            TestOptions::new(2, 256, 93)
3489                .with_null_percent(25)
3490                .with_max_data_page_size(10)
3491                .with_row_filter(),
3492            // Test with nulls and row filter and small pages
3493            TestOptions::new(2, 256, 93)
3494                .with_null_percent(25)
3495                .with_max_data_page_size(10)
3496                .with_row_selections()
3497                .with_row_filter(),
3498            // Test with row selection and no offset index and small pages
3499            TestOptions::new(2, 256, 93)
3500                .with_enabled_statistics(EnabledStatistics::None)
3501                .with_max_data_page_size(10)
3502                .with_row_selections(),
3503        ];
3504
3505        all_options.into_iter().for_each(|opts| {
3506            for writer_version in [WriterVersion::PARQUET_1_0, WriterVersion::PARQUET_2_0] {
3507                for encoding in encodings {
3508                    let opts = TestOptions {
3509                        writer_version,
3510                        encoding: *encoding,
3511                        ..opts.clone()
3512                    };
3513
3514                    single_column_reader_test::<T, _, G>(
3515                        opts,
3516                        rand_max,
3517                        converted_type,
3518                        arrow_type.clone(),
3519                        &converter,
3520                    )
3521                }
3522            }
3523        });
3524    }
3525
3526    /// Create a parquet file and then read it using
3527    /// `ParquetFileArrowReader` using the parameters described in
3528    /// `opts`.
3529    fn single_column_reader_test<T, F, G>(
3530        opts: TestOptions,
3531        rand_max: i32,
3532        converted_type: ConvertedType,
3533        arrow_type: Option<ArrowDataType>,
3534        converter: F,
3535    ) where
3536        T: DataType,
3537        G: RandGen<T>,
3538        F: Fn(&[Option<T::T>]) -> ArrayRef,
3539    {
3540        // Print out options to facilitate debugging failures on CI
3541        println!(
3542            "Running type {:?} single_column_reader_test ConvertedType::{}/ArrowType::{:?} with Options: {:?}",
3543            T::get_physical_type(),
3544            converted_type,
3545            arrow_type,
3546            opts
3547        );
3548
3549        //according to null_percent generate def_levels
3550        let (repetition, def_levels) = match opts.null_percent.as_ref() {
3551            Some(null_percent) => {
3552                let mut rng = rng();
3553
3554                let def_levels: Vec<Vec<i16>> = (0..opts.num_row_groups)
3555                    .map(|_| {
3556                        std::iter::from_fn(|| {
3557                            Some((rng.next_u32() as usize % 100 >= *null_percent) as i16)
3558                        })
3559                        .take(opts.num_rows)
3560                        .collect()
3561                    })
3562                    .collect();
3563                (Repetition::OPTIONAL, Some(def_levels))
3564            }
3565            None => (Repetition::REQUIRED, None),
3566        };
3567
3568        //generate random table data
3569        let values: Vec<Vec<T::T>> = (0..opts.num_row_groups)
3570            .map(|idx| {
3571                let null_count = match def_levels.as_ref() {
3572                    Some(d) => d[idx].iter().filter(|x| **x == 0).count(),
3573                    None => 0,
3574                };
3575                G::gen_vec(rand_max, opts.num_rows - null_count)
3576            })
3577            .collect();
3578
3579        let len = match T::get_physical_type() {
3580            crate::basic::Type::FIXED_LEN_BYTE_ARRAY => rand_max,
3581            crate::basic::Type::INT96 => 12,
3582            _ => -1,
3583        };
3584
3585        let fields = vec![Arc::new(
3586            Type::primitive_type_builder("leaf", T::get_physical_type())
3587                .with_repetition(repetition)
3588                .with_converted_type(converted_type)
3589                .with_length(len)
3590                .build()
3591                .unwrap(),
3592        )];
3593
3594        let schema = Arc::new(
3595            Type::group_type_builder("test_schema")
3596                .with_fields(fields)
3597                .build()
3598                .unwrap(),
3599        );
3600
3601        let arrow_field = arrow_type.map(|t| Field::new("leaf", t, false));
3602
3603        let mut file = tempfile::tempfile().unwrap();
3604
3605        generate_single_column_file_with_data::<T>(
3606            &values,
3607            def_levels.as_ref(),
3608            file.try_clone().unwrap(), // Cannot use &mut File (#1163)
3609            schema,
3610            arrow_field,
3611            &opts,
3612        )
3613        .unwrap();
3614
3615        file.rewind().unwrap();
3616
3617        let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::from(
3618            opts.enabled_statistics == EnabledStatistics::Page,
3619        ));
3620
3621        let mut builder =
3622            ParquetRecordBatchReaderBuilder::try_new_with_options(file, options).unwrap();
3623
3624        let expected_data = match opts.row_selections {
3625            Some((selections, row_count)) => {
3626                let mut without_skip_data = gen_expected_data::<T>(def_levels.as_ref(), &values);
3627
3628                let mut skip_data: Vec<Option<T::T>> = vec![];
3629                let dequeue: VecDeque<RowSelector> = selections.clone().into();
3630                for select in dequeue {
3631                    if select.skip {
3632                        without_skip_data.drain(0..select.row_count);
3633                    } else {
3634                        skip_data.extend(without_skip_data.drain(0..select.row_count));
3635                    }
3636                }
3637                builder = builder.with_row_selection(selections);
3638
3639                assert_eq!(skip_data.len(), row_count);
3640                skip_data
3641            }
3642            None => {
3643                //get flatten table data
3644                let expected_data = gen_expected_data::<T>(def_levels.as_ref(), &values);
3645                assert_eq!(expected_data.len(), opts.num_rows * opts.num_row_groups);
3646                expected_data
3647            }
3648        };
3649
3650        let mut expected_data = match opts.row_filter {
3651            Some(filter) => {
3652                let expected_data = expected_data
3653                    .into_iter()
3654                    .zip(filter.iter())
3655                    .filter_map(|(d, f)| f.then(|| d))
3656                    .collect();
3657
3658                let mut filter_offset = 0;
3659                let filter = RowFilter::new(vec![Box::new(ArrowPredicateFn::new(
3660                    ProjectionMask::all(),
3661                    move |b| {
3662                        let array = BooleanArray::from_iter(
3663                            filter
3664                                .iter()
3665                                .skip(filter_offset)
3666                                .take(b.num_rows())
3667                                .map(|x| Some(*x)),
3668                        );
3669                        filter_offset += b.num_rows();
3670                        Ok(array)
3671                    },
3672                ))]);
3673
3674                builder = builder.with_row_filter(filter);
3675                expected_data
3676            }
3677            None => expected_data,
3678        };
3679
3680        if let Some(offset) = opts.offset {
3681            builder = builder.with_offset(offset);
3682            expected_data = expected_data.into_iter().skip(offset).collect();
3683        }
3684
3685        if let Some(limit) = opts.limit {
3686            builder = builder.with_limit(limit);
3687            expected_data = expected_data.into_iter().take(limit).collect();
3688        }
3689
3690        let mut record_reader = builder
3691            .with_batch_size(opts.record_batch_size)
3692            .build()
3693            .unwrap();
3694
3695        let mut total_read = 0;
3696        loop {
3697            let maybe_batch = record_reader.next();
3698            if total_read < expected_data.len() {
3699                let end = min(total_read + opts.record_batch_size, expected_data.len());
3700                let batch = maybe_batch.unwrap().unwrap();
3701                assert_eq!(end - total_read, batch.num_rows());
3702
3703                let a = converter(&expected_data[total_read..end]);
3704                let b = batch.column(0);
3705
3706                assert_eq!(a.data_type(), b.data_type());
3707                assert_eq!(a.to_data(), b.to_data());
3708                assert_eq!(
3709                    a.as_any().type_id(),
3710                    b.as_any().type_id(),
3711                    "incorrect type ids"
3712                );
3713
3714                total_read = end;
3715            } else {
3716                assert!(maybe_batch.is_none());
3717                break;
3718            }
3719        }
3720    }
3721
3722    fn gen_expected_data<T: DataType>(
3723        def_levels: Option<&Vec<Vec<i16>>>,
3724        values: &[Vec<T::T>],
3725    ) -> Vec<Option<T::T>> {
3726        let data: Vec<Option<T::T>> = match def_levels {
3727            Some(levels) => {
3728                let mut values_iter = values.iter().flatten();
3729                levels
3730                    .iter()
3731                    .flatten()
3732                    .map(|d| match d {
3733                        1 => Some(values_iter.next().cloned().unwrap()),
3734                        0 => None,
3735                        _ => unreachable!(),
3736                    })
3737                    .collect()
3738            }
3739            None => values.iter().flatten().cloned().map(Some).collect(),
3740        };
3741        data
3742    }
3743
3744    fn generate_single_column_file_with_data<T: DataType>(
3745        values: &[Vec<T::T>],
3746        def_levels: Option<&Vec<Vec<i16>>>,
3747        file: File,
3748        schema: TypePtr,
3749        field: Option<Field>,
3750        opts: &TestOptions,
3751    ) -> Result<ParquetMetaData> {
3752        let mut writer_props = opts.writer_props();
3753        if let Some(field) = field {
3754            let arrow_schema = Schema::new(vec![field]);
3755            add_encoded_arrow_schema_to_metadata(&arrow_schema, &mut writer_props);
3756        }
3757
3758        let mut writer = SerializedFileWriter::new(file, schema, Arc::new(writer_props))?;
3759
3760        for (idx, v) in values.iter().enumerate() {
3761            let def_levels = def_levels.map(|d| d[idx].as_slice());
3762            let mut row_group_writer = writer.next_row_group()?;
3763            {
3764                let mut column_writer = row_group_writer
3765                    .next_column()?
3766                    .expect("Column writer is none!");
3767
3768                column_writer
3769                    .typed::<T>()
3770                    .write_batch(v, def_levels, None)?;
3771
3772                column_writer.close()?;
3773            }
3774            row_group_writer.close()?;
3775        }
3776
3777        writer.close()
3778    }
3779
3780    fn get_test_file(file_name: &str) -> File {
3781        let path = PathBuf::from(arrow::util::test_util::arrow_test_data()).join(file_name);
3782
3783        File::open(path.as_path()).expect("File not found!")
3784    }
3785
3786    #[cfg_attr(miri, ignore)] // calls native Zstd code unsupported by Miri
3787    #[test]
3788    fn test_read_structs() {
3789        // This particular test file has columns of struct types where there is
3790        // a column that has the same name as one of the struct fields
3791        // (see: ARROW-11452)
3792        let testdata = arrow::util::test_util::parquet_test_data();
3793        let path = format!("{testdata}/nested_structs.rust.parquet");
3794        let file = File::open(&path).unwrap();
3795        let record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
3796
3797        for batch in record_batch_reader {
3798            batch.unwrap();
3799        }
3800
3801        let file = File::open(&path).unwrap();
3802        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
3803
3804        let mask = ProjectionMask::leaves(builder.parquet_schema(), [3, 8, 10]);
3805        let projected_reader = builder
3806            .with_projection(mask)
3807            .with_batch_size(60)
3808            .build()
3809            .unwrap();
3810
3811        let expected_schema = Schema::new(vec![
3812            Field::new(
3813                "roll_num",
3814                ArrowDataType::Struct(Fields::from(vec![Field::new(
3815                    "count",
3816                    ArrowDataType::UInt64,
3817                    false,
3818                )])),
3819                false,
3820            ),
3821            Field::new(
3822                "PC_CUR",
3823                ArrowDataType::Struct(Fields::from(vec![
3824                    Field::new("mean", ArrowDataType::Int64, false),
3825                    Field::new("sum", ArrowDataType::Int64, false),
3826                ])),
3827                false,
3828            ),
3829        ]);
3830
3831        // Tests for #1652 and #1654
3832        assert_eq!(&expected_schema, projected_reader.schema().as_ref());
3833
3834        for batch in projected_reader {
3835            let batch = batch.unwrap();
3836            assert_eq!(batch.schema().as_ref(), &expected_schema);
3837        }
3838    }
3839
3840    #[cfg_attr(miri, ignore)] // calls native Zstd code unsupported by Miri
3841    #[test]
3842    // same as test_read_structs but constructs projection mask via column names
3843    fn test_read_structs_by_name() {
3844        let testdata = arrow::util::test_util::parquet_test_data();
3845        let path = format!("{testdata}/nested_structs.rust.parquet");
3846        let file = File::open(&path).unwrap();
3847        let record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
3848
3849        for batch in record_batch_reader {
3850            batch.unwrap();
3851        }
3852
3853        let file = File::open(&path).unwrap();
3854        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
3855
3856        let mask = ProjectionMask::columns(
3857            builder.parquet_schema(),
3858            ["roll_num.count", "PC_CUR.mean", "PC_CUR.sum"],
3859        );
3860        let projected_reader = builder
3861            .with_projection(mask)
3862            .with_batch_size(60)
3863            .build()
3864            .unwrap();
3865
3866        let expected_schema = Schema::new(vec![
3867            Field::new(
3868                "roll_num",
3869                ArrowDataType::Struct(Fields::from(vec![Field::new(
3870                    "count",
3871                    ArrowDataType::UInt64,
3872                    false,
3873                )])),
3874                false,
3875            ),
3876            Field::new(
3877                "PC_CUR",
3878                ArrowDataType::Struct(Fields::from(vec![
3879                    Field::new("mean", ArrowDataType::Int64, false),
3880                    Field::new("sum", ArrowDataType::Int64, false),
3881                ])),
3882                false,
3883            ),
3884        ]);
3885
3886        assert_eq!(&expected_schema, projected_reader.schema().as_ref());
3887
3888        for batch in projected_reader {
3889            let batch = batch.unwrap();
3890            assert_eq!(batch.schema().as_ref(), &expected_schema);
3891        }
3892    }
3893
3894    #[test]
3895    fn test_read_maps() {
3896        let testdata = arrow::util::test_util::parquet_test_data();
3897        let path = format!("{testdata}/nested_maps.snappy.parquet");
3898        let file = File::open(path).unwrap();
3899        let record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
3900
3901        for batch in record_batch_reader {
3902            batch.unwrap();
3903        }
3904    }
3905
3906    // test that we can handle the UNKNOWN logical type annotation on any physical type
3907    #[test]
3908    fn test_unknown_logical_type() {
3909        let message_type = "message uk {
3910            OPTIONAL INT32 uki32 (UNKNOWN);
3911            OPTIONAL INT64 uki64 (UNKNOWN);
3912            OPTIONAL INT96 uki96 (UNKNOWN);
3913            OPTIONAL BOOLEAN ukbool (UNKNOWN);
3914            OPTIONAL FLOAT ukfloat (UNKNOWN);
3915            OPTIONAL DOUBLE ukdbl (UNKNOWN);
3916            OPTIONAL BYTE_ARRAY ukbytes (UNKNOWN);
3917            OPTIONAL FIXED_LEN_BYTE_ARRAY(10) ukflba (UNKNOWN);
3918        }";
3919
3920        let schema = Arc::new(parse_message_type(message_type).unwrap());
3921        let file = tempfile::tempfile().unwrap();
3922
3923        let mut writer =
3924            SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default())
3925                .unwrap();
3926
3927        let mut row_group_writer = writer.next_row_group().unwrap();
3928
3929        fn write_nulls<T: DataType>(row_group_writer: &mut SerializedRowGroupWriter<'_, File>) {
3930            let mut column_writer = row_group_writer.next_column().unwrap().unwrap();
3931            // write out a bunch of nulls
3932            column_writer
3933                .typed::<T>()
3934                .write_batch(&[], Some(&[0, 0, 0, 0]), None)
3935                .unwrap();
3936            column_writer.close().unwrap();
3937        }
3938
3939        // INT32
3940        write_nulls::<Int32Type>(&mut row_group_writer);
3941
3942        // INT64
3943        write_nulls::<Int64Type>(&mut row_group_writer);
3944
3945        // INT96
3946        write_nulls::<Int96Type>(&mut row_group_writer);
3947
3948        // BOOLEAN
3949        write_nulls::<BoolType>(&mut row_group_writer);
3950
3951        // FLOAT
3952        write_nulls::<FloatType>(&mut row_group_writer);
3953
3954        // DOUBLE
3955        write_nulls::<DoubleType>(&mut row_group_writer);
3956
3957        // BYTE_ARRAY
3958        write_nulls::<ByteArrayType>(&mut row_group_writer);
3959
3960        // FIXED_LEN_BYTE_ARRAY
3961        write_nulls::<FixedLenByteArrayType>(&mut row_group_writer);
3962
3963        row_group_writer.close().unwrap();
3964
3965        writer.close().unwrap();
3966
3967        let mut reader = ParquetRecordBatchReader::try_new(file, 4).unwrap();
3968        let batch = reader.next().unwrap().unwrap();
3969
3970        for col in batch.columns() {
3971            assert_eq!(col.len(), 4);
3972            assert_eq!(col.logical_null_count(), 4);
3973            assert_eq!(*col.data_type(), ArrowDataType::Null);
3974        }
3975    }
3976
3977    #[test]
3978    fn test_nested_nullability() {
3979        let message_type = "message nested {
3980          OPTIONAL Group group {
3981            REQUIRED INT32 leaf;
3982          }
3983        }";
3984
3985        let file = tempfile::tempfile().unwrap();
3986        let schema = Arc::new(parse_message_type(message_type).unwrap());
3987
3988        {
3989            // Write using low-level parquet API (#1167)
3990            let mut writer =
3991                SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default())
3992                    .unwrap();
3993
3994            {
3995                let mut row_group_writer = writer.next_row_group().unwrap();
3996                let mut column_writer = row_group_writer.next_column().unwrap().unwrap();
3997
3998                column_writer
3999                    .typed::<Int32Type>()
4000                    .write_batch(&[34, 76], Some(&[0, 1, 0, 1]), None)
4001                    .unwrap();
4002
4003                column_writer.close().unwrap();
4004                row_group_writer.close().unwrap();
4005            }
4006
4007            writer.close().unwrap();
4008        }
4009
4010        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
4011        let mask = ProjectionMask::leaves(builder.parquet_schema(), [0]);
4012
4013        let reader = builder.with_projection(mask).build().unwrap();
4014
4015        let expected_schema = Schema::new(vec![Field::new(
4016            "group",
4017            ArrowDataType::Struct(vec![Field::new("leaf", ArrowDataType::Int32, false)].into()),
4018            true,
4019        )]);
4020
4021        let batch = reader.into_iter().next().unwrap().unwrap();
4022        assert_eq!(batch.schema().as_ref(), &expected_schema);
4023        assert_eq!(batch.num_rows(), 4);
4024        assert_eq!(batch.column(0).null_count(), 2);
4025    }
4026
4027    #[test]
4028    fn test_dictionary_preservation() {
4029        let fields = vec![Arc::new(
4030            Type::primitive_type_builder("leaf", PhysicalType::BYTE_ARRAY)
4031                .with_repetition(Repetition::OPTIONAL)
4032                .with_converted_type(ConvertedType::UTF8)
4033                .build()
4034                .unwrap(),
4035        )];
4036
4037        let schema = Arc::new(
4038            Type::group_type_builder("test_schema")
4039                .with_fields(fields)
4040                .build()
4041                .unwrap(),
4042        );
4043
4044        let dict_type = ArrowDataType::Dictionary(
4045            Box::new(ArrowDataType::Int32),
4046            Box::new(ArrowDataType::Utf8),
4047        );
4048
4049        let arrow_field = Field::new("leaf", dict_type, true);
4050
4051        let mut file = tempfile::tempfile().unwrap();
4052
4053        let values = vec![
4054            vec![
4055                ByteArray::from("hello"),
4056                ByteArray::from("a"),
4057                ByteArray::from("b"),
4058                ByteArray::from("d"),
4059            ],
4060            vec![
4061                ByteArray::from("c"),
4062                ByteArray::from("a"),
4063                ByteArray::from("b"),
4064            ],
4065        ];
4066
4067        let def_levels = vec![
4068            vec![1, 0, 0, 1, 0, 0, 1, 1],
4069            vec![0, 0, 1, 1, 0, 0, 1, 0, 0],
4070        ];
4071
4072        let opts = TestOptions {
4073            encoding: Encoding::RLE_DICTIONARY,
4074            ..Default::default()
4075        };
4076
4077        generate_single_column_file_with_data::<ByteArrayType>(
4078            &values,
4079            Some(&def_levels),
4080            file.try_clone().unwrap(), // Cannot use &mut File (#1163)
4081            schema,
4082            Some(arrow_field),
4083            &opts,
4084        )
4085        .unwrap();
4086
4087        file.rewind().unwrap();
4088
4089        let record_reader = ParquetRecordBatchReader::try_new(file, 3).unwrap();
4090
4091        let batches = record_reader
4092            .collect::<Result<Vec<RecordBatch>, _>>()
4093            .unwrap();
4094
4095        assert_eq!(batches.len(), 6);
4096        assert!(batches.iter().all(|x| x.num_columns() == 1));
4097
4098        let row_counts = batches
4099            .iter()
4100            .map(|x| (x.num_rows(), x.column(0).null_count()))
4101            .collect::<Vec<_>>();
4102
4103        assert_eq!(
4104            row_counts,
4105            vec![(3, 2), (3, 2), (3, 1), (3, 1), (3, 2), (2, 2)]
4106        );
4107
4108        let get_dict = |batch: &RecordBatch| batch.column(0).to_data().child_data()[0].clone();
4109
4110        // First and second batch in same row group -> same dictionary
4111        assert_eq!(get_dict(&batches[0]), get_dict(&batches[1]));
4112        // Third batch spans row group -> computed dictionary
4113        assert_ne!(get_dict(&batches[1]), get_dict(&batches[2]));
4114        assert_ne!(get_dict(&batches[2]), get_dict(&batches[3]));
4115        // Fourth, fifth and sixth from same row group -> same dictionary
4116        assert_eq!(get_dict(&batches[3]), get_dict(&batches[4]));
4117        assert_eq!(get_dict(&batches[4]), get_dict(&batches[5]));
4118    }
4119
4120    #[test]
4121    fn test_read_null_list() {
4122        let testdata = arrow::util::test_util::parquet_test_data();
4123        let path = format!("{testdata}/null_list.parquet");
4124        let file = File::open(path).unwrap();
4125        let mut record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
4126
4127        let batch = record_batch_reader.next().unwrap().unwrap();
4128        assert_eq!(batch.num_rows(), 1);
4129        assert_eq!(batch.num_columns(), 1);
4130        assert_eq!(batch.column(0).len(), 1);
4131
4132        let list = batch
4133            .column(0)
4134            .as_any()
4135            .downcast_ref::<ListArray>()
4136            .unwrap();
4137        assert_eq!(list.len(), 1);
4138        assert!(list.is_valid(0));
4139
4140        let val = list.value(0);
4141        assert_eq!(val.len(), 0);
4142    }
4143
4144    #[test]
4145    fn test_null_schema_inference() {
4146        let testdata = arrow::util::test_util::parquet_test_data();
4147        let path = format!("{testdata}/null_list.parquet");
4148        let file = File::open(path).unwrap();
4149
4150        let arrow_field = Field::new(
4151            "emptylist",
4152            ArrowDataType::List(Arc::new(Field::new_list_field(ArrowDataType::Null, true))),
4153            true,
4154        );
4155
4156        let options = ArrowReaderOptions::new().with_skip_arrow_metadata(true);
4157        let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options).unwrap();
4158        let schema = builder.schema();
4159        assert_eq!(schema.fields().len(), 1);
4160        assert_eq!(schema.field(0), &arrow_field);
4161    }
4162
4163    #[test]
4164    fn test_skip_metadata() {
4165        let col = Arc::new(TimestampNanosecondArray::from_iter_values(vec![0, 1, 2]));
4166        let field = Field::new("col", col.data_type().clone(), true);
4167
4168        let schema_without_metadata = Arc::new(Schema::new(vec![field.clone()]));
4169
4170        let metadata = arrow_schema::Metadata::from([("key".to_string(), "value".to_string())]);
4171
4172        let schema_with_metadata = Arc::new(Schema::new(vec![field.with_metadata(metadata)]));
4173
4174        assert_ne!(schema_with_metadata, schema_without_metadata);
4175
4176        let batch =
4177            RecordBatch::try_new(schema_with_metadata.clone(), vec![col as ArrayRef]).unwrap();
4178
4179        let file = |version: WriterVersion| {
4180            let props = WriterProperties::builder()
4181                .set_writer_version(version)
4182                .build();
4183
4184            let file = tempfile().unwrap();
4185            let mut writer =
4186                ArrowWriter::try_new(file.try_clone().unwrap(), batch.schema(), Some(props))
4187                    .unwrap();
4188            writer.write(&batch).unwrap();
4189            writer.close().unwrap();
4190            file
4191        };
4192
4193        let skip_options = ArrowReaderOptions::new().with_skip_arrow_metadata(true);
4194
4195        let v1_reader = file(WriterVersion::PARQUET_1_0);
4196        let v2_reader = file(WriterVersion::PARQUET_2_0);
4197
4198        let arrow_reader =
4199            ParquetRecordBatchReader::try_new(v1_reader.try_clone().unwrap(), 1024).unwrap();
4200        assert_eq!(arrow_reader.schema(), schema_with_metadata);
4201
4202        let reader =
4203            ParquetRecordBatchReaderBuilder::try_new_with_options(v1_reader, skip_options.clone())
4204                .unwrap()
4205                .build()
4206                .unwrap();
4207        assert_eq!(reader.schema(), schema_without_metadata);
4208
4209        let arrow_reader =
4210            ParquetRecordBatchReader::try_new(v2_reader.try_clone().unwrap(), 1024).unwrap();
4211        assert_eq!(arrow_reader.schema(), schema_with_metadata);
4212
4213        let reader = ParquetRecordBatchReaderBuilder::try_new_with_options(v2_reader, skip_options)
4214            .unwrap()
4215            .build()
4216            .unwrap();
4217        assert_eq!(reader.schema(), schema_without_metadata);
4218    }
4219
4220    fn write_parquet_from_iter<I, F>(value: I) -> File
4221    where
4222        I: IntoIterator<Item = (F, ArrayRef)>,
4223        F: AsRef<str>,
4224    {
4225        let batch = RecordBatch::try_from_iter(value).unwrap();
4226        let file = tempfile().unwrap();
4227        let mut writer =
4228            ArrowWriter::try_new(file.try_clone().unwrap(), batch.schema().clone(), None).unwrap();
4229        writer.write(&batch).unwrap();
4230        writer.close().unwrap();
4231        file
4232    }
4233
4234    fn run_schema_test_with_error<I, F>(value: I, schema: SchemaRef, expected_error: &str)
4235    where
4236        I: IntoIterator<Item = (F, ArrayRef)>,
4237        F: AsRef<str>,
4238    {
4239        let file = write_parquet_from_iter(value);
4240        let options_with_schema = ArrowReaderOptions::new().with_schema(schema.clone());
4241        let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
4242            file.try_clone().unwrap(),
4243            options_with_schema,
4244        );
4245        assert_eq!(builder.err().unwrap().to_string(), expected_error);
4246    }
4247
4248    #[test]
4249    fn test_schema_too_few_columns() {
4250        run_schema_test_with_error(
4251            vec![
4252                ("int64", Arc::new(Int64Array::from(vec![0])) as ArrayRef),
4253                ("int32", Arc::new(Int32Array::from(vec![0])) as ArrayRef),
4254            ],
4255            Arc::new(Schema::new(vec![Field::new(
4256                "int64",
4257                ArrowDataType::Int64,
4258                false,
4259            )])),
4260            "Arrow: incompatible arrow schema, expected 2 struct fields got 1",
4261        );
4262    }
4263
4264    #[test]
4265    fn test_schema_too_many_columns() {
4266        run_schema_test_with_error(
4267            vec![("int64", Arc::new(Int64Array::from(vec![0])) as ArrayRef)],
4268            Arc::new(Schema::new(vec![
4269                Field::new("int64", ArrowDataType::Int64, false),
4270                Field::new("int32", ArrowDataType::Int32, false),
4271            ])),
4272            "Arrow: incompatible arrow schema, expected 1 struct fields got 2",
4273        );
4274    }
4275
4276    #[test]
4277    fn test_schema_mismatched_column_names() {
4278        run_schema_test_with_error(
4279            vec![("int64", Arc::new(Int64Array::from(vec![0])) as ArrayRef)],
4280            Arc::new(Schema::new(vec![Field::new(
4281                "other",
4282                ArrowDataType::Int64,
4283                false,
4284            )])),
4285            "Arrow: incompatible arrow schema, expected field named int64 got other",
4286        );
4287    }
4288
4289    #[test]
4290    fn test_schema_incompatible_columns() {
4291        run_schema_test_with_error(
4292            vec![
4293                (
4294                    "col1_invalid",
4295                    Arc::new(Int64Array::from(vec![0])) as ArrayRef,
4296                ),
4297                (
4298                    "col2_valid",
4299                    Arc::new(Int32Array::from(vec![0])) as ArrayRef,
4300                ),
4301                (
4302                    "col3_invalid",
4303                    Arc::new(Date64Array::from(vec![0])) as ArrayRef,
4304                ),
4305            ],
4306            Arc::new(Schema::new(vec![
4307                Field::new("col1_invalid", ArrowDataType::Int32, false),
4308                Field::new("col2_valid", ArrowDataType::Int32, false),
4309                Field::new("col3_invalid", ArrowDataType::Int32, false),
4310            ])),
4311            "Arrow: Incompatible supplied Arrow schema: data type mismatch for field col1_invalid: requested Int32 but found Int64, data type mismatch for field col3_invalid: requested Int32 but found Int64",
4312        );
4313    }
4314
4315    #[test]
4316    fn test_one_incompatible_nested_column() {
4317        let nested_fields = Fields::from(vec![
4318            Field::new("nested1_valid", ArrowDataType::Utf8, false),
4319            Field::new("nested1_invalid", ArrowDataType::Int64, false),
4320        ]);
4321        let nested = StructArray::try_new(
4322            nested_fields,
4323            vec![
4324                Arc::new(StringArray::from(vec!["a"])) as ArrayRef,
4325                Arc::new(Int64Array::from(vec![0])) as ArrayRef,
4326            ],
4327            None,
4328        )
4329        .expect("struct array");
4330        let supplied_nested_fields = Fields::from(vec![
4331            Field::new("nested1_valid", ArrowDataType::Utf8, false),
4332            Field::new("nested1_invalid", ArrowDataType::Int32, false),
4333        ]);
4334        run_schema_test_with_error(
4335            vec![
4336                ("col1", Arc::new(Int64Array::from(vec![0])) as ArrayRef),
4337                ("col2", Arc::new(Int32Array::from(vec![0])) as ArrayRef),
4338                ("nested", Arc::new(nested) as ArrayRef),
4339            ],
4340            Arc::new(Schema::new(vec![
4341                Field::new("col1", ArrowDataType::Int64, false),
4342                Field::new("col2", ArrowDataType::Int32, false),
4343                Field::new(
4344                    "nested",
4345                    ArrowDataType::Struct(supplied_nested_fields),
4346                    false,
4347                ),
4348            ])),
4349            "Arrow: Incompatible supplied Arrow schema: data type mismatch for field nested: \
4350            requested Struct(\"nested1_valid\": non-null Utf8, \"nested1_invalid\": non-null Int32) \
4351            but found Struct(\"nested1_valid\": non-null Utf8, \"nested1_invalid\": non-null Int64)",
4352        );
4353    }
4354
4355    /// Return parquet data with a single column of utf8 strings
4356    fn utf8_parquet() -> Bytes {
4357        let input = StringArray::from_iter_values(vec!["foo", "bar", "baz"]);
4358        let batch = RecordBatch::try_from_iter(vec![("column1", Arc::new(input) as _)]).unwrap();
4359        let props = None;
4360        // write parquet file with non nullable strings
4361        let mut parquet_data = vec![];
4362        let mut writer = ArrowWriter::try_new(&mut parquet_data, batch.schema(), props).unwrap();
4363        writer.write(&batch).unwrap();
4364        writer.close().unwrap();
4365        Bytes::from(parquet_data)
4366    }
4367
4368    #[test]
4369    fn test_schema_error_bad_types() {
4370        // verify incompatible schemas error on read
4371        let parquet_data = utf8_parquet();
4372
4373        // Ask to read it back with an incompatible schema (int vs string)
4374        let input_schema: SchemaRef = Arc::new(Schema::new(vec![Field::new(
4375            "column1",
4376            arrow::datatypes::DataType::Int32,
4377            false,
4378        )]));
4379
4380        // read it back out
4381        let reader_options = ArrowReaderOptions::new().with_schema(input_schema.clone());
4382        let err =
4383            ParquetRecordBatchReaderBuilder::try_new_with_options(parquet_data, reader_options)
4384                .unwrap_err();
4385        assert_eq!(
4386            err.to_string(),
4387            "Arrow: Incompatible supplied Arrow schema: data type mismatch for field column1: requested Int32 but found Utf8"
4388        )
4389    }
4390
4391    #[test]
4392    fn test_schema_error_bad_nullability() {
4393        // verify incompatible schemas error on read
4394        let parquet_data = utf8_parquet();
4395
4396        // Ask to read it back with an incompatible schema (nullability mismatch)
4397        let input_schema: SchemaRef = Arc::new(Schema::new(vec![Field::new(
4398            "column1",
4399            arrow::datatypes::DataType::Utf8,
4400            true,
4401        )]));
4402
4403        // read it back out
4404        let reader_options = ArrowReaderOptions::new().with_schema(input_schema.clone());
4405        let err =
4406            ParquetRecordBatchReaderBuilder::try_new_with_options(parquet_data, reader_options)
4407                .unwrap_err();
4408        assert_eq!(
4409            err.to_string(),
4410            "Arrow: Incompatible supplied Arrow schema: nullability mismatch for field column1: expected true but found false"
4411        )
4412    }
4413
4414    #[test]
4415    fn test_read_binary_as_utf8() {
4416        let file = write_parquet_from_iter(vec![
4417            (
4418                "binary_to_utf8",
4419                Arc::new(BinaryArray::from(vec![
4420                    b"one".as_ref(),
4421                    b"two".as_ref(),
4422                    b"three".as_ref(),
4423                ])) as ArrayRef,
4424            ),
4425            (
4426                "large_binary_to_large_utf8",
4427                Arc::new(LargeBinaryArray::from(vec![
4428                    b"one".as_ref(),
4429                    b"two".as_ref(),
4430                    b"three".as_ref(),
4431                ])) as ArrayRef,
4432            ),
4433            (
4434                "binary_view_to_utf8_view",
4435                Arc::new(BinaryViewArray::from(vec![
4436                    b"one".as_ref(),
4437                    b"two".as_ref(),
4438                    b"three".as_ref(),
4439                ])) as ArrayRef,
4440            ),
4441        ]);
4442        let supplied_fields = Fields::from(vec![
4443            Field::new("binary_to_utf8", ArrowDataType::Utf8, false),
4444            Field::new(
4445                "large_binary_to_large_utf8",
4446                ArrowDataType::LargeUtf8,
4447                false,
4448            ),
4449            Field::new("binary_view_to_utf8_view", ArrowDataType::Utf8View, false),
4450        ]);
4451
4452        let options = ArrowReaderOptions::new().with_schema(Arc::new(Schema::new(supplied_fields)));
4453        let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
4454            file.try_clone().unwrap(),
4455            options,
4456        )
4457        .expect("reader builder with schema")
4458        .build()
4459        .expect("reader with schema");
4460
4461        let batch = arrow_reader.next().unwrap().unwrap();
4462        assert_eq!(batch.num_columns(), 3);
4463        assert_eq!(batch.num_rows(), 3);
4464        assert_eq!(
4465            batch
4466                .column(0)
4467                .as_string::<i32>()
4468                .iter()
4469                .collect::<Vec<_>>(),
4470            vec![Some("one"), Some("two"), Some("three")]
4471        );
4472
4473        assert_eq!(
4474            batch
4475                .column(1)
4476                .as_string::<i64>()
4477                .iter()
4478                .collect::<Vec<_>>(),
4479            vec![Some("one"), Some("two"), Some("three")]
4480        );
4481
4482        assert_eq!(
4483            batch.column(2).as_string_view().iter().collect::<Vec<_>>(),
4484            vec![Some("one"), Some("two"), Some("three")]
4485        );
4486    }
4487
4488    #[test]
4489    #[should_panic(expected = "Invalid UTF8 sequence at")]
4490    fn test_read_non_utf8_binary_as_utf8() {
4491        let file = write_parquet_from_iter(vec![(
4492            "non_utf8_binary",
4493            Arc::new(BinaryArray::from(vec![
4494                b"\xDE\x00\xFF".as_ref(),
4495                b"\xDE\x01\xAA".as_ref(),
4496                b"\xDE\x02\xFF".as_ref(),
4497            ])) as ArrayRef,
4498        )]);
4499        let supplied_fields = Fields::from(vec![Field::new(
4500            "non_utf8_binary",
4501            ArrowDataType::Utf8,
4502            false,
4503        )]);
4504
4505        let options = ArrowReaderOptions::new().with_schema(Arc::new(Schema::new(supplied_fields)));
4506        let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
4507            file.try_clone().unwrap(),
4508            options,
4509        )
4510        .expect("reader builder with schema")
4511        .build()
4512        .expect("reader with schema");
4513        arrow_reader.next().unwrap().unwrap_err();
4514    }
4515
4516    #[test]
4517    fn test_with_schema() {
4518        let nested_fields = Fields::from(vec![
4519            Field::new("utf8_to_dict", ArrowDataType::Utf8, false),
4520            Field::new("int64_to_ts_nano", ArrowDataType::Int64, false),
4521        ]);
4522
4523        let nested_arrays: Vec<ArrayRef> = vec![
4524            Arc::new(StringArray::from(vec!["a", "a", "a", "b"])) as ArrayRef,
4525            Arc::new(Int64Array::from(vec![1, 2, 3, 4])) as ArrayRef,
4526        ];
4527
4528        let nested = StructArray::try_new(nested_fields, nested_arrays, None).unwrap();
4529
4530        let file = write_parquet_from_iter(vec![
4531            (
4532                "int32_to_ts_second",
4533                Arc::new(Int32Array::from(vec![0, 1, 2, 3])) as ArrayRef,
4534            ),
4535            (
4536                "date32_to_date64",
4537                Arc::new(Date32Array::from(vec![0, 1, 2, 3])) as ArrayRef,
4538            ),
4539            ("nested", Arc::new(nested) as ArrayRef),
4540        ]);
4541
4542        let supplied_nested_fields = Fields::from(vec![
4543            Field::new(
4544                "utf8_to_dict",
4545                ArrowDataType::Dictionary(
4546                    Box::new(ArrowDataType::Int32),
4547                    Box::new(ArrowDataType::Utf8),
4548                ),
4549                false,
4550            ),
4551            Field::new(
4552                "int64_to_ts_nano",
4553                ArrowDataType::Timestamp(
4554                    arrow::datatypes::TimeUnit::Nanosecond,
4555                    Some("+10:00".into()),
4556                ),
4557                false,
4558            ),
4559        ]);
4560
4561        let supplied_schema = Arc::new(Schema::new(vec![
4562            Field::new(
4563                "int32_to_ts_second",
4564                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Second, Some("+01:00".into())),
4565                false,
4566            ),
4567            Field::new("date32_to_date64", ArrowDataType::Date64, false),
4568            Field::new(
4569                "nested",
4570                ArrowDataType::Struct(supplied_nested_fields),
4571                false,
4572            ),
4573        ]));
4574
4575        let options = ArrowReaderOptions::new().with_schema(supplied_schema.clone());
4576        let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
4577            file.try_clone().unwrap(),
4578            options,
4579        )
4580        .expect("reader builder with schema")
4581        .build()
4582        .expect("reader with schema");
4583
4584        assert_eq!(arrow_reader.schema(), supplied_schema);
4585        let batch = arrow_reader.next().unwrap().unwrap();
4586        assert_eq!(batch.num_columns(), 3);
4587        assert_eq!(batch.num_rows(), 4);
4588        assert_eq!(
4589            batch
4590                .column(0)
4591                .as_any()
4592                .downcast_ref::<TimestampSecondArray>()
4593                .expect("downcast to timestamp second")
4594                .value_as_datetime_with_tz(0, "+01:00".parse().unwrap())
4595                .map(|v| v.to_string())
4596                .expect("value as datetime"),
4597            "1970-01-01 01:00:00 +01:00"
4598        );
4599        assert_eq!(
4600            batch
4601                .column(1)
4602                .as_any()
4603                .downcast_ref::<Date64Array>()
4604                .expect("downcast to date64")
4605                .value_as_date(0)
4606                .map(|v| v.to_string())
4607                .expect("value as date"),
4608            "1970-01-01"
4609        );
4610
4611        let nested = batch
4612            .column(2)
4613            .as_any()
4614            .downcast_ref::<StructArray>()
4615            .expect("downcast to struct");
4616
4617        let nested_dict = nested
4618            .column(0)
4619            .as_any()
4620            .downcast_ref::<Int32DictionaryArray>()
4621            .expect("downcast to dictionary");
4622
4623        assert_eq!(
4624            nested_dict
4625                .values()
4626                .as_any()
4627                .downcast_ref::<StringArray>()
4628                .expect("downcast to string")
4629                .iter()
4630                .collect::<Vec<_>>(),
4631            vec![Some("a"), Some("b")]
4632        );
4633
4634        assert_eq!(
4635            nested_dict.keys().iter().collect::<Vec<_>>(),
4636            vec![Some(0), Some(0), Some(0), Some(1)]
4637        );
4638
4639        assert_eq!(
4640            nested
4641                .column(1)
4642                .as_any()
4643                .downcast_ref::<TimestampNanosecondArray>()
4644                .expect("downcast to timestamp nanosecond")
4645                .value_as_datetime_with_tz(0, "+10:00".parse().unwrap())
4646                .map(|v| v.to_string())
4647                .expect("value as datetime"),
4648            "1970-01-01 10:00:00.000000001 +10:00"
4649        );
4650    }
4651
4652    #[test]
4653    fn test_empty_projection() {
4654        let testdata = arrow::util::test_util::parquet_test_data();
4655        let path = format!("{testdata}/alltypes_plain.parquet");
4656        let file = File::open(path).unwrap();
4657
4658        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
4659        let file_metadata = builder.metadata().file_metadata();
4660        let expected_rows = file_metadata.num_rows() as usize;
4661
4662        let mask = ProjectionMask::leaves(builder.parquet_schema(), []);
4663        let batch_reader = builder
4664            .with_projection(mask)
4665            .with_batch_size(2)
4666            .build()
4667            .unwrap();
4668
4669        let mut total_rows = 0;
4670        for maybe_batch in batch_reader {
4671            let batch = maybe_batch.unwrap();
4672            total_rows += batch.num_rows();
4673            assert_eq!(batch.num_columns(), 0);
4674            assert!(batch.num_rows() <= 2);
4675        }
4676
4677        assert_eq!(total_rows, expected_rows);
4678    }
4679
4680    fn test_row_group_batch(row_group_size: usize, batch_size: usize) {
4681        let schema = Arc::new(Schema::new(vec![Field::new(
4682            "list",
4683            ArrowDataType::List(Arc::new(Field::new_list_field(ArrowDataType::Int32, true))),
4684            true,
4685        )]));
4686
4687        let mut buf = Vec::with_capacity(1024);
4688
4689        let mut writer = ArrowWriter::try_new(
4690            &mut buf,
4691            schema.clone(),
4692            Some(
4693                WriterProperties::builder()
4694                    .set_max_row_group_row_count(Some(row_group_size))
4695                    .build(),
4696            ),
4697        )
4698        .unwrap();
4699        for _ in 0..2 {
4700            let mut list_builder = ListBuilder::new(Int32Builder::with_capacity(batch_size));
4701            for _ in 0..(batch_size) {
4702                list_builder.append(true);
4703            }
4704            let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(list_builder.finish())])
4705                .unwrap();
4706            writer.write(&batch).unwrap();
4707        }
4708        writer.close().unwrap();
4709
4710        let mut record_reader =
4711            ParquetRecordBatchReader::try_new(Bytes::from(buf), batch_size).unwrap();
4712        assert_eq!(
4713            batch_size,
4714            record_reader.next().unwrap().unwrap().num_rows()
4715        );
4716        assert_eq!(
4717            batch_size,
4718            record_reader.next().unwrap().unwrap().num_rows()
4719        );
4720    }
4721
4722    #[test]
4723    #[cfg_attr(miri, ignore)] // Takes too long
4724    fn test_row_group_exact_multiple() {
4725        const BATCH_SIZE: usize = REPETITION_LEVELS_BATCH_SIZE;
4726        test_row_group_batch(8, 8);
4727        test_row_group_batch(10, 8);
4728        test_row_group_batch(8, 10);
4729        test_row_group_batch(BATCH_SIZE, BATCH_SIZE);
4730        test_row_group_batch(BATCH_SIZE + 1, BATCH_SIZE);
4731        test_row_group_batch(BATCH_SIZE, BATCH_SIZE + 1);
4732        test_row_group_batch(BATCH_SIZE, BATCH_SIZE - 1);
4733        test_row_group_batch(BATCH_SIZE - 1, BATCH_SIZE);
4734    }
4735
4736    /// Given a RecordBatch containing all the column data, return the expected batches given
4737    /// a `batch_size` and `selection`
4738    fn get_expected_batches(
4739        column: &RecordBatch,
4740        selection: &RowSelection,
4741        batch_size: usize,
4742    ) -> Vec<RecordBatch> {
4743        let mut expected_batches = vec![];
4744
4745        let mut selection: VecDeque<_> = selection.clone().into();
4746        let mut row_offset = 0;
4747        let mut last_start = None;
4748        while row_offset < column.num_rows() && !selection.is_empty() {
4749            let mut batch_remaining = batch_size.min(column.num_rows() - row_offset);
4750            while batch_remaining > 0 && !selection.is_empty() {
4751                let (to_read, skip) = match selection.front_mut() {
4752                    Some(selection) if selection.row_count > batch_remaining => {
4753                        selection.row_count -= batch_remaining;
4754                        (batch_remaining, selection.skip)
4755                    }
4756                    Some(_) => {
4757                        let select = selection.pop_front().unwrap();
4758                        (select.row_count, select.skip)
4759                    }
4760                    None => break,
4761                };
4762
4763                batch_remaining -= to_read;
4764
4765                match skip {
4766                    true => {
4767                        if let Some(last_start) = last_start.take() {
4768                            expected_batches.push(column.slice(last_start, row_offset - last_start))
4769                        }
4770                        row_offset += to_read
4771                    }
4772                    false => {
4773                        last_start.get_or_insert(row_offset);
4774                        row_offset += to_read
4775                    }
4776                }
4777            }
4778        }
4779
4780        if let Some(last_start) = last_start.take() {
4781            expected_batches.push(column.slice(last_start, row_offset - last_start))
4782        }
4783
4784        // Sanity check, all batches except the final should be the batch size
4785        for batch in &expected_batches[..expected_batches.len() - 1] {
4786            assert_eq!(batch.num_rows(), batch_size);
4787        }
4788
4789        expected_batches
4790    }
4791
4792    fn create_test_selection(
4793        step_len: usize,
4794        total_len: usize,
4795        skip_first: bool,
4796    ) -> (RowSelection, usize) {
4797        let mut remaining = total_len;
4798        let mut skip = skip_first;
4799        let mut vec = vec![];
4800        let mut selected_count = 0;
4801        while remaining != 0 {
4802            let step = if remaining > step_len {
4803                step_len
4804            } else {
4805                remaining
4806            };
4807            vec.push(RowSelector {
4808                row_count: step,
4809                skip,
4810            });
4811            remaining -= step;
4812            if !skip {
4813                selected_count += step;
4814            }
4815            skip = !skip;
4816        }
4817        (vec.into(), selected_count)
4818    }
4819
4820    #[test]
4821    #[cfg_attr(miri, ignore)] // Takes too long
4822    fn test_scan_row_with_selection() {
4823        let testdata = arrow::util::test_util::parquet_test_data();
4824        let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
4825        let test_file = File::open(&path).unwrap();
4826
4827        let mut serial_reader =
4828            ParquetRecordBatchReader::try_new(File::open(&path).unwrap(), 7300).unwrap();
4829        let data = serial_reader.next().unwrap().unwrap();
4830
4831        let do_test = |batch_size: usize, selection_len: usize| {
4832            for skip_first in [false, true] {
4833                let selections = create_test_selection(batch_size, data.num_rows(), skip_first).0;
4834
4835                let expected = get_expected_batches(&data, &selections, batch_size);
4836                let skip_reader = create_skip_reader(&test_file, batch_size, selections);
4837                assert_eq!(
4838                    skip_reader.collect::<Result<Vec<_>, _>>().unwrap(),
4839                    expected,
4840                    "batch_size: {batch_size}, selection_len: {selection_len}, skip_first: {skip_first}"
4841                );
4842            }
4843        };
4844
4845        // total row count 7300
4846        // 1. test selection len more than one page row count
4847        do_test(1000, 1000);
4848
4849        // 2. test selection len less than one page row count
4850        do_test(20, 20);
4851
4852        // 3. test selection_len less than batch_size
4853        do_test(20, 5);
4854
4855        // 4. test selection_len more than batch_size
4856        // If batch_size < selection_len
4857        do_test(20, 5);
4858
4859        fn create_skip_reader(
4860            test_file: &File,
4861            batch_size: usize,
4862            selections: RowSelection,
4863        ) -> ParquetRecordBatchReader {
4864            let options =
4865                ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
4866            let file = test_file.try_clone().unwrap();
4867            ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
4868                .unwrap()
4869                .with_batch_size(batch_size)
4870                .with_row_selection(selections)
4871                .build()
4872                .unwrap()
4873        }
4874    }
4875
4876    #[test]
4877    fn test_batch_size_overallocate() {
4878        let testdata = arrow::util::test_util::parquet_test_data();
4879        // `alltypes_plain.parquet` only have 8 rows
4880        let path = format!("{testdata}/alltypes_plain.parquet");
4881        let test_file = File::open(path).unwrap();
4882
4883        let builder = ParquetRecordBatchReaderBuilder::try_new(test_file).unwrap();
4884        let num_rows = builder.metadata.file_metadata().num_rows();
4885        let reader = builder
4886            .with_batch_size(1024)
4887            .with_projection(ProjectionMask::all())
4888            .build()
4889            .unwrap();
4890        assert_ne!(1024, num_rows);
4891        assert_eq!(reader.read_plan.batch_size(), num_rows as usize);
4892    }
4893
4894    #[test]
4895    #[cfg_attr(miri, ignore)] // Takes too long
4896    fn test_read_with_page_index_enabled() {
4897        let testdata = arrow::util::test_util::parquet_test_data();
4898
4899        {
4900            // `alltypes_tiny_pages.parquet` has page index
4901            let path = format!("{testdata}/alltypes_tiny_pages.parquet");
4902            let test_file = File::open(path).unwrap();
4903            let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
4904                test_file,
4905                ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
4906            )
4907            .unwrap();
4908            let page_index = builder
4909                .metadata()
4910                .page_index()
4911                .expect("page index should be present");
4912            let num_columns = builder.metadata().row_group(0).num_columns();
4913            let offset_indexes = page_index.offset_indexes_for_rowgroup(0);
4914            assert!(offset_indexes.is_some_and(|ois| ois.len() == num_columns));
4915            let column_indexes = page_index.offset_indexes_for_rowgroup(0);
4916            assert!(column_indexes.is_some_and(|cis| cis.len() == num_columns));
4917            assert!(page_index.offset_index(0, 0).is_some());
4918            assert!(page_index.column_index(0, 0).is_some());
4919            assert!(page_index.page_locations(0, 0).is_some());
4920            assert_eq!(page_index.num_data_pages(0, 0), Some(325));
4921            let reader = builder.build().unwrap();
4922            let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
4923            assert_eq!(batches.len(), 8);
4924        }
4925
4926        {
4927            // `alltypes_plain.parquet` doesn't have page index
4928            let path = format!("{testdata}/alltypes_plain.parquet");
4929            let test_file = File::open(path).unwrap();
4930            let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
4931                test_file,
4932                ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
4933            )
4934            .unwrap();
4935            // Although `Vec<Vec<PageLoacation>>` of each row group is empty,
4936            // we should read the file successfully.
4937            assert!(builder.metadata().page_index().is_none());
4938            let reader = builder.build().unwrap();
4939            let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
4940            assert_eq!(batches.len(), 1);
4941        }
4942    }
4943
4944    #[test]
4945    fn test_raw_repetition() {
4946        const MESSAGE_TYPE: &str = "
4947            message Log {
4948              OPTIONAL INT32 eventType;
4949              REPEATED INT32 category;
4950              REPEATED group filter {
4951                OPTIONAL INT32 error;
4952              }
4953            }
4954        ";
4955        let schema = Arc::new(parse_message_type(MESSAGE_TYPE).unwrap());
4956        let props = Default::default();
4957
4958        let mut buf = Vec::with_capacity(1024);
4959        let mut writer = SerializedFileWriter::new(&mut buf, schema, props).unwrap();
4960        let mut row_group_writer = writer.next_row_group().unwrap();
4961
4962        // column 0
4963        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4964        col_writer
4965            .typed::<Int32Type>()
4966            .write_batch(&[1], Some(&[1]), None)
4967            .unwrap();
4968        col_writer.close().unwrap();
4969        // column 1
4970        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4971        col_writer
4972            .typed::<Int32Type>()
4973            .write_batch(&[1, 1], Some(&[1, 1]), Some(&[0, 1]))
4974            .unwrap();
4975        col_writer.close().unwrap();
4976        // column 2
4977        let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4978        col_writer
4979            .typed::<Int32Type>()
4980            .write_batch(&[1], Some(&[1]), Some(&[0]))
4981            .unwrap();
4982        col_writer.close().unwrap();
4983
4984        let rg_md = row_group_writer.close().unwrap();
4985        assert_eq!(rg_md.num_rows(), 1);
4986        writer.close().unwrap();
4987
4988        let bytes = Bytes::from(buf);
4989
4990        let mut no_mask = ParquetRecordBatchReader::try_new(bytes.clone(), 1024).unwrap();
4991        let full = no_mask.next().unwrap().unwrap();
4992
4993        assert_eq!(full.num_columns(), 3);
4994
4995        for idx in 0..3 {
4996            let b = ParquetRecordBatchReaderBuilder::try_new(bytes.clone()).unwrap();
4997            let mask = ProjectionMask::leaves(b.parquet_schema(), [idx]);
4998            let mut reader = b.with_projection(mask).build().unwrap();
4999            let projected = reader.next().unwrap().unwrap();
5000
5001            assert_eq!(projected.num_columns(), 1);
5002            assert_eq!(full.column(idx), projected.column(0));
5003        }
5004    }
5005
5006    #[test]
5007    fn test_read_lz4_raw() {
5008        let testdata = arrow::util::test_util::parquet_test_data();
5009        let path = format!("{testdata}/lz4_raw_compressed.parquet");
5010        let file = File::open(path).unwrap();
5011
5012        let batches = ParquetRecordBatchReader::try_new(file, 1024)
5013            .unwrap()
5014            .collect::<Result<Vec<_>, _>>()
5015            .unwrap();
5016        assert_eq!(batches.len(), 1);
5017        let batch = &batches[0];
5018
5019        assert_eq!(batch.num_columns(), 3);
5020        assert_eq!(batch.num_rows(), 4);
5021
5022        // https://github.com/apache/parquet-testing/pull/18
5023        let a: &Int64Array = batch.column(0).as_any().downcast_ref().unwrap();
5024        assert_eq!(
5025            a.values(),
5026            &[1593604800, 1593604800, 1593604801, 1593604801]
5027        );
5028
5029        let a: &BinaryArray = batch.column(1).as_any().downcast_ref().unwrap();
5030        let a: Vec<_> = a.iter().flatten().collect();
5031        assert_eq!(a, &[b"abc", b"def", b"abc", b"def"]);
5032
5033        let a: &Float64Array = batch.column(2).as_any().downcast_ref().unwrap();
5034        assert_eq!(a.values(), &[42.000000, 7.700000, 42.125000, 7.700000]);
5035    }
5036
5037    // This test is to ensure backward compatibility, it test 2 files containing the LZ4 CompressionCodec
5038    // but different algorithms: LZ4_HADOOP and LZ4_RAW.
5039    // 1. hadoop_lz4_compressed.parquet -> It is a file with LZ4 CompressionCodec which uses
5040    //    LZ4_HADOOP algorithm for compression.
5041    // 2. non_hadoop_lz4_compressed.parquet -> It is a file with LZ4 CompressionCodec which uses
5042    //    LZ4_RAW algorithm for compression. This fallback is done to keep backward compatibility with
5043    //    older parquet-cpp versions.
5044    //
5045    // For more information, check: https://github.com/apache/arrow-rs/issues/2988
5046    #[test]
5047    #[cfg_attr(miri, ignore)] // Takes too long
5048    fn test_read_lz4_hadoop_fallback() {
5049        for file in [
5050            "hadoop_lz4_compressed.parquet",
5051            "non_hadoop_lz4_compressed.parquet",
5052        ] {
5053            let testdata = arrow::util::test_util::parquet_test_data();
5054            let path = format!("{testdata}/{file}");
5055            let file = File::open(path).unwrap();
5056            let expected_rows = 4;
5057
5058            let batches = ParquetRecordBatchReader::try_new(file, expected_rows)
5059                .unwrap()
5060                .collect::<Result<Vec<_>, _>>()
5061                .unwrap();
5062            assert_eq!(batches.len(), 1);
5063            let batch = &batches[0];
5064
5065            assert_eq!(batch.num_columns(), 3);
5066            assert_eq!(batch.num_rows(), expected_rows);
5067
5068            let a: &Int64Array = batch.column(0).as_any().downcast_ref().unwrap();
5069            assert_eq!(
5070                a.values(),
5071                &[1593604800, 1593604800, 1593604801, 1593604801]
5072            );
5073
5074            let b: &BinaryArray = batch.column(1).as_any().downcast_ref().unwrap();
5075            let b: Vec<_> = b.iter().flatten().collect();
5076            assert_eq!(b, &[b"abc", b"def", b"abc", b"def"]);
5077
5078            let c: &Float64Array = batch.column(2).as_any().downcast_ref().unwrap();
5079            assert_eq!(c.values(), &[42.0, 7.7, 42.125, 7.7]);
5080        }
5081    }
5082
5083    #[test]
5084    #[cfg_attr(miri, ignore)] // Takes too long
5085    fn test_read_lz4_hadoop_large() {
5086        let testdata = arrow::util::test_util::parquet_test_data();
5087        let path = format!("{testdata}/hadoop_lz4_compressed_larger.parquet");
5088        let file = File::open(path).unwrap();
5089        let expected_rows = 10000;
5090
5091        let batches = ParquetRecordBatchReader::try_new(file, expected_rows)
5092            .unwrap()
5093            .collect::<Result<Vec<_>, _>>()
5094            .unwrap();
5095        assert_eq!(batches.len(), 1);
5096        let batch = &batches[0];
5097
5098        assert_eq!(batch.num_columns(), 1);
5099        assert_eq!(batch.num_rows(), expected_rows);
5100
5101        let a: &StringArray = batch.column(0).as_any().downcast_ref().unwrap();
5102        let a: Vec<_> = a.iter().flatten().collect();
5103        assert_eq!(a[0], "c7ce6bef-d5b0-4863-b199-8ea8c7fb117b");
5104        assert_eq!(a[1], "e8fb9197-cb9f-4118-b67f-fbfa65f61843");
5105        assert_eq!(a[expected_rows - 2], "ab52a0cc-c6bb-4d61-8a8f-166dc4b8b13c");
5106        assert_eq!(a[expected_rows - 1], "85440778-460a-41ac-aa2e-ac3ee41696bf");
5107    }
5108
5109    #[test]
5110    #[cfg(feature = "snap")]
5111    fn test_read_nested_lists() {
5112        let testdata = arrow::util::test_util::parquet_test_data();
5113        let path = format!("{testdata}/nested_lists.snappy.parquet");
5114        let file = File::open(path).unwrap();
5115
5116        let f = file.try_clone().unwrap();
5117        let mut reader = ParquetRecordBatchReader::try_new(f, 60).unwrap();
5118        let expected = reader.next().unwrap().unwrap();
5119        assert_eq!(expected.num_rows(), 3);
5120
5121        let selection = RowSelection::from(vec![
5122            RowSelector::skip(1),
5123            RowSelector::select(1),
5124            RowSelector::skip(1),
5125        ]);
5126        let mut reader = ParquetRecordBatchReaderBuilder::try_new(file)
5127            .unwrap()
5128            .with_row_selection(selection)
5129            .build()
5130            .unwrap();
5131
5132        let actual = reader.next().unwrap().unwrap();
5133        assert_eq!(actual.num_rows(), 1);
5134        assert_eq!(actual.column(0), &expected.column(0).slice(1, 1));
5135    }
5136
5137    #[test]
5138    fn test_arbitrary_decimal() {
5139        let values = [1, 2, 3, 4, 5, 6, 7, 8];
5140        let decimals_19_0 = Decimal128Array::from_iter_values(values)
5141            .with_precision_and_scale(19, 0)
5142            .unwrap();
5143        let decimals_12_0 = Decimal128Array::from_iter_values(values)
5144            .with_precision_and_scale(12, 0)
5145            .unwrap();
5146        let decimals_17_10 = Decimal128Array::from_iter_values(values)
5147            .with_precision_and_scale(17, 10)
5148            .unwrap();
5149
5150        let written = RecordBatch::try_from_iter([
5151            ("decimal_values_19_0", Arc::new(decimals_19_0) as ArrayRef),
5152            ("decimal_values_12_0", Arc::new(decimals_12_0) as ArrayRef),
5153            ("decimal_values_17_10", Arc::new(decimals_17_10) as ArrayRef),
5154        ])
5155        .unwrap();
5156
5157        let mut buffer = Vec::with_capacity(1024);
5158        let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
5159        writer.write(&written).unwrap();
5160        writer.close().unwrap();
5161
5162        let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 8)
5163            .unwrap()
5164            .collect::<Result<Vec<_>, _>>()
5165            .unwrap();
5166
5167        assert_eq!(&written.slice(0, 8), &read[0]);
5168    }
5169
5170    #[test]
5171    fn test_list_skip() {
5172        let mut list = ListBuilder::new(Int32Builder::new());
5173        list.append_value([Some(1), Some(2)]);
5174        list.append_value([Some(3)]);
5175        list.append_value([Some(4)]);
5176        let list = list.finish();
5177        let batch = RecordBatch::try_from_iter([("l", Arc::new(list) as _)]).unwrap();
5178
5179        // First page contains 2 values but only 1 row
5180        let props = WriterProperties::builder()
5181            .set_data_page_row_count_limit(1)
5182            .set_write_batch_size(2)
5183            .build();
5184
5185        let mut buffer = Vec::with_capacity(1024);
5186        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), Some(props)).unwrap();
5187        writer.write(&batch).unwrap();
5188        writer.close().unwrap();
5189
5190        let selection = vec![RowSelector::skip(2), RowSelector::select(1)];
5191        let mut reader = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer))
5192            .unwrap()
5193            .with_row_selection(selection.into())
5194            .build()
5195            .unwrap();
5196        let out = reader.next().unwrap().unwrap();
5197        assert_eq!(out.num_rows(), 1);
5198        assert_eq!(out, batch.slice(2, 1));
5199    }
5200
5201    fn test_decimal32_roundtrip() {
5202        let d = |values: Vec<i32>, p: u8| {
5203            let iter = values.into_iter();
5204            PrimitiveArray::<Decimal32Type>::from_iter_values(iter)
5205                .with_precision_and_scale(p, 2)
5206                .unwrap()
5207        };
5208
5209        let d1 = d(vec![1, 2, 3, 4, 5], 9);
5210        let batch = RecordBatch::try_from_iter([("d1", Arc::new(d1) as ArrayRef)]).unwrap();
5211
5212        let mut buffer = Vec::with_capacity(1024);
5213        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
5214        writer.write(&batch).unwrap();
5215        writer.close().unwrap();
5216
5217        let builder = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer)).unwrap();
5218        let t1 = builder.parquet_schema().columns()[0].physical_type();
5219        assert_eq!(t1, PhysicalType::INT32);
5220
5221        let mut reader = builder.build().unwrap();
5222        assert_eq!(batch.schema(), reader.schema());
5223
5224        let out = reader.next().unwrap().unwrap();
5225        assert_eq!(batch, out);
5226    }
5227
5228    fn test_decimal64_roundtrip() {
5229        // Precision <= 9 -> INT32
5230        // Precision <= 18 -> INT64
5231
5232        let d = |values: Vec<i64>, p: u8| {
5233            let iter = values.into_iter();
5234            PrimitiveArray::<Decimal64Type>::from_iter_values(iter)
5235                .with_precision_and_scale(p, 2)
5236                .unwrap()
5237        };
5238
5239        let d1 = d(vec![1, 2, 3, 4, 5], 9);
5240        let d2 = d(vec![1, 2, 3, 4, 10.pow(10) - 1], 10);
5241        let d3 = d(vec![1, 2, 3, 4, 10.pow(18) - 1], 18);
5242
5243        let batch = RecordBatch::try_from_iter([
5244            ("d1", Arc::new(d1) as ArrayRef),
5245            ("d2", Arc::new(d2) as ArrayRef),
5246            ("d3", Arc::new(d3) as ArrayRef),
5247        ])
5248        .unwrap();
5249
5250        let mut buffer = Vec::with_capacity(1024);
5251        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
5252        writer.write(&batch).unwrap();
5253        writer.close().unwrap();
5254
5255        let builder = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer)).unwrap();
5256        let t1 = builder.parquet_schema().columns()[0].physical_type();
5257        assert_eq!(t1, PhysicalType::INT32);
5258        let t2 = builder.parquet_schema().columns()[1].physical_type();
5259        assert_eq!(t2, PhysicalType::INT64);
5260        let t3 = builder.parquet_schema().columns()[2].physical_type();
5261        assert_eq!(t3, PhysicalType::INT64);
5262
5263        let mut reader = builder.build().unwrap();
5264        assert_eq!(batch.schema(), reader.schema());
5265
5266        let out = reader.next().unwrap().unwrap();
5267        assert_eq!(batch, out);
5268    }
5269
5270    fn test_decimal_roundtrip<T: DecimalType>() {
5271        // Precision <= 9 -> INT32
5272        // Precision <= 18 -> INT64
5273        // Precision > 18 -> FIXED_LEN_BYTE_ARRAY
5274
5275        let d = |values: Vec<usize>, p: u8| {
5276            let iter = values.into_iter().map(T::Native::usize_as);
5277            PrimitiveArray::<T>::from_iter_values(iter)
5278                .with_precision_and_scale(p, 2)
5279                .unwrap()
5280        };
5281
5282        let d1 = d(vec![1, 2, 3, 4, 5], 9);
5283        let d2 = d(vec![1, 2, 3, 4, 10.pow(10) - 1], 10);
5284        let d3 = d(vec![1, 2, 3, 4, 10.pow(18) - 1], 18);
5285        let d4 = d(vec![1, 2, 3, 4, 10.pow(19) - 1], 19);
5286
5287        let batch = RecordBatch::try_from_iter([
5288            ("d1", Arc::new(d1) as ArrayRef),
5289            ("d2", Arc::new(d2) as ArrayRef),
5290            ("d3", Arc::new(d3) as ArrayRef),
5291            ("d4", Arc::new(d4) as ArrayRef),
5292        ])
5293        .unwrap();
5294
5295        let mut buffer = Vec::with_capacity(1024);
5296        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
5297        writer.write(&batch).unwrap();
5298        writer.close().unwrap();
5299
5300        let builder = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer)).unwrap();
5301        let t1 = builder.parquet_schema().columns()[0].physical_type();
5302        assert_eq!(t1, PhysicalType::INT32);
5303        let t2 = builder.parquet_schema().columns()[1].physical_type();
5304        assert_eq!(t2, PhysicalType::INT64);
5305        let t3 = builder.parquet_schema().columns()[2].physical_type();
5306        assert_eq!(t3, PhysicalType::INT64);
5307        let t4 = builder.parquet_schema().columns()[3].physical_type();
5308        assert_eq!(t4, PhysicalType::FIXED_LEN_BYTE_ARRAY);
5309
5310        let mut reader = builder.build().unwrap();
5311        assert_eq!(batch.schema(), reader.schema());
5312
5313        let out = reader.next().unwrap().unwrap();
5314        assert_eq!(batch, out);
5315    }
5316
5317    #[test]
5318    fn test_decimal() {
5319        test_decimal32_roundtrip();
5320        test_decimal64_roundtrip();
5321        test_decimal_roundtrip::<Decimal128Type>();
5322        test_decimal_roundtrip::<Decimal256Type>();
5323    }
5324
5325    #[test]
5326    #[cfg_attr(miri, ignore)] // Takes too long
5327    fn test_list_selection() {
5328        let schema = Arc::new(Schema::new(vec![Field::new_list(
5329            "list",
5330            Field::new_list_field(ArrowDataType::Utf8, true),
5331            false,
5332        )]));
5333        let mut buf = Vec::with_capacity(1024);
5334
5335        let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None).unwrap();
5336
5337        for i in 0..2 {
5338            let mut list_a_builder = ListBuilder::new(StringBuilder::new());
5339            for j in 0..1024 {
5340                list_a_builder.values().append_value(format!("{i} {j}"));
5341                list_a_builder.append(true);
5342            }
5343            let batch =
5344                RecordBatch::try_new(schema.clone(), vec![Arc::new(list_a_builder.finish())])
5345                    .unwrap();
5346            writer.write(&batch).unwrap();
5347        }
5348        let _metadata = writer.close().unwrap();
5349
5350        let buf = Bytes::from(buf);
5351        let reader = ParquetRecordBatchReaderBuilder::try_new(buf)
5352            .unwrap()
5353            .with_row_selection(RowSelection::from(vec![
5354                RowSelector::skip(100),
5355                RowSelector::select(924),
5356                RowSelector::skip(100),
5357                RowSelector::select(924),
5358            ]))
5359            .build()
5360            .unwrap();
5361
5362        let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
5363        let batch = concat_batches(&schema, &batches).unwrap();
5364
5365        assert_eq!(batch.num_rows(), 924 * 2);
5366        let list = batch.column(0).as_list::<i32>();
5367
5368        for w in list.value_offsets().windows(2) {
5369            assert_eq!(w[0] + 1, w[1])
5370        }
5371        let mut values = list.values().as_string::<i32>().iter();
5372
5373        for i in 0..2 {
5374            for j in 100..1024 {
5375                let expected = format!("{i} {j}");
5376                assert_eq!(values.next().unwrap().unwrap(), &expected);
5377            }
5378        }
5379    }
5380
5381    #[test]
5382    #[cfg_attr(miri, ignore)] // Takes too long
5383    fn test_list_selection_fuzz() {
5384        let mut rng = rng();
5385        let schema = Arc::new(Schema::new(vec![Field::new_list(
5386            "list",
5387            Field::new_list(
5388                Field::LIST_FIELD_DEFAULT_NAME,
5389                Field::new_list_field(ArrowDataType::Int32, true),
5390                true,
5391            ),
5392            true,
5393        )]));
5394        let mut buf = Vec::with_capacity(1024);
5395        let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None).unwrap();
5396
5397        let mut list_a_builder = ListBuilder::new(ListBuilder::new(Int32Builder::new()));
5398
5399        for _ in 0..2048 {
5400            if rng.random_bool(0.2) {
5401                list_a_builder.append(false);
5402                continue;
5403            }
5404
5405            let list_a_len = rng.random_range(0..10);
5406            let list_b_builder = list_a_builder.values();
5407
5408            for _ in 0..list_a_len {
5409                if rng.random_bool(0.2) {
5410                    list_b_builder.append(false);
5411                    continue;
5412                }
5413
5414                let list_b_len = rng.random_range(0..10);
5415                let int_builder = list_b_builder.values();
5416                for _ in 0..list_b_len {
5417                    match rng.random_bool(0.2) {
5418                        true => int_builder.append_null(),
5419                        false => int_builder.append_value(rng.random()),
5420                    }
5421                }
5422                list_b_builder.append(true)
5423            }
5424            list_a_builder.append(true);
5425        }
5426
5427        let array = Arc::new(list_a_builder.finish());
5428        let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
5429
5430        writer.write(&batch).unwrap();
5431        let _metadata = writer.close().unwrap();
5432
5433        let buf = Bytes::from(buf);
5434
5435        let cases = [
5436            vec![
5437                RowSelector::skip(100),
5438                RowSelector::select(924),
5439                RowSelector::skip(100),
5440                RowSelector::select(924),
5441            ],
5442            vec![
5443                RowSelector::select(924),
5444                RowSelector::skip(100),
5445                RowSelector::select(924),
5446                RowSelector::skip(100),
5447            ],
5448            vec![
5449                RowSelector::skip(1023),
5450                RowSelector::select(1),
5451                RowSelector::skip(1023),
5452                RowSelector::select(1),
5453            ],
5454            vec![
5455                RowSelector::select(1),
5456                RowSelector::skip(1023),
5457                RowSelector::select(1),
5458                RowSelector::skip(1023),
5459            ],
5460        ];
5461
5462        for batch_size in [100, 1024, 2048] {
5463            for selection in &cases {
5464                let selection = RowSelection::from(selection.clone());
5465                let reader = ParquetRecordBatchReaderBuilder::try_new(buf.clone())
5466                    .unwrap()
5467                    .with_row_selection(selection.clone())
5468                    .with_batch_size(batch_size)
5469                    .build()
5470                    .unwrap();
5471
5472                let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
5473                let actual = concat_batches(batch.schema_ref(), &batches).unwrap();
5474                assert_eq!(actual.num_rows(), selection.row_count());
5475
5476                let mut batch_offset = 0;
5477                let mut actual_offset = 0;
5478                for selector in selection.iter() {
5479                    if selector.skip {
5480                        batch_offset += selector.row_count;
5481                        continue;
5482                    }
5483
5484                    assert_eq!(
5485                        batch.slice(batch_offset, selector.row_count),
5486                        actual.slice(actual_offset, selector.row_count)
5487                    );
5488
5489                    batch_offset += selector.row_count;
5490                    actual_offset += selector.row_count;
5491                }
5492            }
5493        }
5494    }
5495
5496    #[test]
5497    fn test_read_old_nested_list() {
5498        use arrow::datatypes::DataType;
5499        use arrow::datatypes::ToByteSlice;
5500
5501        let testdata = arrow::util::test_util::parquet_test_data();
5502        // message my_record {
5503        //     REQUIRED group a (LIST) {
5504        //         REPEATED group array (LIST) {
5505        //             REPEATED INT32 array;
5506        //         }
5507        //     }
5508        // }
5509        // should be read as list<list<int32>>
5510        let path = format!("{testdata}/old_list_structure.parquet");
5511        let test_file = File::open(path).unwrap();
5512
5513        // create expected ListArray
5514        let a_values = Int32Array::from(vec![1, 2, 3, 4]);
5515
5516        // Construct a buffer for value offsets, for the nested array: [[1, 2], [3, 4]]
5517        let a_value_offsets = arrow::buffer::Buffer::from([0, 2, 4].to_byte_slice());
5518
5519        // Construct a list array from the above two
5520        let a_list_data = ArrayData::builder(DataType::List(Arc::new(Field::new(
5521            "array",
5522            DataType::Int32,
5523            false,
5524        ))))
5525        .len(2)
5526        .add_buffer(a_value_offsets)
5527        .add_child_data(a_values.into_data())
5528        .build()
5529        .unwrap();
5530        let a = ListArray::from(a_list_data);
5531
5532        let builder = ParquetRecordBatchReaderBuilder::try_new(test_file).unwrap();
5533        let mut reader = builder.build().unwrap();
5534        let out = reader.next().unwrap().unwrap();
5535        assert_eq!(out.num_rows(), 1);
5536        assert_eq!(out.num_columns(), 1);
5537        // grab first column
5538        let c0 = out.column(0);
5539        let c0arr = c0.as_any().downcast_ref::<ListArray>().unwrap();
5540        // get first row: [[1, 2], [3, 4]]
5541        let r0 = c0arr.value(0);
5542        let r0arr = r0.as_any().downcast_ref::<ListArray>().unwrap();
5543        assert_eq!(r0arr, &a);
5544    }
5545
5546    #[test]
5547    fn test_read_row_numbers() {
5548        let file = write_parquet_from_iter(vec![(
5549            "value",
5550            Arc::new(Int64Array::from(vec![1, 2, 3])) as ArrayRef,
5551        )]);
5552        let supplied_fields = Fields::from(vec![Field::new("value", ArrowDataType::Int64, false)]);
5553
5554        let row_number_field = Arc::new(
5555            Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5556        );
5557
5558        let options = ArrowReaderOptions::new()
5559            .with_schema(Arc::new(Schema::new(supplied_fields)))
5560            .with_virtual_columns(vec![row_number_field.clone()])
5561            .unwrap();
5562        let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
5563            file.try_clone().unwrap(),
5564            options,
5565        )
5566        .expect("reader builder with schema")
5567        .build()
5568        .expect("reader with schema");
5569
5570        let batch = arrow_reader.next().unwrap().unwrap();
5571        let schema = Arc::new(Schema::new(vec![
5572            Field::new("value", ArrowDataType::Int64, false),
5573            (*row_number_field).clone(),
5574        ]));
5575
5576        assert_eq!(batch.schema(), schema);
5577        assert_eq!(batch.num_columns(), 2);
5578        assert_eq!(batch.num_rows(), 3);
5579        assert_eq!(
5580            batch
5581                .column(0)
5582                .as_primitive::<types::Int64Type>()
5583                .iter()
5584                .collect::<Vec<_>>(),
5585            vec![Some(1), Some(2), Some(3)]
5586        );
5587        assert_eq!(
5588            batch
5589                .column(1)
5590                .as_primitive::<types::Int64Type>()
5591                .iter()
5592                .collect::<Vec<_>>(),
5593            vec![Some(0), Some(1), Some(2)]
5594        );
5595    }
5596
5597    #[test]
5598    fn test_read_only_row_numbers() {
5599        let file = write_parquet_from_iter(vec![(
5600            "value",
5601            Arc::new(Int64Array::from(vec![1, 2, 3])) as ArrayRef,
5602        )]);
5603        let row_number_field = Arc::new(
5604            Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5605        );
5606        let options = ArrowReaderOptions::new()
5607            .with_virtual_columns(vec![row_number_field.clone()])
5608            .unwrap();
5609        let metadata = ArrowReaderMetadata::load(&file, options).unwrap();
5610        let num_columns = metadata
5611            .metadata
5612            .file_metadata()
5613            .schema_descr()
5614            .num_columns();
5615
5616        let mut arrow_reader = ParquetRecordBatchReaderBuilder::new_with_metadata(file, metadata)
5617            .with_projection(ProjectionMask::none(num_columns))
5618            .build()
5619            .expect("reader with schema");
5620
5621        let batch = arrow_reader.next().unwrap().unwrap();
5622        let schema = Arc::new(Schema::new(vec![row_number_field]));
5623
5624        assert_eq!(batch.schema(), schema);
5625        assert_eq!(batch.num_columns(), 1);
5626        assert_eq!(batch.num_rows(), 3);
5627        assert_eq!(
5628            batch
5629                .column(0)
5630                .as_primitive::<types::Int64Type>()
5631                .iter()
5632                .collect::<Vec<_>>(),
5633            vec![Some(0), Some(1), Some(2)]
5634        );
5635    }
5636
5637    #[test]
5638    fn test_read_row_numbers_row_group_order() -> Result<()> {
5639        // Make a parquet file with 100 rows split across 2 row groups
5640        let array = Int64Array::from_iter_values(5000..5100);
5641        let batch = RecordBatch::try_from_iter([("col", Arc::new(array) as ArrayRef)])?;
5642        let mut buffer = Vec::new();
5643        let options = WriterProperties::builder()
5644            .set_max_row_group_row_count(Some(50))
5645            .build();
5646        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema().clone(), Some(options))?;
5647        // write in 10 row batches as the size limits are enforced after each batch
5648        for batch_chunk in (0..10).map(|i| batch.slice(i * 10, 10)) {
5649            writer.write(&batch_chunk)?;
5650        }
5651        writer.close()?;
5652
5653        let row_number_field = Arc::new(
5654            Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5655        );
5656
5657        let buffer = Bytes::from(buffer);
5658
5659        let options =
5660            ArrowReaderOptions::new().with_virtual_columns(vec![row_number_field.clone()])?;
5661
5662        // read out with normal options
5663        let arrow_reader =
5664            ParquetRecordBatchReaderBuilder::try_new_with_options(buffer.clone(), options.clone())?
5665                .build()?;
5666
5667        assert_eq!(
5668            ValuesAndRowNumbers {
5669                values: (5000..5100).collect(),
5670                row_numbers: (0..100).collect()
5671            },
5672            ValuesAndRowNumbers::new_from_reader(arrow_reader)
5673        );
5674
5675        // Now read, out of order row groups
5676        let arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(buffer, options)?
5677            .with_row_groups(vec![1, 0])
5678            .build()?;
5679
5680        assert_eq!(
5681            ValuesAndRowNumbers {
5682                values: (5050..5100).chain(5000..5050).collect(),
5683                row_numbers: (50..100).chain(0..50).collect(),
5684            },
5685            ValuesAndRowNumbers::new_from_reader(arrow_reader)
5686        );
5687
5688        Ok(())
5689    }
5690
5691    /// A file with *mixed* row-group ordinal metadata (spec-valid — the
5692    /// `RowGroup.ordinal` thrift field is optional; Go parquet writers emit
5693    /// such files) must read fine without row numbers, and must fail
5694    /// deterministically with them — even when every *selected* row group
5695    /// carries an ordinal. See <https://github.com/apache/arrow-rs/issues/10381>.
5696    #[test]
5697    fn test_mixed_row_group_ordinals() -> Result<()> {
5698        use crate::file::metadata::{ParquetMetaDataReader, RowGroupMetaData};
5699
5700        // 100 rows split across 4 row groups of 25
5701        let array = Int64Array::from_iter_values(5000..5100);
5702        let batch = RecordBatch::try_from_iter([("col", Arc::new(array) as ArrayRef)])?;
5703        let mut buffer = Vec::new();
5704        let props = WriterProperties::builder()
5705            .set_max_row_group_row_count(Some(25))
5706            .build();
5707        let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema().clone(), Some(props))?;
5708        for batch_chunk in (0..10).map(|i| batch.slice(i * 10, 10)) {
5709            writer.write(&batch_chunk)?;
5710        }
5711        writer.close()?;
5712        let buffer = Bytes::from(buffer);
5713
5714        // Strip the ordinal from row group 1 to simulate a mixed-ordinal
5715        // writer (the builder starts with no ordinal; copy everything else).
5716        let metadata = ParquetMetaDataReader::new().parse_and_finish(&buffer)?;
5717        let schema_descr = metadata.file_metadata().schema_descr_ptr();
5718        let mut row_groups = metadata.row_groups().to_vec();
5719        let stripped = row_groups[1].clone();
5720        let mut builder = RowGroupMetaData::builder(schema_descr)
5721            .set_num_rows(stripped.num_rows())
5722            .set_total_byte_size(stripped.total_byte_size())
5723            .set_sorting_columns(stripped.sorting_columns().cloned())
5724            .set_column_metadata(stripped.columns().to_vec());
5725        if let Some(offset) = stripped.file_offset() {
5726            builder = builder.set_file_offset(offset);
5727        }
5728        row_groups[1] = builder.build()?;
5729        assert_eq!(row_groups[1].ordinal(), None);
5730        let metadata = Arc::new(metadata.into_builder().set_row_groups(row_groups).build());
5731
5732        // Plain read (no row numbers): succeeds and returns all values.
5733        let arrow_metadata =
5734            ArrowReaderMetadata::try_new(Arc::clone(&metadata), ArrowReaderOptions::new())?;
5735        let reader =
5736            ParquetRecordBatchReaderBuilder::new_with_metadata(buffer.clone(), arrow_metadata)
5737                .build()?;
5738        let values: Vec<i64> = reader
5739            .flat_map(|batch| {
5740                let batch = batch.expect("could not read batch");
5741                batch
5742                    .column(0)
5743                    .as_primitive::<types::Int64Type>()
5744                    .values()
5745                    .to_vec()
5746            })
5747            .collect();
5748        assert_eq!(values, (5000..5100).collect::<Vec<_>>());
5749
5750        // Row-number read: fails deterministically, even when selecting only
5751        // row groups that DO carry ordinals.
5752        let row_number_field = Arc::new(
5753            Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5754        );
5755        let options = ArrowReaderOptions::new().with_virtual_columns(vec![row_number_field])?;
5756        let arrow_metadata = ArrowReaderMetadata::try_new(Arc::clone(&metadata), options)?;
5757        let result = ParquetRecordBatchReaderBuilder::new_with_metadata(buffer, arrow_metadata)
5758            .with_row_groups(vec![0]) // row group 0 has an ordinal
5759            .build()
5760            .and_then(|mut reader| reader.next().transpose().map_err(|e| e.into()));
5761        let err = result.expect_err("row numbers over mixed ordinals must fail");
5762        assert!(
5763            err.to_string().contains("inconsistent row-group ordinals"),
5764            "unexpected error: {err}"
5765        );
5766
5767        Ok(())
5768    }
5769
5770    #[derive(Debug, PartialEq)]
5771    struct ValuesAndRowNumbers {
5772        values: Vec<i64>,
5773        row_numbers: Vec<i64>,
5774    }
5775    impl ValuesAndRowNumbers {
5776        fn new_from_reader(reader: ParquetRecordBatchReader) -> Self {
5777            let mut values = vec![];
5778            let mut row_numbers = vec![];
5779            for batch in reader {
5780                let batch = batch.expect("Could not read batch");
5781                values.extend(
5782                    batch
5783                        .column_by_name("col")
5784                        .expect("Could not get col column")
5785                        .as_primitive::<arrow::datatypes::Int64Type>()
5786                        .iter()
5787                        .map(|v| v.expect("Could not get value")),
5788                );
5789
5790                row_numbers.extend(
5791                    batch
5792                        .column_by_name("row_number")
5793                        .expect("Could not get row_number column")
5794                        .as_primitive::<arrow::datatypes::Int64Type>()
5795                        .iter()
5796                        .map(|v| v.expect("Could not get row number"))
5797                        .collect::<Vec<_>>(),
5798                );
5799            }
5800            Self {
5801                values,
5802                row_numbers,
5803            }
5804        }
5805    }
5806
5807    #[test]
5808    fn test_with_virtual_columns_rejects_non_virtual_fields() {
5809        // Try to pass a regular field (not a virtual column) to with_virtual_columns
5810        let regular_field = Arc::new(Field::new("regular_column", ArrowDataType::Int64, false));
5811        assert_eq!(
5812            ArrowReaderOptions::new()
5813                .with_virtual_columns(vec![regular_field])
5814                .unwrap_err()
5815                .to_string(),
5816            "Parquet error: Field 'regular_column' is not a virtual column. Virtual columns must have extension type names starting with 'arrow.virtual.'"
5817        );
5818    }
5819
5820    #[test]
5821    #[cfg_attr(miri, ignore)] // Takes too long
5822    fn test_row_numbers_with_multiple_row_groups() {
5823        test_row_numbers_with_multiple_row_groups_helper(
5824            false,
5825            |path, selection, _row_filter, batch_size| {
5826                let file = File::open(path).unwrap();
5827                let row_number_field = Arc::new(
5828                    Field::new("row_number", ArrowDataType::Int64, false)
5829                        .with_extension_type(RowNumber),
5830                );
5831                let options = ArrowReaderOptions::new()
5832                    .with_virtual_columns(vec![row_number_field])
5833                    .unwrap();
5834                let reader = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
5835                    .unwrap()
5836                    .with_row_selection(selection)
5837                    .with_batch_size(batch_size)
5838                    .build()
5839                    .expect("Could not create reader");
5840                reader
5841                    .collect::<Result<Vec<_>, _>>()
5842                    .expect("Could not read")
5843            },
5844        );
5845    }
5846
5847    #[test]
5848    #[cfg_attr(miri, ignore)] // Takes too long
5849    fn test_row_numbers_with_multiple_row_groups_and_filter() {
5850        test_row_numbers_with_multiple_row_groups_helper(
5851            true,
5852            |path, selection, row_filter, batch_size| {
5853                let file = File::open(path).unwrap();
5854                let row_number_field = Arc::new(
5855                    Field::new("row_number", ArrowDataType::Int64, false)
5856                        .with_extension_type(RowNumber),
5857                );
5858                let options = ArrowReaderOptions::new()
5859                    .with_virtual_columns(vec![row_number_field])
5860                    .unwrap();
5861                let reader = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
5862                    .unwrap()
5863                    .with_row_selection(selection)
5864                    .with_batch_size(batch_size)
5865                    .with_row_filter(row_filter.expect("No filter"))
5866                    .build()
5867                    .expect("Could not create reader");
5868                reader
5869                    .collect::<Result<Vec<_>, _>>()
5870                    .expect("Could not read")
5871            },
5872        );
5873    }
5874
5875    #[test]
5876    fn test_read_row_group_indices() {
5877        // create a parquet file with 3 row groups, 2 rows each
5878        let array1 = Int64Array::from(vec![1, 2]);
5879        let array2 = Int64Array::from(vec![3, 4]);
5880        let array3 = Int64Array::from(vec![5, 6]);
5881
5882        let batch1 =
5883            RecordBatch::try_from_iter(vec![("value", Arc::new(array1) as ArrayRef)]).unwrap();
5884        let batch2 =
5885            RecordBatch::try_from_iter(vec![("value", Arc::new(array2) as ArrayRef)]).unwrap();
5886        let batch3 =
5887            RecordBatch::try_from_iter(vec![("value", Arc::new(array3) as ArrayRef)]).unwrap();
5888
5889        let mut buffer = Vec::new();
5890        let options = WriterProperties::builder()
5891            .set_max_row_group_row_count(Some(2))
5892            .build();
5893        let mut writer = ArrowWriter::try_new(&mut buffer, batch1.schema(), Some(options)).unwrap();
5894        writer.write(&batch1).unwrap();
5895        writer.write(&batch2).unwrap();
5896        writer.write(&batch3).unwrap();
5897        writer.close().unwrap();
5898
5899        let file = Bytes::from(buffer);
5900        let row_group_index_field = Arc::new(
5901            Field::new("row_group_index", ArrowDataType::Int64, false)
5902                .with_extension_type(RowGroupIndex),
5903        );
5904
5905        let options = ArrowReaderOptions::new()
5906            .with_virtual_columns(vec![row_group_index_field.clone()])
5907            .unwrap();
5908        let mut arrow_reader =
5909            ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options)
5910                .expect("reader builder with virtual columns")
5911                .build()
5912                .expect("reader with virtual columns");
5913
5914        let batch = arrow_reader.next().unwrap().unwrap();
5915
5916        assert_eq!(batch.num_columns(), 2);
5917        assert_eq!(batch.num_rows(), 6);
5918
5919        assert_eq!(
5920            batch
5921                .column(0)
5922                .as_primitive::<types::Int64Type>()
5923                .iter()
5924                .collect::<Vec<_>>(),
5925            vec![Some(1), Some(2), Some(3), Some(4), Some(5), Some(6)]
5926        );
5927
5928        assert_eq!(
5929            batch
5930                .column(1)
5931                .as_primitive::<types::Int64Type>()
5932                .iter()
5933                .collect::<Vec<_>>(),
5934            vec![Some(0), Some(0), Some(1), Some(1), Some(2), Some(2)]
5935        );
5936    }
5937
5938    #[test]
5939    fn test_read_only_row_group_indices() {
5940        let array1 = Int64Array::from(vec![1, 2, 3]);
5941        let array2 = Int64Array::from(vec![4, 5]);
5942
5943        let batch1 =
5944            RecordBatch::try_from_iter(vec![("value", Arc::new(array1) as ArrayRef)]).unwrap();
5945        let batch2 =
5946            RecordBatch::try_from_iter(vec![("value", Arc::new(array2) as ArrayRef)]).unwrap();
5947
5948        let mut buffer = Vec::new();
5949        let options = WriterProperties::builder()
5950            .set_max_row_group_row_count(Some(3))
5951            .build();
5952        let mut writer = ArrowWriter::try_new(&mut buffer, batch1.schema(), Some(options)).unwrap();
5953        writer.write(&batch1).unwrap();
5954        writer.write(&batch2).unwrap();
5955        writer.close().unwrap();
5956
5957        let file = Bytes::from(buffer);
5958        let row_group_index_field = Arc::new(
5959            Field::new("row_group_index", ArrowDataType::Int64, false)
5960                .with_extension_type(RowGroupIndex),
5961        );
5962
5963        let options = ArrowReaderOptions::new()
5964            .with_virtual_columns(vec![row_group_index_field.clone()])
5965            .unwrap();
5966        let metadata = ArrowReaderMetadata::load(&file, options).unwrap();
5967        let num_columns = metadata
5968            .metadata
5969            .file_metadata()
5970            .schema_descr()
5971            .num_columns();
5972
5973        let mut arrow_reader = ParquetRecordBatchReaderBuilder::new_with_metadata(file, metadata)
5974            .with_projection(ProjectionMask::none(num_columns))
5975            .build()
5976            .expect("reader with virtual columns only");
5977
5978        let batch = arrow_reader.next().unwrap().unwrap();
5979        let schema = Arc::new(Schema::new(vec![(*row_group_index_field).clone()]));
5980
5981        assert_eq!(batch.schema(), schema);
5982        assert_eq!(batch.num_columns(), 1);
5983        assert_eq!(batch.num_rows(), 5);
5984
5985        assert_eq!(
5986            batch
5987                .column(0)
5988                .as_primitive::<types::Int64Type>()
5989                .iter()
5990                .collect::<Vec<_>>(),
5991            vec![Some(0), Some(0), Some(0), Some(1), Some(1)]
5992        );
5993    }
5994
5995    #[test]
5996    fn test_read_row_group_indices_with_selection() -> Result<()> {
5997        let mut buffer = Vec::new();
5998        let options = WriterProperties::builder()
5999            .set_max_row_group_row_count(Some(10))
6000            .build();
6001
6002        let schema = Arc::new(Schema::new(vec![Field::new(
6003            "value",
6004            ArrowDataType::Int64,
6005            false,
6006        )]));
6007
6008        let mut writer = ArrowWriter::try_new(&mut buffer, schema.clone(), Some(options))?;
6009
6010        // write out 3 batches of 10 rows each
6011        for i in 0..3 {
6012            let start = i * 10;
6013            let array = Int64Array::from_iter_values(start..start + 10);
6014            let batch = RecordBatch::try_from_iter(vec![("value", Arc::new(array) as ArrayRef)])?;
6015            writer.write(&batch)?;
6016        }
6017        writer.close()?;
6018
6019        let file = Bytes::from(buffer);
6020        let row_group_index_field = Arc::new(
6021            Field::new("rg_idx", ArrowDataType::Int64, false).with_extension_type(RowGroupIndex),
6022        );
6023
6024        let options =
6025            ArrowReaderOptions::new().with_virtual_columns(vec![row_group_index_field])?;
6026
6027        // test row groups are read in reverse order
6028        let arrow_reader =
6029            ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options.clone())?
6030                .with_row_groups(vec![2, 1, 0])
6031                .build()?;
6032
6033        let batches: Vec<_> = arrow_reader.collect::<Result<Vec<_>, _>>()?;
6034        let combined = concat_batches(&batches[0].schema(), &batches)?;
6035
6036        let values = combined.column(0).as_primitive::<types::Int64Type>();
6037        let first_val = values.value(0);
6038        let last_val = values.value(combined.num_rows() - 1);
6039        // first row from rg 2
6040        assert_eq!(first_val, 20);
6041        // the last row from rg 0
6042        assert_eq!(last_val, 9);
6043
6044        let rg_indices = combined.column(1).as_primitive::<types::Int64Type>();
6045        assert_eq!(rg_indices.value(0), 2);
6046        assert_eq!(rg_indices.value(10), 1);
6047        assert_eq!(rg_indices.value(20), 0);
6048
6049        Ok(())
6050    }
6051
6052    pub(crate) fn test_row_numbers_with_multiple_row_groups_helper<F>(
6053        use_filter: bool,
6054        test_case: F,
6055    ) where
6056        F: FnOnce(PathBuf, RowSelection, Option<RowFilter>, usize) -> Vec<RecordBatch>,
6057    {
6058        let seed: u64 = random();
6059        println!("test_row_numbers_with_multiple_row_groups seed: {seed}");
6060        let mut rng = StdRng::seed_from_u64(seed);
6061
6062        use tempfile::TempDir;
6063        let tempdir = TempDir::new().expect("Could not create temp dir");
6064
6065        let (bytes, metadata) = generate_file_with_row_numbers(&mut rng);
6066
6067        let path = tempdir.path().join("test.parquet");
6068        std::fs::write(&path, bytes).expect("Could not write file");
6069
6070        let mut case = vec![];
6071        let mut remaining = metadata.file_metadata().num_rows();
6072        while remaining > 0 {
6073            let row_count = rng.random_range(1..=remaining);
6074            remaining -= row_count;
6075            case.push(RowSelector {
6076                row_count: row_count as usize,
6077                skip: rng.random_bool(0.5),
6078            });
6079        }
6080
6081        let filter = use_filter.then(|| {
6082            let filter = (0..metadata.file_metadata().num_rows())
6083                .map(|_| rng.random_bool(0.99))
6084                .collect::<Vec<_>>();
6085            let mut filter_offset = 0;
6086            RowFilter::new(vec![Box::new(ArrowPredicateFn::new(
6087                ProjectionMask::all(),
6088                move |b| {
6089                    let array = BooleanArray::from_iter(
6090                        filter
6091                            .iter()
6092                            .skip(filter_offset)
6093                            .take(b.num_rows())
6094                            .map(|x| Some(*x)),
6095                    );
6096                    filter_offset += b.num_rows();
6097                    Ok(array)
6098                },
6099            ))])
6100        });
6101
6102        let selection = RowSelection::from(case);
6103        let batches = test_case(path, selection.clone(), filter, rng.random_range(1..4096));
6104
6105        if selection.skipped_row_count() == metadata.file_metadata().num_rows() as usize {
6106            assert!(batches.into_iter().all(|batch| batch.num_rows() == 0));
6107            return;
6108        }
6109        let actual = concat_batches(batches.first().expect("No batches").schema_ref(), &batches)
6110            .expect("Failed to concatenate");
6111        // assert_eq!(selection.row_count(), actual.num_rows());
6112        let values = actual
6113            .column(0)
6114            .as_primitive::<types::Int64Type>()
6115            .iter()
6116            .collect::<Vec<_>>();
6117        let row_numbers = actual
6118            .column(1)
6119            .as_primitive::<types::Int64Type>()
6120            .iter()
6121            .collect::<Vec<_>>();
6122        assert_eq!(
6123            row_numbers
6124                .into_iter()
6125                .map(|number| number.map(|number| number + 1))
6126                .collect::<Vec<_>>(),
6127            values
6128        );
6129    }
6130
6131    fn generate_file_with_row_numbers(rng: &mut impl Rng) -> (Bytes, ParquetMetaData) {
6132        let schema = Arc::new(Schema::new(Fields::from(vec![Field::new(
6133            "value",
6134            ArrowDataType::Int64,
6135            false,
6136        )])));
6137
6138        let mut buf = Vec::with_capacity(1024);
6139        let mut writer =
6140            ArrowWriter::try_new(&mut buf, schema.clone(), None).expect("Could not create writer");
6141
6142        let mut values = 1..=rng.random_range(1..4096);
6143        while !values.is_empty() {
6144            let batch_values = values
6145                .by_ref()
6146                .take(rng.random_range(1..4096))
6147                .collect::<Vec<_>>();
6148            let array = Arc::new(Int64Array::from(batch_values)) as ArrayRef;
6149            let batch =
6150                RecordBatch::try_from_iter([("value", array)]).expect("Could not create batch");
6151            writer.write(&batch).expect("Could not write batch");
6152            writer.flush().expect("Could not flush");
6153        }
6154        let metadata = writer.close().expect("Could not close writer");
6155
6156        (Bytes::from(buf), metadata)
6157    }
6158}