Skip to main content

parquet/arrow/array_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//! Logic for reading into arrow arrays: [`ArrayReader`] and [`RowGroups`]
19
20use crate::errors::Result;
21use arrow_array::ArrayRef;
22use arrow_schema::DataType as ArrowType;
23use std::any::Any;
24use std::sync::Arc;
25
26use crate::arrow::record_reader::GenericRecordReader;
27use crate::arrow::record_reader::buffer::ValuesBuffer;
28use crate::column::page::PageIterator;
29use crate::column::reader::decoder::ColumnValueDecoder;
30use crate::file::metadata::ParquetMetaData;
31use crate::file::reader::{FilePageIterator, FileReader};
32
33mod builder;
34mod byte_array;
35mod byte_array_dictionary;
36mod byte_view_array;
37mod cached_array_reader;
38mod empty_array;
39mod fixed_len_byte_array;
40mod fixed_size_list_array;
41mod list_array;
42mod list_view_array;
43mod map_array;
44mod null_array;
45mod primitive_array;
46mod row_group_cache;
47mod row_group_index;
48mod row_number;
49mod struct_array;
50
51#[cfg(test)]
52pub(crate) mod test_util;
53
54// Note that this crate is public under the `experimental` feature flag.
55use crate::file::metadata::RowGroupMetaData;
56pub use builder::{ArrayReaderBuilder, CacheOptions, CacheOptionsBuilder};
57pub use byte_array::make_byte_array_reader;
58// Re-exported (beyond the `experimental` feature) so `file::metadata::dictionary`
59// can PLAIN-decode a raw dictionary page without duplicating this logic.
60pub(crate) use byte_array::ByteArrayDecoderPlain;
61pub use byte_array_dictionary::make_byte_array_dictionary_reader;
62#[cfg_attr(not(feature = "experimental"), expect(unused_imports))]
63pub use byte_view_array::make_byte_view_array_reader;
64#[cfg_attr(not(feature = "experimental"), expect(unused_imports))]
65pub use fixed_len_byte_array::make_fixed_len_byte_array_reader;
66pub use fixed_size_list_array::FixedSizeListArrayReader;
67pub use list_array::ListArrayReader;
68pub use list_view_array::ListViewArrayReader;
69pub use map_array::MapArrayReader;
70pub use null_array::NullArrayReader;
71pub use primitive_array::PrimitiveArrayReader;
72pub use row_group_cache::RowGroupCache;
73pub use struct_array::StructArrayReader;
74
75/// Reads Parquet data into Arrow Arrays.
76///
77/// This is an internal implementation detail of the Parquet reader, and is not
78/// intended for public use.
79///
80/// This is the core trait for reading encoded Parquet data directly into Arrow
81/// Arrays efficiently. There are various specializations of this trait for
82/// different combinations of encodings and arrays, such as
83/// [`PrimitiveArrayReader`], [`ListArrayReader`], etc.
84///
85/// Each `ArrayReader` logically contains the following state
86/// 1. A handle to the encoded Parquet data
87/// 2. An in progress buffered Array
88///
89/// Data can either be read in batches using [`ArrayReader::next_batch`] or
90/// incrementally using [`ArrayReader::read_records`] and [`ArrayReader::skip_records`].
91pub trait ArrayReader: Send {
92    // TODO: this function is never used, and the trait is not public. Perhaps this should be
93    // removed.
94    fn as_any(&self) -> &dyn Any;
95
96    /// Returns the arrow type of this array reader.
97    fn get_data_type(&self) -> &ArrowType;
98
99    /// Reads at most `batch_size` records into an arrow array and return it.
100    #[cfg(any(feature = "experimental", test))]
101    fn next_batch(&mut self, batch_size: usize) -> Result<ArrayRef> {
102        self.read_records(batch_size)?;
103        self.consume_batch()
104    }
105
106    /// Reads at most `batch_size` records' bytes into buffer
107    ///
108    /// Returns the number of records read, which can be less than `batch_size` if
109    /// pages is exhausted.
110    fn read_records(&mut self, batch_size: usize) -> Result<usize>;
111
112    /// Consume all currently stored buffer data
113    /// into an arrow array and return it.
114    fn consume_batch(&mut self) -> Result<ArrayRef>;
115
116    /// Skips over `num_records` records, returning the number of rows skipped
117    ///
118    /// Note that calling `skip_records` with large values of `num_records` is
119    /// efficient as it avoids decoding data into the the in-progress array.
120    /// However, there is overhead to calling this function, so for small values of
121    /// `num_records`, it can be more efficient to call read_records and apply
122    /// a filter to the resulting array.
123    fn skip_records(&mut self, num_records: usize) -> Result<usize>;
124
125    /// If this array has a non-zero definition level, i.e. has a nullable parent
126    /// array, returns the definition levels of data from the last call of `next_batch`
127    ///
128    /// Otherwise returns None
129    ///
130    /// This is used by parent [`ArrayReader`] to compute their null bitmaps
131    fn get_def_levels(&self) -> Option<&[i16]>;
132
133    /// If this array has a non-zero repetition level, i.e. has a repeated parent
134    /// array, returns the repetition levels of data from the last call of `next_batch`
135    ///
136    /// Otherwise returns None
137    ///
138    /// This is used by parent [`ArrayReader`] to compute their array offsets
139    fn get_rep_levels(&self) -> Option<&[i16]>;
140
141    /// Returns the maximum definition level for this reader as defined by
142    /// the Parquet schema. For leaf readers this is the column's max def level;
143    /// for composite readers it is the level at which the composite itself is
144    /// fully defined.
145    ///
146    /// The default panics. Synthetic readers that are not backed by a Parquet
147    /// column (e.g. row-number or empty-struct readers) may rely on this
148    /// default because they are never used as children of schema-driven
149    /// composite readers (list, map, struct), which are the only callers.
150    fn max_def_level(&self) -> i16 {
151        panic!("max_def_level called on a reader that does not track definition levels")
152    }
153}
154
155/// Interface for reading data pages from the columns of one or more RowGroups.
156pub trait RowGroups {
157    /// Get the number of rows in this collection
158    fn num_rows(&self) -> usize;
159
160    /// Returns a [`PageIterator`] for all pages in the specified column chunk
161    /// across all row groups in this collection.
162    fn column_chunks(&self, i: usize) -> Result<Box<dyn PageIterator>>;
163
164    /// Returns an iterator over the row groups in this collection
165    ///
166    /// Note this may not include all row groups in [`Self::metadata`].
167    fn row_groups(&self) -> Box<dyn Iterator<Item = &RowGroupMetaData> + '_>;
168
169    /// Returns the parquet metadata
170    fn metadata(&self) -> &ParquetMetaData;
171}
172
173impl RowGroups for Arc<dyn FileReader> {
174    fn num_rows(&self) -> usize {
175        FileReader::metadata(self.as_ref())
176            .file_metadata()
177            .num_rows() as usize
178    }
179
180    fn column_chunks(&self, column_index: usize) -> Result<Box<dyn PageIterator>> {
181        let iterator = FilePageIterator::new(column_index, Arc::clone(self))?;
182        Ok(Box::new(iterator))
183    }
184
185    fn row_groups(&self) -> Box<dyn Iterator<Item = &RowGroupMetaData> + '_> {
186        Box::new(FileReader::metadata(self.as_ref()).row_groups().iter())
187    }
188
189    fn metadata(&self) -> &ParquetMetaData {
190        FileReader::metadata(self.as_ref())
191    }
192}
193
194/// Uses `record_reader` to read up to `batch_size` records from `pages`
195///
196/// Returns the number of records read, which can be less than `batch_size` if
197/// pages is exhausted.
198fn read_records<V, CV>(
199    record_reader: &mut GenericRecordReader<V, CV>,
200    pages: &mut dyn PageIterator,
201    batch_size: usize,
202) -> Result<usize>
203where
204    V: ValuesBuffer,
205    CV: ColumnValueDecoder<Buffer = V>,
206{
207    let mut records_read = 0usize;
208    while records_read < batch_size {
209        let records_to_read = batch_size - records_read;
210
211        let records_read_once = record_reader.read_records(records_to_read)?;
212        records_read += records_read_once;
213
214        // Record reader exhausted
215        if records_read_once < records_to_read {
216            if let Some(page_reader) = pages.next() {
217                // Read from new page reader (i.e. column chunk)
218                record_reader.set_page_reader(page_reader?)?;
219            } else {
220                // Page reader also exhausted
221                break;
222            }
223        }
224    }
225    Ok(records_read)
226}
227
228/// Uses `record_reader` to skip up to `batch_size` records from `pages`
229///
230/// Returns the number of records skipped, which can be less than `batch_size` if
231/// pages is exhausted
232fn skip_records<V, CV>(
233    record_reader: &mut GenericRecordReader<V, CV>,
234    pages: &mut dyn PageIterator,
235    batch_size: usize,
236) -> Result<usize>
237where
238    V: ValuesBuffer,
239    CV: ColumnValueDecoder<Buffer = V>,
240{
241    let mut records_skipped = 0usize;
242    while records_skipped < batch_size {
243        let records_to_read = batch_size - records_skipped;
244
245        let records_skipped_once = record_reader.skip_records(records_to_read)?;
246        records_skipped += records_skipped_once;
247
248        // Record reader exhausted
249        if records_skipped_once < records_to_read {
250            if let Some(page_reader) = pages.next() {
251                // Read from new page reader (i.e. column chunk)
252                record_reader.set_page_reader(page_reader?)?;
253            } else {
254                // Page reader also exhausted
255                break;
256            }
257        }
258    }
259    Ok(records_skipped)
260}